ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python混合调度架构:定时任务与事件驱动的高效实践

Python混合调度架构:定时任务与事件驱动的高效实践 “大内密探·案卷 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 或 CeleryA 任务看起来用系统 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 apscheduler3.10.4 redis5.0.0 sqlalchemy2.0.23 psycopg2-binary2.9.9 requests2.31.0 aiohttp3.9.1 python-dotenv1.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(levellogging.INFO) redis_client aioredis.from_url( redis://localhost:6379/0, max_connections10, decode_responsesTrue, ) async def main(): scheduler AsyncIOScheduler(timezoneAsia/Shanghai) scheduler.add_job( scan_task, CronTrigger(minute*/3), idscan_job, replace_existingTrue, ) scheduler.start() listener_task asyncio.create_task(consume_events()) try: await asyncio.Event().wait() except (KeyboardInterrupt, SystemExit): pass finally: listener_task.cancel() scheduler.shutdown(waitFalse) 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, nxTrue, ex180) 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}, timeout5) 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, jsondata, timeoutaiohttp.ClientTimeout(total5), ) 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, reasonresult.get(msg)) except Exception as exc: await db_update_status(record_id, failed, reasonstr(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_SIZE500 条超过这个量事务时间会变长锁表风险上升CONCURRENCY4表间并发度再高就会出现锁等待EXTERNAL_TIMEOUT5 秒外部接口超时必须显式设置RETRY_COUNT3 次失败重试次数超过就进失败数据包REDIS_POOL_SIZE10连接池大小满足日常峰值即可PGBOUNCER_MODEtransaction事务级连接池适合短事务场景这些值不是拍脑袋定的。批大小 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 LOCKEDRedis 消息丢失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”虽然版本号小但思路放在任何后台调度系统里都成立。你要是也在做类似的东西可以从公交和专车的划分开始试起大部分调度难题都会变得好解很多。
返回列表