我手头这套设备,说起来有点寒酸:一台 4 核 8G 的旧台式机当主力,两台只有 2G 内存的老笔记本在旁边待命,硬盘还是机械盘。干的事也相当朴素——帮朋友和自己批量处理视频素材,转码、抽帧、做目标检测、归档去重。一开始我觉得这活写个 Python 脚本轮询目录就行,丢到 crond 里定时跑,简单粗暴。直到有一天凌晨三点被电话吵醒,说客户等着要的片子卡在“队列”里出不来,我才第一次认真坐下来想:我这些脚本,能叫系统吗?后来真正救了我的,不是换更贵的机器,而是消息队列。
这篇文章想聊的,就是我自己从“脚本堆砌”到“消息队列 + 微服务”的完整觉醒过程。如果你手头也有一堆不算新的机器,或者正在用定时任务、共享目录、数据库表硬撑一些异步任务,那这篇内容应该能让你少走很多弯路。我会把我当时的错误设计、思考过程、选型对比、最后的落地结构和踩过的重复消费坑,全部摆出来讲。
1. 卡死的转码任务:从共享目录脚本到半夜被电话叫醒
1.1 最初的“任务系统”长什么样
我最开始的设计非常原始:一台机器上建了一个incoming目录,别的机器通过 Samba/NFS 往里丢视频文件;主机器上跑一个 Python 脚本,每 5 分钟扫描一次目录,发现新文件就往一个超时队列里丢,然后用 multiprocessing 同时起 8 个子进程去执行“转码 → 抽帧 → 识别 → 归档”这条流水线。
任务状态怎么标记?全靠改文件名后缀。处理中就把文件重命名成.processing,成功了改成.done,失败了改成.fail。听起来好像能工作,对吧?在小批量情况下确实能工作。但问题是,这套东西没有任何真正的“任务对象”,任务的状态散落在文件名、进程内存和人的记忆里。
1.2 压垮它的三个典型场景
第一次出问题是在一个周五晚上,朋友一次性给我塞了 200 个视频文件。脚本一次性扫出了 200 个任务,8 个 worker 全部启动,每台设备的 CPU 直接飙到 100%,内存开始吃 swap,然后几个转码进程被 OOM Killer 杀掉。被杀的进程留下了一堆.processing文件,脚本重启后不知道这些文件到哪一步了,只能人工去猜“这个文件到底处理完没有”。
第二次是进程崩溃后没有重试。一个视频在抽帧阶段因为内存不够崩了,文件停在.processing,后面的任务全被堵住。客户那边只看到“怎么这么慢”,我在电话这头只能对着日志一条条翻,手动把这些文件恢复到待处理状态。
第三次更恶心,两台机器同时在跑消费脚本,同时扫描到了同一个文件,于是各自为战,各处理了一遍。最后生成了两份抽帧目录、两次识别结果,我根本分不清哪一份是对的。
1.3 加机器为什么没用
出了这么多问题,我的第一反应是加机器。我把那两台笔记本也拉进来一起跑脚本,以为并发多了就快了。结果呢?文件锁问题更严重,多机同时抢同一批文件的概率更高;任务状态也彻底失控,日志分散在各台机器上,排查问题要在三台机器之间来回跳。
直到那个时候我才意识到:瓶颈根本不是 CPU,也不是内存,而是“任务”这个概念在我的系统里根本不存在。我需要一个让任务真正排队、可以追踪、可以重试、可以分配给不同机器的中间层。这东西有一个名字,叫消息队列。
2. 为什么视频处理这类任务天生“吃队列”:三个刚需与两种伪方案
2.1 视频处理任务的三个硬特征
很多人一聊消息队列就想到双十一秒杀、高并发下单,觉得自己个人项目用不上。但视频处理这类任务,其实天生就适合消息队列,因为它有三个绕不开的特征。
第一,单个任务耗时极长。一个 10 分钟的视频转码,快则几分钟,慢则十几二十分钟。如果调用方是同步等待的,HTTP 请求早就超时了。第二,任务到达速率极不均匀。平时一天没几个文件,一到客户批量交付就是几百个文件疯狂涌进来。这跟秒杀的流量洪峰本质是一回事,只不过秒杀是秒级的峰值,视频处理是小时级的峰值。第三,失败率不低。编码参数不合法、内存不足、输入文件损坏、第三方识别接口超时,各种状况都能让任务失败,而且失败了通常不能直接丢弃,得重试。
这三个特征,决定了你没法用“同步调用 + 简单脚本”来承载。它需要有一个中间层,把“任务的产生”和“任务的处理”彻底隔离开。你可以把消息队列想成饭店里的点单环节:客人不直接进厨房对厨师喊菜,而是把需求写在点菜单上,传到后厨窗口;厨师按顺序做,做完一道端走一道。客人不用站着等厨房炒完,后厨也不会因为前面客人多了就乱作一团。视频处理也一样,上层业务只负责“把需求写下来”,底层 worker 只负责“从队列里取任务去干”。两边互不关心对方怎么运作。
2.2 数据库表当任务队列:省事但坑多
在没有真正消息队列前,很多人会走一条很自然的弯路:用数据库表来模拟任务队列。我也试过。建一张task表,字段是主键、任务类型、参数、状态、重试次数,然后多个 worker 定时去 SELECT 状态为 pending 的记录,处理完再 UPDATE 成 done。
这个方案有几个致命问题。第一,多个 worker 同时 SELECT 会拿到同一条记录,你得靠SELECT ... FOR UPDATE SKIP LOCKED来加锁,但在 2G 内存的破机器上,数据库锁一多,查询性能立刻下降。第二,worker 处理到一半崩溃了怎么办?这条记录还是 pending,但任务实际已经半执行状态,重跑会重复,不重跑会丢失。你得自己发明一个“处理中”状态和“超时重新入队”的逻辑,相当于手搓一个半成品消息队列。第三,数据库表轮询本身是低效的,你没法像消息队列那样做长轮询阻塞待命,只能每秒钟扫一次表,垃圾佬机器上的磁盘 IO 可经不起这么折腾。
2.3 Redis List 离“真队列”还差什么
后来我升级了一点,用 Redis List 当队列,LPUSH生产,BRPOP消费,靠一个常驻进程去处理。这一下子解决了“数据库锁竞争”和“轮询开销”的问题,也比数据库表舒服很多。但用着用着我又发现,Redis List 只解决了“排队”本身,没有解决“可靠性”。
BRPOP把消息从列表里弹出来,如果消费者进程当场崩了,这条消息就永远消失了——因为列表里已经没有了。更麻烦的是,你没法做“消息确认”,没法知道哪条消息被谁消费到一半,没法做重试,也没法设置消息的延迟处理。当时我还自己加了一个 backup list,处理前先把消息复制过去,处理完再删掉,逻辑非常别扭。直到后来我认真研究了 Redis Stream,才意识到自己一直在用“假队列”,而真正能支撑业务闭环的“真队列”,至少要有 ACK 确认、消息重投递、死信隔离这几个基本能力。
3. 破烂机器上的架构成型:选型、部署与量化验证
3.1 为什么最终选了 Redis Stream,而不是 RabbitMQ / Kafka
既然决定要上真消息队列,接下来就是选型。我当时认真对比了三个主流方案:RabbitMQ、Kafka、Redis Stream。
| 维度 | RabbitMQ | Kafka | Redis Stream |
|---|---|---|---|
| 功能完整度 | Exchange/路由/死信/延迟队列全都有 | 高吞吐、日志存储强,但功能偏底层 | 消费组、ACK、Pending 列表都有,延迟队列要自己处理 |
| 资源占用 | Erlang VM 在 1G 内存机器上偏吃紧,管理插件更重 | 一个 broker 纯启动就占用不小,磁盘日志更吃空间 | 本身占用极低,复用已有 Redis 即可 |
| 运维复杂度 | 需要单独维护一个服务,配置较多 | 至少 3 节点才能玩得舒服,垃圾佬的机器根本跑不动 | 命令简单,几乎零额外运维 |
| 适用场景 | 复杂路由、企业级集成 | 大数据量、跨系统长期存储 | 个人项目、中等规模异步任务、低配机器 |
有人会说,RabbitMQ 才像“正经消息队列”,Redis Stream 是不是有点寒酸?我的回答是:对我这台 4 核 8G + 2G 笔记本的环境来说,RabbitMQ 的运维成本和内存开销已经成为一个负担。Kafka 就更不用提了,3 个节点在机械硬盘上起步,我可能还没等任务处理完,先被磁盘 IO 卡死了。Redis Stream 基于 Redis 5.0 引入,它具备消费组、ACK、Pending Entries List 这些核心能力,对我每天几千条任务的量级来说,绰绰有余。
当然,如果你的业务将来要跨多个团队、要复杂的消息路由、要支撑更高的吞吐,那 RabbitMQ 或 Kafka 是值得考虑的。但对我们这种“垃圾佬”来说,先让系统跑起来、跑得稳,比追逐“主流”重要得多。
3.2 任务消息、消费组与 ACK 的最小闭环
选型定了之后,我开始把“视频处理任务”建模成消息。先跑一个 Redis 实例,设置好内存限制和持久化策略,然后在上面创建 Stream 和消费组。这里有一个很小的闭环,我贴出来给大家参考。
一个典型的消息,我用 JSON 来承载。包含任务 ID、任务类型、输入文件路径、参数哈希、回调地址等。任务 ID 不是文件名,而是视频文件的唯一标识,后面讲幂等时会提到它有多关键。
生产端往 stream 里丢消息:
XADD video_tasks MAXLEN ~ 10000 * task_id "a1f9" type "transcode" input "/data/in/a.mp4" params "{}"创建消费组,从尚未消费的消息开始读取:
XGROUP CREATE video_tasks video_worker 0消费端循环,阻塞读取新任务,处理完执行 XACK:
import redis import json r = redis.Redis(host='localhost', port=6379, decode_responses=True) while True: # 阻塞 5 秒读取一条消息 results = r.xreadgroup( groupname='video_worker', consumername='worker-1', streams={'video_tasks': '>'}, count=1, block=5000 ) if not results: continue stream, messages = results[0] for msg_id, fields in messages: try: process(fields) # 转码、抽帧、识别等 r.xack('video_tasks', 'video_worker', msg_id) except Exception as e: # 记录失败次数,超过阈值转死信 handle_failed(msg_id, fields, e)这段代码最核心的两点:xreadgroup里的>表示读取“从未投递给任何消费者”的新消息;xack则是处理成功后告诉 Redis“这条消息我搞定了,可以从 Pending 里清掉了”。如果 worker 在process()中途崩溃,没有调用xack,消息就会一直留在 Pending Entries List 里。等它重新上线,或者被其他消费者认领,就可以继续处理。
Redis 的持久化我也做了取舍。垃圾佬机器没有 UPS,断电是常有的事,所以我开了appendonly yes,并且设置appendfsync everysec,这样最多丢一秒的消息。对视频处理来说,任务丢了可以重新入队,但任务状态和结果表不能胡写。内存上限我设了maxmemory 512mb,避免 Redis 把整个机器的内存吃光。
3.3 用 200 个任务压出来的真实数据
架构改完,我心里其实也没底。于是做了一个简单的压测:一次性往队列里塞 200 个真实转码任务,三台机器同时跑 worker,每台机器 2 个消费进程(因为你得给转码进程留内存,旧笔记本跑 3 个并发很容易 OOM)。
我当时最担心的不是“能不能跑完”,而是“会不会比原来更慢”。毕竟原来虽然会 OOM,但至少 8 个进程同时在跑,现在总共才 6 个 worker,数量还少了。跑完之后我把数据列出来对比了一下,结果很有意思。
| 指标 | 旧脚本方案 | Redis Stream 方案 |
|---|---|---|
| 200 个任务实际耗时 | 约 11 小时(期间 OOM 两次,人工重试 5 次) | 约 3.6 小时(中间无人工干预) |
| 处理期间最高内存 | 台式机 3.2G,接近崩溃 | 每台机器平均 1.2G,稳定 |
| 状态可追踪性 | 靠文件名后缀猜 | 队列长度、消费者状态、Pending 条数随时可查 |
| 失败重试 | 手动处理 | 自动 Pending 重投递,失败次数记录 |
有人会说,为什么 worker 少了,反而更快?关键原因是“任务不落地、无状态了”。以前脚本崩了之后,后面的任务全堵在文件系统里,日志乱成一团,每次恢复都要人工判断进度;现在每个任务从入队到结束都有明确状态,worker 崩了消息还在 Pending 里,换个消费者立刻续跑。整个系统的有效工作时间大幅提高,反而是真正的提速。
当然我也得说实话:如果瓶颈是 CPU 本身,消息队列不可能让 6 个 worker 比 60 个 worker 还快。它解决的核心问题不是“计算变快”,而是“让计算不被打断、让资源不被浪费、让协作变成可能”。这一点,在破机器上尤其明显。
4. 微服务不是口号:消息队列如何逼我完成第一次服务化拆分
4.1 先弄清楚:微服务微的不是框架,而是边界
一提到微服务,很多人想到的是 Spring Cloud、注册中心、网关、配置中心那一大套。但这其实是工具层面的东西。微服务真正的含义,是把一个又大又笼统的“系统”,按照业务边界拆成若干个可以独立部署、独立演进、独立扩展的服务。
我最初的代码就是一个大 Python 脚本,把所有功能混在一起:扫描目录、转码、抽帧、目标检测、写数据库、发通知。看起来是一个文件搞定所有事,但每次想改识别算法,要担心会不会影响转码;每次想加一台机器,要在三种日志里来回找。这些耦合问题,本质上跟代码在哪里没关系,时间长了,你根本分不清哪个模块该为哪个故障负责。
4.2 队列作为服务间契约:替代 HTTP 长连接
在引入 Redis Stream 之后,我发现一个自然的副产品:我不得不把整个处理流程拆成三段,因为它们对“生命周期”的要求完全不同。
第一段是生产者,它只负责把“待处理任务”写成一条消息,发完就结束,完全不关心后续转码要跑多久。第二段是消费者,它常驻运行,从队列里取消息,执行真正的重计算。第三段是结果收尾,负责把处理结果写入数据库、发通知、清理临时文件。
这个拆分不是我想“微服务化”才做的,而是队列逼出来的:生产者不需要知道自己发出的消息会被哪个进程处理;消费者不需要知道任务是从哪里来的;收尾服务也不需要知道处理过程到底经历了多少次重试。三者之间唯一的契约就是“消息格式”。我后来把“消息格式”定义为一张 JSON Schema,这个 Schema 就是服务间协议的雏形。
如果当初用 HTTP 同步调用来做这件事,会怎么样?上游把一个转码请求 POST 给下游,下游得在几十分钟里保持连接不超时,上游还得维护一堆没完成的 HTTP 调用状态;下游过载了,上游得自己写重试逻辑。这些本来该由中间件解决的问题,会全变成业务代码里的痛苦。而换了队列之后,上游只要保证“消息能发出去”,下游只要保证“消息能处理完”,两边都轻松。
这一段经验让我彻底理解了一个道理:做微服务,不要在 Spring Cloud 那一套东西里陷得太深,先把服务之间的“通信边界”想清楚。用队列做异步边界,往往是个人项目里最低成本、最实用的微服务化方式。你甚至不需要 RPC 框架,只需要一个可靠的队列和一个稳定的消息格式。
4.3 两次拆分才看懂的服务边界
拆服务这件事,光靠想是想不明白的,得靠踩。
我第一次拆分,是按“代码功能”拆的。把视频处理脚本拆成了“转码服务”和“识别服务”,但两者还是强耦合,因为转码完马上要调用识别,识别失败转码结果也没意义。这个拆分实际上只是把文件拆开,并没有形成独立的边界,运维复杂度反而上去了。
第二次拆分,才真正按“生命周期”拆。转码和基础抽帧属于“重计算层”,需要一堆常驻 worker 去消费;识别结果属于“数据沉淀层”,需要另一个消费者去做落库和通知;而“任务来源”则可以独立扩展,随时接受新需求。这样一来,三个服务之间完全是靠队列通信的,互不阻塞,任何一个服务挂了,另外两个还能暂时照常运作。
我后来画过一张很简陋的图,不值得放上来,但脑子里很清楚:生产源(采集目录/API 接口)→ 任务队列(Redis Stream)→ 重计算 Worker 组 → 结果队列 → 落库与通知服务。整条链路里唯一的“记忆”就是队列和数据库,每个服务自己都是无状态的,可以随时加机器扩展。
5. 重复消费问题:第一次让我怀疑人生的 Bug 与最终的幂等设计
5.1 线上事故:同一个视频为什么跑了两遍
讲完架构,必须单独聊一个绕不开的问题——消息队列重复消费。
事情发生在架构切换后第二周。某天早上我查数据库,发现同一条视频记录出现了两条结果,文件系统里也多了两个一模一样的抽帧目录。第一反应是代码有 Bug,赶紧去看日志,顺着消息 ID 一路追,最后定位到了原因:
消费者 A 从队列里取到任务,开始转码,转到一半,进程因为内存不足退出了。消息已经投递给了 A,但 A 没有来得及执行XACK,所以消息一直躺在 Pending Entries List 里。消费者 B 上线后,通过 Redis 的认领机制把这条 Pending 消息重新取走,重新执行了一遍完整的转码流程。于是,同一个视频被处理了两次。
这个现象在消息队列里有一个专业术语,叫“至少一次投递”(At-Least-Once Delivery)。意思是,消息队列保证每条消息至少会被消费一次,但不保证不会重复。这不是队列的缺陷,而是所有消息队列的默认语义。RabbitMQ 有,Kafka 也有,Redis Stream 同样有。既然队列本身就是这个语义,那业务上就必须自己处理“如果同一件事被执行了两次,结果仍然正确”。这就是幂等设计。
5.2 幂等设计的三个层次:接口、存储、文件系统
吃了一次亏之后,我给自己定了一个规矩:所有从队列消费的任务,都必须做到“可重入”。也就是,就算同一任务被执行两次,最终效果跟执行一次完全一致。具体到视频处理这个场景,我做了三个层面的处理。
第一层,是接口层的检查。消费者在真正开始干活前,先根据任务 ID(视频文件的唯一哈希)去查一下结果表。如果发现这个任务已经是 done 状态,直接XACK确认掉,不再重复处理。这个检查非常便宜,但能挡掉大部分重复消费。
第二层,是存储层的唯一约束。结果表里,任务 ID 字段必须加唯一索引。插入结果时用INSERT ... ON CONFLICT DO NOTHING,就算两个 worker 同时在结果表里插入同一条记录,数据库也会保证只有一条能成功。这层兜底的作用很大,因为不管代码怎么漏,数据库的唯一性约束是最后一道防线。
第三层,是文件系统层面的隔离。每个任务的处理结果都放在以任务 ID 命名的目录里,比如/output/a1f9/。如果任务重复执行,转码输出会直接覆盖到同一个目录,而不是新建一个a1f9-copy目录。这样文件系统层面也不会因为重复执行而膨胀出垃圾文件。
5.3 落实在 Redis Stream 上的实现细节
把幂等落实到 Redis Stream 上,有几个具体细节值得说一说。
关于 Pending 的认领。消费者崩溃后,消息会留在队列的 Pending 里。如果你用的是XCLAIM或XAUTOCLAIM去认领这些消息,一定要给每条消息设置一个合理的“最小空闲时间”(min-idle-time)。比如一个转码任务起码要跑 3 分钟,你就不能在它刚崩溃 5 秒后立刻让另一个 worker 接管,否则两个 worker 可能同时在跑同一个任务,幂等就没意义了。我把这个值设成了 600 秒,宁可让任务晚一点重跑,也不要让它并发跑。
关于 XAUTOCLAIM 的注意点。Redis 6.2 以后建议用XAUTOCLAIM替代XCLAIM,它的 IDLE 管理更合理。同时,认领消息时不要只认领一条就完事,因为它会返回一批消息,你需要循环处理完再发XACK。
关于死信队列。我给消费者加了失败次数重试机制:每个消息在 Redis 里维护一个failed:task_id的计数器;达到 5 次仍失败,就把它另写到一个dead_tasksstream 里,同时把原消息XACK掉,防止它继续无限重投递卡死整个队列。死信不是终点,我每天会扫一眼dead_tasks的条数,用来发现那些“一直处理不过去的脏任务”——比如一个损坏的视频文件,或是一个永远不合法的编码参数。
关于监控。Redis Stream 没有现成的 UI,但几个命令足够用了。XLEN video_tasks查看积压消息数;XPENDING video_tasks video_worker查看有多少条消息在 Pending 里卡着;XINFO GROUPS video_tasks可以看每个消费组的速度和积压差距。我把这些命令包成了一个监控脚本,每 5 分钟跑一次,把关键指标打到日志里。队列长度超过 500 或者 Pending 数量超过 20 的时候,我就知道该加 worker 或者该检查是不是有任务卡死了。
关于延迟队列。视频处理里偶尔会需要“等 10 分钟后检查某文件是否存在”这类逻辑,Redis Stream 本身不支持延迟消息,我的做法是另建一个delay_tasksstream,里面放好消息要投递的时间戳,由一个轻量消费者扫描,到了时间再把它 XADD 到video_tasks里。这种“自己搓延迟队列”的办法对个人项目完全够用。如果你用 RabbitMQ,可以用内置的 TTL + 死信路由实现更优雅的延迟队列,但在 Redis Stream 里,先接受这个朴素方案也无妨。
写在最后的一点经验
如果你也想把手头的旧电脑变成一台讲道理的“微型分布式系统”,我的建议是:别一上来就学 Spring Cloud、Nacos、Sentinel 那一整套,那是已经明确知道自己要微服务化之后的事。一个真正能让你成长的分水岭,是先把你那个“轮询目录”的脚本,改成 XADD / XREADGROUP / XACK 的闭环。改完之后你会发现,队列会逼你重新思考很多问题:任务之间怎么解耦、状态怎么可追踪、崩溃之后怎么恢复、同一个任务能不能被安全地执行多次。
这些东西想清楚了,你再去看任何微服务框架,都会觉得亲切很多。消息队列不是一个需要仰望的技术符号,它是你系统里最朴素的“分工与协作”的显性表达。我用破烂电脑悟到的,不是某个框架的配置,而是这个。