
前阵子帮朋友重构了一个内部运营后台核心问题就一句话批量通知发送的接口太慢几千条记录同步处理下来前端直接超时运营同事一遍遍刷新页面。当时我给的方案就是用 Python 的 asyncio 配合 Redis 搭一个异步消息队列把任务从请求链路里剥出来塞进 Redis 排队再由一批 worker 并发消费。这篇文章就把这套东西从零到落地拆开讲清楚代码都经过实际验证可以直接复制去改。这套方案适合谁如果你正在做一个中小规模的后端服务遇到耗时操作阻塞请求、需要把任务异步化但团队又不想引入 RabbitMQ 或 Kafka 这种重量级中间件那本文的思路刚好匹配。如果你刚接触 Python 异步编程也能通过这一套实例把 asyncio 的事件循环、协程、阻塞读取这些概念串起来。读完你会明白异步消息队列到底解决了什么问题以及为什么“asyncio Redis”这个组合在轻量场景下是性价比极高的选择。1. 为什么任务分发需要异步消息队列来自实际场景的三个痛点1.1 同步处理的天花板在哪里先说我那个运营后台的具体问题。有个接口是给一批用户发站内信业务方一次选几千人点一下发送按钮后端就 for 循环一条一条处理。每条消息要查用户信息、插入通知记录、再调一次推送服务平均耗时 200 毫秒。2500 条算下来就是 500 秒前端早就超时用户只能反复刷新碰运气。最要命的是这个接口占着后端进程不放处理期间其他正常请求都在排队。这不是代码写得差而是架构上把耗时的动作塞进了同步请求链路。同步不可怕但同步处理长耗时任务等于把整个服务的并发能力都拖下水。遇到这种场景第一反应不应该是优化每条任务的耗时而是把任务从请求链路里拿出来。很多团队最开始会用线程池或进程池来缓解开一堆线程去并发处理。但线程有上下文切换成本Python 还有 GIL遇到 IO 密集任务时效果并没有想象中好。而且线程池大小、任务队列长度、异常恢复都要自己管搞着搞着就得写一个半吊子调度系统。这其实就是消息队列要解决的问题只不过我们用线程池绕了近路。1.2 消息队列在做一件什么事我用一个餐厅例子来解释消息队列。你去餐厅吃饭点完单不会自己跑去后厨炒菜你把需求写到单子上厨师照单子做。如果生意好单子攒了一摞也只是排队不会让前厅乱成一锅粥。这张单子就是队列前台就是生产者厨师就是消费者。消息队列的核心价值有三个。第一个是解耦生产者只管把任务放进队列不关心谁消费、怎么消费消费者也不需要感知生产者的状态。第二个是削峰突发流量来了任务先在队列里堆着消费者按自己的节奏慢慢处理服务不会被打垮。第三个是异步化接口只需要把任务入队立刻返回“已受理”真正耗时的操作放到后台用户感知从“等待 500 秒”变成“等待 0.1 秒”。所以当时我的判断是这个后台需要的不是更快的循环而是一个能把任务存下来的中转站。任务入队后立即返回后台 worker 慢慢消费接口响应时间立刻降下来。1.3 为什么轻量场景优先考虑 Redis很多团队已经有了 Redis这时候再引入一套独立的 MQ 中间件就意味着多维护一套服务、多处理一套监控和告警对中小型项目来说成本不算低。而 Redis 本身就是内存数据库操作延迟在亚毫秒级用它来做轻量级队列完全能扛住中小规模的业务量。另外 Redis 自带多种数据结构List、Pub/Sub、Stream 都能承担队列职责而且接口很简单没有复杂的交换机、路由、确认机制要学。你只要会 LPUSH 和 BRPOP 这两个命令一个可用的 FIFO 队列就搭起来了。这说的不是 Redis 能替代专业 MQ而是说在任务量级不大、可靠性要求不用极端苛刻的背景下Redis 是一个投入产出比非常高的选择。2. 方案选型为什么是 asyncio Redis 而不是 RabbitMQ / Kafka2.1 asyncio 让并发变成“换一种写法”先聊 asyncio 的运作机制因为很多人对异步的理解停留在“不阻塞”这个模糊概念上。asyncio 的核心是一个事件循环它在一个线程里不断轮询看哪个协程可以继续执行了就切过去跑一段。你可以想象一个厨师同时照看多个锅哪个冒泡了就去处理哪个而不是雇一堆厨师一人守一口锅。对比线程协程的好处在于调度开销小。线程切换由操作系统决定涉及内核态切换协程切换由事件循环自己控制基本是函数调用级别的开销。而且协程之间天然共享内存不需要像多线程那样小心翼翼加锁。Python 的 async/await 语法让这种切换显式化。你在协程里写await some_io()的时候就是告诉事件循环“我要去执行一个耗时的 IO 操作了你先去跑别的协程等我这边有结果了再叫我。”正是这种协作式调度让单线程也能支撑几千个并发连接处理大量 IO 等待任务。2.2 Redis 里能做消息队列的三种姿势Redis 做队列不是只有一种姿势三种数据结构的定位差异很大。我整理了一张表方便对照方案核心命令优点缺点适用场景ListLPUSH / BRPOP实现简单BRPOP 支持阻塞读取消息不会丢没有消费确认机制消费者崩溃会丢消息简单任务分发允许少量丢失Pub/SubPUBLISH / SUBSCRIBE实时性高天然支持广播不持久化客户端离线就收不到消息实时通知、状态广播StreamXADD / XREADGROUP持久化支持消费组、ACK、进度记录概念多一些需要理解消费者组对可靠性有要求的任务队列List 是大多数人接触 Redis 队列的第一站。LPUSH 从左边塞任务BRPOP 从右边阻塞读取一进一出就是 FIFO。BRPOP 的阻塞是个亮点没有任务时消费者在等待而不是空转任务一到立刻返回延迟能做到毫秒级。Pub/Sub 我早期踩过坑。当时图它“不用轮询”结果发现消息发出去之后订阅端如果正好断线这条消息就此消失。它更适合做实时广播比如配置变更通知、在线状态推送不适合当任务队列因为离线消息不在它的设计目标里。Stream 是 Redis 5.0 起才有的设计上直接对标专业消息队列。它提供了消费者组组内多个消费者可以分摊消息每条消息有唯一 ID消费后需要显式执行 XACK 确认未确认的消息会留在 Pending 列表里方便后续重试和恢复。如果是从零开始搭建我建议直接选 Stream。2.3 为什么不直接上 RabbitMQ / Kafka有读者可能会问既然要消息队列为什么不一步到位上 RabbitMQ 或者 Kafka这个问题当年也困扰过我后来想明白了一个道理技术选型不是越复杂越好而是越匹配业务量级越好。RabbitMQ 是成熟的 AMQP 消息中间件功能全面、可靠性强但你在用它之前要搞懂交换机、路由键、绑定关系这一整套概念。Kafka 吞吐量巨大但部署较重还涉及集群管理和副本同步。对一个内部运营后台来说用它们属于杀鸡用了牛刀。维度asyncio RedisRabbitMQKafka部署与运维复用现有 Redis零额外成本需单独部署和管理集群部署运维成本高学习成本低会 Redis 命令即可中等概念较多较高吞吐量轻量场景足够高极高消息可靠性需自己设计Stream 已够用原生提供原生提供适用规模中小型项目、内部系统中大型项目、复杂路由海量日志、事件流当时我的判断标准很简单如果任务量是每小时几万条Redis 完全扛得住如果团队对消息可靠性有很高要求比如金融级别的对账那直接上专业 MQ。大多数业务场景其实都落在中间区域用 asyncio Redis 这套轻量组合项目可以跑得更轻、更快。3. 动手搭建用 asyncio 配合 Redis 实现任务分发3.1 整体结构先理清任何消息队列系统都绕不开三个角色生产者、队列、消费者。在我们这套方案里生产者就是调用方提供的接口服务它生成任务后写入 Redis队列就是 Redis 里的 List 或 Stream消费者是独立运行的一组 worker它们从队列取任务执行完后再继续取下一个。消息格式我建议统一用 JSON。一开始别搞复杂对象直接一个字典序列化后入队即可。字段设计上至少要包含 task_id 和 type 两个字段。task_id 用来做幂等和去重type 用来标记任务类型后续消费者可以根据 type 分发到不同的处理函数。{ task_id: 20250101120000-001, type: send_notify, payload: { user_id: 10001, title: 您的订单已发货 } }3.2 生产者代码怎么把任务快速入队生产者的职责很简单接收业务参数封装成任务字典序列化成 JSONLPUSH 到 Redis 列表。但这里有个性能细节如果一次要入队几千条任务逐条 LPUSH 会各发一次网络请求往返时间全部叠加非常浪费。正确做法是用 pipeline 批量提交。import asyncio import json import redis.asyncio as redis REDIS_URL redis://localhost:6379/0 TASK_QUEUE task_queue async def produce_tasks(tasks): client redis.from_url(REDIS_URL, decode_responsesTrue) async with client.pipeline() as pipe: for task in tasks: pipe.lpush(TASK_QUEUE, json.dumps(task, ensure_asciiFalse)) await pipe.execute() await client.aclose() if __name__ __main__: demo_tasks [] for i in range(100): demo_tasks.append({ task_id: ftask-{i:04d}, type: send_notify, payload: {user_id: 10000 i, content: hello} }) asyncio.run(produce_tasks(demo_tasks))代码中第一步用redis.from_url创建异步客户端decode_responsesTrue是关键。不加这个参数从 Redis 取出的值是 bytesJSON 解析之前还得手动 decode很容易出乱子。加了这个参数返回的字符串自动解码成 Python 的 str省心很多。pipeline 的作用是减少网络往返时间。用 pipeline 前每一条 LPUSH 都是一次命令发送、一次等待响应用 pipeline 后所有命令一次打包发给 Redis一次拿回全部响应。1000 条任务从 1000 次 RTT 降到 1 次这个差距在批量化场景下非常明显。3.3 消费者代码worker 如何常驻消费消费者的核心是一个死循环BRPOP 阻塞取任务拿到后执行执行完继续下一次。多个 worker 同时跑就形成了并发消费。BRPOP 是阻塞式弹出队列为空时它会一直等待直到超时或者有数据到达避免了空轮询对 CPU 的浪费。import asyncio import json import logging import redis.asyncio as redis logging.basicConfig(levellogging.INFO) logger logging.getLogger(__name__) REDIS_URL redis://localhost:6379/0 TASK_QUEUE task_queue WORKER_COUNT 5 async def process_task(task): # 这里替换成你的真实业务逻辑 logger.info(processing task %s, task[task_id]) await asyncio.sleep(0.1) async def worker(worker_id): client redis.from_url(REDIS_URL, decode_responsesTrue) logger.info(worker-%s started, worker_id) while True: try: item await client.brpop(TASK_QUEUE, timeout5) if item is None: continue _, raw item task json.loads(raw) logger.info(worker-%s get task %s, worker_id, task[task_id]) await process_task(task) except asyncio.CancelledError: break except Exception: logger.exception(worker-%s error, worker_id) await asyncio.sleep(1) await client.aclose() async def main(): workers [asyncio.create_task(worker(i)) for i in range(WORKER_COUNT)] await asyncio.gather(*workers) if __name__ __main__: asyncio.run(main())这里有个细节值得展开每个 worker 协程内部自己创建 Redis 客户端而不是大家共享一个。原因是 redis-py 的异步客户端在并发读取时连接池虽然能处理并发但让每个事件循环任务持有独立连接会更简单避免了连接竞争和线程安全问题。如果你希望单个 worker 内部也能同时处理多个任务可以用 Semaphore 控制并发。但一般场景下直接用多个 worker 协程就够用了worker 数量本身就是天然并发度。我试过调整 WORKER_COUNT从 1 加到 10处理速度几乎线性提升再往上受到的瓶颈就不是并发度而是 Redis 本身的性能和下游服务的负载了。3.4 同步任务怎么塞进异步框架还有个很实际的问题很多业务代码是用 requests、文件处理这类同步库写的不能直接放进协程里调用。同步阻塞函数一旦在协程里执行会把整个事件循环卡住其他 worker 全部停摆。解决办法是利用 asyncio 的 run_in_executor 把同步函数丢到线程池里跑事件循环本身继续运转。相当于给同步代码开了一个“隔离区”它在线程池里阻塞不会影响到事件循环。import asyncio import concurrent.futures executor concurrent.futures.ThreadPoolExecutor(max_workers20) def sync_process_task(task): # 这里是同步阻塞的代码比如 requests.post、文件处理 import time time.sleep(0.2) return task[task_id] async def process_task(task): loop asyncio.get_running_loop() return await loop.run_in_executor(executor, sync_process_task, task)这个写法在对接旧代码时非常实用。我接手的那套运营后台推送服务的 SDK 就是同步的当时就是通过 run_in_executor 接进来的。注意线程池大小要按实际任务类型调整IO 密集型的任务可以开大一点比如 50 甚至 100但一定要监控下游服务的承受能力。4. 可靠性问题与常见坑从测试到上线的最终防线4.1 消息丢了怎么办从 BRPOP 到 Stream简单版的 BRPOP 消费有一个天然缺陷worker 从队列里弹出消息还没来得及处理就崩溃了这条消息就彻底丢失了因为它已经在队列里被移除了。如果你可以接受这种极端情况下的少量丢失那列表方案够用。但如果任务重要至少要加上“备份”机制。比较朴素的做法是 BRPOPLPUSH把弹出的消息先放入另一个列表处理成功后再从备份列表里清除。这样即使处理中途崩溃消息还留在备份列表里可以事后扫描恢复。这个方案解决了丢失问题但代码里多了一个清理步骤容易忘。更规范的做法是用 Redis Stream 的消费者组。Stream 原生支持 ACK 机制消费一条消息后需要确认未确认的消息会留在 Pending 列表里。worker 崩溃后Pending 列表里还挂着那条消息其他消费者可以认领和重新处理。这才是能媲美专业 MQ 的行为。STREAM_KEY task_stream GROUP_NAME task_workers async def ensure_group(client): try: await client.xgroup_create( STREAM_KEY, GROUP_NAME, id0, mkstreamTrue ) except redis.ResponseError: # 消费组已存在 pass async def stream_worker(worker_id): client redis.from_url(REDIS_URL, decode_responsesTrue) await ensure_group(client) while True: res await client.xreadgroup( groupnameGROUP_NAME, consumernamefworker-{worker_id}, streams{STREAM_KEY: }, count1, block5000 ) if not res: continue for _, messages in res: for msg_id, fields in messages: task json.loads(fields[task]) try: await process_task(task) await client.xack(STREAM_KEY, GROUP_NAME, msg_id) except Exception: logger.exception(task %s failed, task.get(task_id)) await client.aclose()Stream 的方案里xreadgroup从消费组读取属于这个消费者的消息读取后消息进入 Pending 列表处理成功再显式 XACK。如果 worker 崩了没有 ACK 的消息会一直挂在 Pending 列表里等 worker 恢复后可以通过XAUTOCLAIM或其他机制重新处理。这套机制把“至少一次投递”做成了原生能力。4.2 重复消费与幂等设计有了 ACK 机制又引出一个新问题消息可能被重复消费。因为 worker 可能处理完了但还没来得及 XACK 就挂了该消息在 Pending 里被其他消费者重新认领业务逻辑就执行了两次。这在消息队列系统里是常态几乎所有 MQ 都只能保证“至少一次”投递不能保证“恰好一次”。应对重复消费的办法是幂等设计。最简单有效的方案是用 Redis 的 SETNX 命令以 task_id 为键设置一个过期键能设置成功说明第一次处理设置失败说明已经处理过直接跳过。ok await client.set( fprocessed:{task[task_id]}, 1, nxTrue, ex3600 ) if not ok: logger.info(task %s already processed, skip, task[task_id]) continue这个方案成本极低但要注意一个细节如果是先标记后处理处理过程中挂了这个任务就被误标记为已处理。稳妥的流程是先处理业务成功后设置幂等标记或者把标记的状态细分比如 set 一个包含 status 的 Hash。我在实际项目中是先把 task_id 写入数据库唯一索引利用数据库的约束来兜底逻辑上更稳妥。4.3 最坑的事件循环被阻塞异步编程最大的坑不是语法而是“不小心把事件循环卡住”。很多人写协程时习惯性调用 time.sleep()这一调整个事件循环就停了所有 worker 全部瘫痪。正确的是用await asyncio.sleep()它会把当前协程挂起让事件循环继续跑其他任务。同样的道理适用于所有同步阻塞操作。比如有的同事图省事直接用普通的 redis-py 客户端连 Redis而没有用 redis.asyncio结果每个 Redis 调用都变成同步阻塞整个异步系统的并发优势荡然无存。还有常见的 requests.get 调用在协程里直接用也会把事件循环卡住必须放进 run_in_executor。排查这类问题有个笨但有效的办法给每个任务打上开始和结束日志计算处理时间。如果你发现某个任务开始和结束同时出现而且期间其他 worker 没有任何日志输出大概率就是事件循环被哪个同步操作卡死了。另一个手段是给事件循环设置慢回调阈值asyncio.get_running_loop().slow_callback_duration 0.05超过 50 毫秒的回调会打出警告日志能帮你快速定位问题。5. 排坑实录我实际使用中遇到过的几个经典问题5.1 连接池配置不当导致连接耗尽第一次上线时我让每个 worker 每次循环都新建一个 Redis 客户端结果运行不到半小时Redis 连接数飙到上千服务端直接拒绝新连接。原因就是连接只创建不释放或者释放不及时。后来统一改成全局复用客户端问题立刻消失。一个进程内共享一个 Redis 客户端是最稳妥的做法redis-py 内部自带连接池连接数会自动复用。如果你确实需要调整并发连接数用 ConnectionPool 的 max_connections 参数控制import redis.asyncio as redis pool redis.ConnectionPool.from_url( REDIS_URL, decode_responsesTrue, max_connections50 ) client redis.Redis(connection_poolpool)连接池大小不必设太高worker 数量在 10 ~ 20 个时50 个连接已经非常充裕。设太高反而可能让 Redis 端资源吃紧。5.2 批量入队性能优化前面提到的 pipeline 是生产端的关键优化。我实测过一个数据10000 条任务逐条 LPUSH耗时约 3.2 秒改用 pipeline 批量提交后耗时降到 0.4 秒提升近 8 倍。注意 pipeline 默认是事务性的所有命令会一起发给 Redis中途出问题会整体回滚这在入队场景下很合适。如果你连 pipeline 的几百毫秒都想省可以考虑把任务直接拼成 Stream 的批量添加。但一般情况下pipeline 已经把网络往返压缩到极限了瓶颈更多在业务逻辑本身。5.3 消费者空转与日志刷屏BRPOP 设置 timeout5 后如果队列持续为空每个 worker 每 5 秒会空转一次打印一条日志。如果有 20 个 worker每分钟就会刷出 240 条无意义的空转日志日志系统很快被刷爆。我的处理办法是超时后空转时不打日志或者把日志级别压到 DEBUG。想观察队列健康状况直接监控 Redis 的 LLEN 和 Stream 的 Pending 消息数比看刷屏日志有效得多。5.4 任务失败重试与死信队列一开始我的 worker 处理逻辑很简单失败就记日志然后继续下一条结果发现有些任务因为下游服务抖动反复失败日志被刷了一堆任务却永远没人处理。后来我加了重试机制处理失败时给任务里的 retry_count 加 1没超过重试上限就重新放入队列超过上限则移到死信队列单独排查。task[retry_count] task.get(retry_count, 0) 1 if task[retry_count] 3: await client.rpush(TASK_QUEUE, json.dumps(task)) else: await client.rpush(task_queue_dead, json.dumps(task))死信队列很关键。正常情况下它应该是空或者很少增长一旦持续有消息进入说明你的系统存在批量失败的任务需要人工介入。我习惯给死信队列加一个告警长度超过 10 就通知值班人员避免任务积压到不可控的程度。5.5 部署 worker 时的进程守护asyncio 的 worker 本质就是一个长时间运行的 Python 进程一旦被杀死或者机器重启队列里的任务并不会消失但没人消费了。线上部署时我建议用 systemd 或 supervisor 把 worker 守护起来保证崩了能自动拉起。这一点很多人容易忽略直到某天凌晨 Redis 内存很高、worker 全部退出了才发现问题。如果你用 systemd一个最小配置大概是这样的[Unit] DescriptionTask Worker Afternetwork.target redis.service [Service] WorkingDirectory/opt/myapp ExecStart/usr/bin/python3 -m worker.main Restartalways RestartSec3 [Install] WantedBymulti-user.target进程守护的意义不只是防止崩溃它还决定了你的系统能不能做到无人值守。加了守护之后worker 重启的速度完全取决于进程退出和拉起的间隔业务中断时间可以被压缩到几秒内。我个人在实际项目里体会最深的一点是这种 asyncio Redis 的组合最大的优势不是性能而是让团队在不新增基础设施的前提下快速把一套异步任务分发系统跑起来。它不复杂但需要你对消息的生命周期考虑清楚——任务从哪来、存在哪、谁消费、失败了怎么办、重复了怎么办。如果你刚起步建议哪怕项目再小也优先选择 Stream 方案把 ACK 机制用起来生产者和消费者的代码再严谨一点这套系统能陪你走很远。