“大内密探·案卷 0.5:时钟坐公交,数据打专车。”这句话是我在一套 Python 调度系统改造时顺手写在白板上的备注,后来发现它比任何架构文档都好用。有一类任务像时钟一样,到了点就必须发车,沿着固定线路批量跑;另一类任务像专车,数据一到就得立刻响应,一单一单处理。这套系统基于 Python 3.11 搭了一个混合调度骨架,把定时批量扫描任务和事件驱动实时处理任务分开设计,最终把原来动辄几秒钟的外部系统调用压到了几百毫秒量级,数据库压力也明显降了下来。这篇文章想把“两套车”的设计思路、工具选型和踩坑过程完整记录下来,给正在做定时任务、数据同步、回调处理这类后台系统的朋友一个可以参考的样本。你只需要有一点 Python 基础,就能看懂大部分实现。
1. 先看清两类任务的“调性”,再决定怎么调度
1.1 “时钟坐公交”到底是什么任务
项目里的 A 任务,核心工作是定时扫描业务表里的合同状态、库存数量、截止日期这些字段,把符合条件的记录聚合成待办数据,再调用外部系统接口完成同步或提醒。这类任务有一个共同特征:时间驱动、批量处理、允许排队积压。每天跑多少次、每次扫多少数据,基本可以提前估算,晚跑几分钟也不会造成灾难性后果。
用公交车类比非常贴切——它按时刻表发车,把一批人统一运到某个站点,偶尔晚点,但整体节奏稳定。所以这类任务我们内部叫“公交任务”,它追求的是吞吐量和可靠性,而不是单条数据的毫秒级延迟。A 任务的延迟目标可以放宽到秒级甚至分钟级,只要别把数据库连接池占死就行。
1.2 “数据打专车”又是什么任务
B 任务是上游系统实时推送过来的数据处理,典型场景是支付回调、库存变更、状态流转通知。这类数据的特点是突发性强、时效性高、单条处理链路短。来一条数据就要处理一条,不能说“等下一班公交车一起走吧”——支付回调晚处理几秒钟,用户侧就能感知到异常。
专车任务要做的动作其实也不复杂:收到数据、解析、校验、写库、通知下游、更新状态。但它的难点在于不能被其他任务拖累。如果公交任务在扫库时把表锁了,专车任务的写库就会被卡住;如果公交任务在调用外部接口时超时,整个进程的并发处理能力都会下降。C 任务则是 B 任务的“后排乘客”,每处理满 N 条数据,就顺手更新一次统计信息、扫描失败数据包准备重放,所以它也属于专车体系。
1.3 为什么不能把所有任务塞进一个大循环
改造之前,这套系统的做法非常朴素:一个大循环,把所有待办任务都拉出来挨个处理。初期数据量小的时候确实能跑,但数据量上来之后,问题一个接一个冒出来。
- 互相阻塞:A 任务调用外部系统接口时,如果接口响应慢,整个循环被拖住,实时数据只能等在内存队列里,像一群乘客站在雨里等一辆堵在半路的公交。
- 重复扫描:没有统一锁和批次边界,一到调度时刻就把所有数据捞一遍,重复处理引发重复写库,下游拿到重复数据还要再过滤。
- 资源没法隔离:一次大批量扫描把数据库连接池占满,实时回调数据连进库的资格都没有,最终只能靠人工补数据。
这些问题的根源,就是把两种故障模型完全不同的任务放进了同一个执行通道。解决方案也很直白:把“时刻表逻辑”和“按需发车逻辑”彻底分开,公交走公交的道,专车走专车的道。
2. 工具选型:两套车要配两套底盘
2.1 为什么选 APScheduler 而不是 Cron 或 Celery
A 任务看起来用系统 Cron 就能实现,但在实际开发里,Cron 有几个硬伤:不好做依赖管理,重试逻辑要自己写,多实例部署时会出现“每个节点都触发一次”的问题。Celery 又太重了,为了一个低频定时任务引入 broker、worker、beat,运维成本完全划不来。
APScheduler 正好卡在中间。它可以内嵌在 Python 进程里,把任务直接注册成函数,支持 interval、cron、date 三种触发器,还能通过 jobstore 做任务持久化。我们最终选的是APScheduler 3.10,而不是 4.x,原因很现实:4.x 的 API 变化很大,迁移成本高,社区里大量资料和示例还停留在 3.x,生产系统没必要为了追新把自己搭进去。
2.2 存储层:PostgreSQL 加 pgbouncer 的搭配
A、B、C 三类任务都要读写业务表,如果每个任务各自建一个连接池,并发一起来数据库连接数会迅速膨胀。所以我们在应用层与 PostgreSQL 之间加了一层pgbouncer,用事务级连接池模式,任务事务生命周期短,连接复用率高,连接数也能稳稳压住。
表结构上至少需要两张核心表:一张叫扫描任务表,记录每次调度批次的状态;一张叫任务数据表,存每条待处理记录的明细和状态。这两张表是公交和专车的“交汇点”,A 任务负责往数据表里塞待处理数据,B 任务负责把数据处理完并更新状态。并发策略是表内串行、表间并发,这个后面会展开讲。
2.3 事件通道为什么用 Redis 而不是消息队列
B 任务的数据流是:上游回调进接口 → 写库打上待处理标记 → 发布一条 Redis 消息 → 异步消费者订阅并处理。有人会问,这个场景直接用 RabbitMQ 或 Kafka 不更正规吗?
正规是正规,但没必要。我们的消息量级远达不到消息队列的容量需求,Redis 的 Pub/Sub 足够支撑,还没有额外的中间件依赖。代价是 Pub/Sub 的消息不持久化,Redis 重启或消费者掉线时消息会丢。所以我们做了一个兜底设计:消费失败的数据会标记成失败状态,由 C 任务定期扫描失败数据包进行重放。说白了,用 Redis 省了运维复杂度,但必须用数据库状态来兜底。
3. 整体架构与三条任务链路
3.1 公交路线:定时扫描任务链路
A 任务的完整链路是:调度器触发 → 抢 Redis 分布式锁 → 从扫描任务表取批次 → 用SELECT ... FOR UPDATE SKIP LOCKED锁定待处理记录 → 按数据来源分组 → 分批调用外部系统 API → 更新任务状态。每次扫描开始前先抢锁,抢不到就说明另一个实例已经在跑,当前实例直接退出,避免多实例重复执行。
在调用外部系统时,我们坚持批内并发、批间串行。一批数据内部可以同时发出多个请求,但批次之间必须排好队,避免瞬时请求量把外部系统打崩。这里的外部调用已经从数据库直连改成了服务化 API,单次调用耗时从原来的秒级降到了几百毫秒,具体效果后面说。
3.2 专车路线:事件驱动任务链路
B 任务的链路是:上游回调进入 API 接口 → 先把数据落库并标记待处理 → 使用 Redis 发布一条事件通知 → 常驻协程订阅到消息后立即处理 → 处理成功更新状态,失败打上失败标记等待重放。这里的关键点在于,落库操作在发布消息之前完成,这样哪怕消息丢失,数据库里还留着一条待处理记录,不会出现数据凭空消失的情况。
因为 B 任务的核心是“来一个处理一个”,所以代码里必须用异步实现。同步阻塞会让 Redis 订阅一直闲着,数据多起来之后延迟会直线上升。异步协程可以在等待外部接口响应的同时,继续接收新的消息,这才是专车该有的服务体验。
3.3 后排乘客:C 任务统计与失败重放
C 任务不算独立链路,它搭在 B 任务后面。每处理成功 N 条数据,就触发一次统计更新,把当前处理总量、失败量、平均耗时写入统计表。同时,C 任务还负责定时扫描失败数据包,把这些数据重新投递到处理队列里。失败的记录不会整表重试,只重放失败的那几条,这样外部系统的调用量能大幅降下来——这就是标题里“数据打专车”的完整含义。
4. 核心代码与配置实现
原项目在验证后被格式化清掉了,我按当时的设计模式重新搭了一个最小可跑版本,用到的库和流程基本一致。下面这段可以当成一个骨架照着改。
4.1 环境准备与依赖清单
Python 版本用3.11+,理由很直接:asyncio 在 3.11 里的异常处理更友好,整体异步生态也更成熟。依赖项尽量精简:
pip install apscheduler==3.10.4 redis==5.0.0 sqlalchemy==2.0.23 psycopg2-binary==2.9.9 requests==2.31.0 aiohttp==3.9.1 python-dotenv==1.0.0SQLAlchemy 2.0 的异步查询能力和 asyncio 搭配起来很顺手,数据库驱动用的 psycopg2-binary,如果追求更高性能可以换成 asyncpg,但连接池和事务写法会有点差异。
4.2 主调度器入口
import asyncio import json import logging import redis.asyncio as aioredis from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger logging.basicConfig(level=logging.INFO) redis_client = aioredis.from_url( "redis://localhost:6379/0", max_connections=10, decode_responses=True, ) async def main(): scheduler = AsyncIOScheduler(timezone="Asia/Shanghai") scheduler.add_job( scan_task, CronTrigger(minute="*/3"), id="scan_job", replace_existing=True, ) scheduler.start() listener_task = asyncio.create_task(consume_events()) try: await asyncio.Event().wait() except (KeyboardInterrupt, SystemExit): pass finally: listener_task.cancel() scheduler.shutdown(wait=False) if __name__ == "__main__": asyncio.run(main())这里有个容易踩的坑:B 任务的消费者协程不要用scheduler.add_job(... interval ...)去注册,否则每次调度都会重新创建一个订阅协程,多个订阅者会同时消费同一频道,导致消息被分散处理,逻辑无法保证顺序。正确做法是像上面这样,单独asyncio.create_task(consume_events()),只创建一次。
4.3 公交任务核心实现
async def db_fetch_pending(limit: int) -> list[dict]: from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession query = text( "SELECT id, source, payload, status " "FROM zw_scan_time_task " "WHERE status = 'pending' " "ORDER BY id " "LIMIT :limit " "FOR UPDATE SKIP LOCKED" ) async with AsyncSession(engine) as session: result = await session.execute(query, {"limit": limit}) return [dict(row) for row in result.mappings()] async def scan_task(): lock_key = "lock:scan" got_lock = await redis_client.set(lock_key, "1", nx=True, ex=180) if not got_lock: logging.info("另一个实例正在执行扫描,本次跳过") return try: rows = await db_fetch_pending(500) groups = {} for row in rows: groups.setdefault(row["source"], []).append(row) for source, batch in groups.items(): # 批内并发,批间串行,用 asyncio.to_thread 避免阻塞事件循环 await asyncio.to_thread(call_external_api, source, batch) batch_ids = [row["id"] for row in batch] await db_mark_done(batch_ids) finally: await redis_client.delete(lock_key) def call_external_api(source: str, batch: list[dict]): # 这里可以使用 requests,因为它在子线程里运行 # 注意设置 timeout,绝对不能省略 resp = requests.post(EXTERNAL_API_URL, json={"source": source, "batch": batch}, timeout=5) resp.raise_for_status()SELECT ... FOR UPDATE SKIP LOCKED是这段代码的灵魂。它让多个实例并发扫表时,每个实例只取到属于自己的那批记录,而不是互相等锁。Redis 分布式锁解决的是“同一时刻只允许一个实例执行扫描”,SKIP LOCKED 解决的是“多个实例扫到同一批数据时自动错开”,两者并不冲突。我建议即使单实例部署也把锁加上,后面扩实例的时候就不用改代码了。
4.4 专车任务核心实现
async def consume_events(): pubsub = redis_client.pubsub() await pubsub.subscribe("event:data") async for message in pubsub.listen(): if message["type"] != "message": continue data = json.loads(message["data"]) await handle_one(data) async def handle_one(data: dict): record_id = data["record_id"] try: async with session.post( EXTERNAL_API_URL, json=data, timeout=aiohttp.ClientTimeout(total=5), ) as resp: result = await resp.json() if result.get("code") == 0: await db_update_status(record_id, "success") else: await db_update_status(record_id, "failed", reason=result.get("msg")) except Exception as exc: await db_update_status(record_id, "failed", reason=str(exc)) processed_count = await redis_client.incr("stat:processed") if processed_count % STAT_BATCH_SIZE == 0: asyncio.create_task(update_statistics())这里两个细节值得说。第一,aiohttp.ClientSession一定要全局复用,不要每次请求都新建一个,否则 TCP 连接握手开销会吞掉异步带来的性能收益。第二,asyncio.create_task(update_statistics())是异步统计的标准做法,不能让统计阻塞主处理流程,它是真的“后排乘客”——不干扰司机开车。
4.5 关键参数实测推荐值
| 参数 | 建议值 | 说明 |
|---|---|---|
| SCAN_INTERVAL | 每 3 分钟 | A 任务 cron 触发器,可以根据数据增长量调整 |
| SCAN_BATCH_SIZE | 500 条 | 超过这个量事务时间会变长,锁表风险上升 |
| CONCURRENCY | 4 | 表间并发度,再高就会出现锁等待 |
| EXTERNAL_TIMEOUT | 5 秒 | 外部接口超时,必须显式设置 |
| RETRY_COUNT | 3 次 | 失败重试次数,超过就进失败数据包 |
| REDIS_POOL_SIZE | 10 | 连接池大小,满足日常峰值即可 |
| PGBOUNCER_MODE | transaction | 事务级连接池,适合短事务场景 |
这些值不是拍脑袋定的。批大小 500 是因为实测中事务控制在 1 秒内完成,外部调用即使有一两条超时重试,整体仍然可接受。并发度 4 是因为数据库 CPU 在并发 4 时达到拐点,再高收益锐减。调优的原则是:找到瓶颈,然后让瓶颈变成可控参数,而不是盲目堆并发。
5. 踩坑记录:三次拍大腿和一张速查表
5.1 Redis 连接池被榨干
第一次压测的时候,B 任务处理速度一上来,Redis 立刻报连接超时。排查发现代码里到处是redis.Redis()现用现建,没走连接池,每次publish都要新建连接,量一大就把 Redis 的连接数打满了。后来改成启动时用aioredis.from_url创建全局连接池,所有协程复用同一个客户端,问题立刻消失。这类问题用一句话总结就是:连接必须池化,对象必须复用。
5.2 长事务卡住公交站,专车也进不了站
A 任务第二次改造时,外部接口偶尔响应要 10 秒,但requests.post没设置 timeout,请求就一直在那里等。更糟的是,外部调用在数据库事务里,A 任务的事务一直不提交,B 任务要更新的那张表就被锁住了。用户反馈“数据不实时了”,一查才发现公交把站台堵死了,专车只能在外面等。
解决方式是把外部调用的逻辑和数据库事务拆开。事务里只做状态更新,外部调用放到事务外,再给 HTTP 请求加 5 秒超时。这样就算外部系统抽风,也不会连累数据库事务。
5.3 DBLINK 直连外部库,10 秒变 200 毫秒
这个项目最开始的版本,A 任务通过数据库 DBLINK 直接查询外部系统的业务库,单次查询用时经常在 7 到 10 秒。改造后,我们把外部系统封装成 API 服务,应用层通过 HTTP 调用。同样的数据,单次调用降到了 200 毫秒左右;加上批次合并和失败数据包重放机制,外部系统总调用量下降了近 60%,数据库压力也肉眼可见地降了下来。这里我想说的是,能用服务化接口就别跨库直连,DBLINK 虽然是数据库自带的功能,但它把两个系统的耦合直接埋进了 SQL 里,排查问题时你根本不知道瓶颈在哪一端。
5.4 常见问题速查表
| 现象 | 可能原因 | 解决方式 |
|---|---|---|
| 定时任务偶尔不执行 | cron 触发器时区没配置 | AsyncIOScheduler 显式设置 timezone |
| 同一批数据被重复处理 | Redis 锁没加,或 SKIP LOCKED 缺失 | 抢锁 + FOR UPDATE SKIP LOCKED |
| Redis 消息丢失 | Pub/Sub 本身不持久化 | 落库兜底,C 任务扫描失败数据包重放 |
| 数据库连接打满 | pgbouncer 池太小或任务并发过高 | 调大 pgbouncer 池,必要时限制任务并发数 |
| 外部接口响应慢拖垮主流程 | HTTP 请求没有设置超时 | 所有外部调用统一 5 秒超时 + 3 次重试 |
6. 效果复盘与下一步还能怎么改
6.1 这套方案跑出来的实际数据
以这套系统当时的量级做参照,改造完成后的效果很直观。A 任务扫描一轮从原先的 30 秒左右缩短到 6 秒,因为 DBLINK 换成了 API 调用,批内并行也起来了;B 任务单条数据的处理延迟从平均 1.2 秒降到了 200 毫秒,用户基本感知不到回调延迟;数据库连接峰值从 80 左右降到了 20 上下,深夜低峰期甚至可以到个位数。数据库压力的下降,也直接带来了排查问题时的宁静——半夜不会再被锁表告警吵醒了。
6.2 如果让我再做一次,我会优先改三件事
第一,把事件通道从 Redis Pub/Sub 换成 Redis Stream 或真正的消息队列。Pub/Sub 的丢消息问题虽然用数据库状态兜底解决了,但兜底机制始终是事后补偿,如果能从源头做到消息持久化,整个系统的健壮性会再上一个台阶。
第二,给每个任务加一个trace_id,从 Redis 消息到数据库写入,再到外部系统调用,全部串起来。这次改造中排查问题的最大痛点就是链路不透明,一条数据失败了要翻三个日志文件才能找到根因。有了 trace_id,后续排查效率能翻倍。
第三,把 A 任务的多实例锁从 Redis 升级成数据库 jobstore,让任务本身也具备失败记录和恢复能力。简单说就是:不要满足于“系统能跑”,要让它“跑挂了还能自己爬起来”。
如果让我用一句话总结这次改造,就是:任务调度不是把代码塞进循环里就完事,而是要先区分时间驱动和事件驱动,因为这两种任务的故障模型完全不同。时间驱动任务最怕“该跑的时候没跑”,事件驱动任务最怕“该处理的时候被堵住”。分开了,故障范围就隔离了,排查问题的边界也就清晰了。这套“案卷 0.5”虽然版本号小,但思路放在任何后台调度系统里都成立。你要是也在做类似的东西,可以从公交和专车的划分开始试起,大部分调度难题都会变得好解很多。