当我第一次把Madeira这个词填进仓库目录名的时候,只是觉得它读起来很有辨识度。后来有人在 PR 评论区里追问:这个项目为什么叫 Madeira?我想了想,给出的答案倒也简单:马德拉酒要经历漫长的熟成过程才有味道,异步任务的链路也是如此,速度快不算本事,重点是每一步都能被追踪、能被重试、能被信任。
Madeira是一套异步任务调度与事件回执分发服务,主要部署在内网,解决一个非常具体的实际问题:业务系统把事件或任务交给外部回调地址,但外部接口不一定时刻可用,网络会抖动,接口会超时,数据会被重复推送。如果上下游每次都走同步调用,任何一端的抖动都会顺着链路传导放大;如果只是把任务丢进消息队列就不管后续,又会在两个系统之间留下数据黑洞。Madeira 就是夹在上下游之间的一层“异步缓冲带”,把任务的接收、路由、投递、重试、死信全部管起来。
项目做完之后回看,最值得说的不是某个复杂算法,而是那些藏在细节里的取舍:幂等键放在哪一层、延迟重试怎么设计、监控里到底该看哪些指标、死信如何人工干预。这些内容我会在这篇里全部拆开讲,也把当时踩过的坑一并记录下来,希望对正在搭类似中间层的朋友有帮助。
1. 项目背后:为什么是 Madeira,它在解决什么问题
1.1 命名由来与项目定位
起初,“Madeira”只是我临时起的代号。我习惯给工程项目取一个有辨识度的名字,方便在监控面板和日志前缀里一眼认出来。没想到这个名字一直留了下来,后来成了这套异步任务调度服务的正式代号。
准确地说,Madeira 是一个部署在内网的异步任务调度与事件回执分发服务。它的日常工作可以概括成三件事:第一,接收上游业务系统发来的任务事件;第二,把事件按照既定规则投递给下游回调地址;第三,在投递失败时按照策略重试、隔离,并把最终失败的事件送到死信队列等待人工处理。整个过程全部异步化,上游系统发完消息就可以继续干自己的事,不需要暂停在那里等回调结果。
这事听起来不太起眼,但解决过的麻烦非常真实。我踩过最深的坑是:两个系统之间采用同步 HTTP 调用,上游产出数据后直接请求下游接口,下游偶尔抖动就会出现请求超时。上游为了确保成功,又加了一层“超时重试”,结果下游恢复后瞬间收到爆量请求,直接把服务打挂。后来我意识到,同步调用链条越长,系统的短板效应越明显,任何一个环节的单点故障都会顺着调用链传染。与其在两边反复调超时和重试参数,不如在中间插入一个异步缓冲层,让上下游在时间上解耦。Madeira 就是这个想法的落地形态。
1.2 核心痛点与使用场景
Madeira 要解决的核心问题,可以拆成四个词:解耦、削峰、可靠、可观测。
解耦指的是把“产生事件”和“处理事件”拆开。上游只负责把事件丢给 Madeira,至于下游是谁、有几个、什么时候处理完成,上游并不需要关心。削峰指的是当上游短时间产生大量事件时,Madeira 可以把流量挂在队列里慢慢消费,避免下游被瞬时高峰冲垮。可靠指的是所有事件最终都要有结果,要么投递成功,要么进入死信队列并明确记录失败原因,不允许事件凭空消失。可观测指的是每条事件从进入到完成的全过程都有轨迹,团队可以随时查询一条事件当前处于什么状态、卡在哪个环节、被重试了几次。
实战中我用得最多的场景有两个。一个是回调分发:比如用户在业务系统里发起一个异步任务,任务完成后需要通知到外部平台,外部平台要求必须带签名、必须幂等、不能重复通知。另一个是批量补偿:比如凌晨跑批出来的账务数据要批量推送给合作方,合作方接口能力有限,每秒只能接受几十个请求,Madeira 就在队列消费时做均匀限速,把几万条数据以可控速率推送完。这两个场景听起来普通,但细节特别多,后面几节我再逐个拆开讲。
2. 设计选型:一条事件要走多久才能被信任
2.1 技术栈盘点与选型理由
Madeira 的核心技术栈不算复杂:Python 写业务、Redis 做队列与缓存、PostgreSQL 做最终存储、Docker Compose 负责本地和单机部署。我在重写版本里选了 FastAPI 框架,主要原因是团队当时对 Python 更熟,开发速度更快。如果完全从纯性能出发,Go 会更折腾一些,但 Madeira 的瓶颈通常不在计算上,而在网络调用和下游接口的响应速度上,所以 FastAPI 的异步接口完全够用。
队列选型上我最初考虑过 RabbitMQ,后来还是用了 Redis Streams。原因有几个:Redis 在团队里已经是现成的基础设施,不需要额外维护一套 Erlang 生态;Redis Streams 支持消费组、消息确认、Pending 列表,能力上恰好覆盖需求;而且用 Redis 还能顺便承接幂等标记、限流计数器这些附带功能,减少组件数量。只有一点要提醒:Redis 必须开启 AOF 持久化,并且最好用单独实例,别和业务缓存混用同一份内存,否则消息丢失和缓存淘汰会互相干扰。
PostgreSQL 用来落地业务数据和审计日志。一开始我甚至觉得审计日志可以省掉,直接信 Redis 就能查到所有事件。后来一次意外改变了这个想法:有人误操作清掉了 Redis 的部分 key,导致一批在途任务失去了来源上下文,排查了很久。从那时起,我规定所有已经进入终态的事件都要落 PostgreSQL,Redis 只承担“短时流转”和“消费结算”,不承担“永久存储”。
2.2 关键链路设计:入站、路由、出站
整体链路我把它分成三段:入站、路由、出站。
入站链路负责暴露 HTTP 接口并接收事件。上游系统调用 Madeira 的 API,把事件内容、回调地址、自定义参数一起送过来。接入层要做三件事:签名校验、格式校验、幂等校验。签名校验保证请求确实来自可信的上游;格式校验保证必填字段齐全;幂等校验保证同一个事件 ID 只进队列一次。三层校验都通过,事件才会被写入 Redis Stream,然后立刻返回给上游一个“已受理”的回执。
路由链路是中间的核心。消费者从 Redis Stream 里拉出事件,根据事件里的路由键查出一组目标规则。这个阶段不真正请求下游,只负责解析、补全模板、决定走哪条投递通道,然后把“待投递”格式的消息写进另一组 Stream。这样做有一个好处:入站和出站完全分离,入站高峰不会直接压给出站,出站改造也不会影响入站写入。
出站链路负责真正的 HTTP 调用。出站消费者根据目标地址发送请求,等待响应,判断状态码。如果成功,就把事件标记为完成;如果失败,按照配置的退避策略安排重试;超过最大重试次数后,事件进入死信队列。出站链路的节奏可以单独配置,比如每个消费者只允许并发 10 个请求,或者平均每秒不超过 30 个,好配合下游接口的容量。整条链路由多个独立消费者组成,每个消费者只做一件事,出了问题也只需要重启对应的进程,不会拖垮全局。
2.3 数据模型与消息格式
消息格式我最终收敛成了一套统一的事件信封。每个事件在传输层都有以下字段:
event_id:全局唯一,通常由上游生成,也是幂等键。event_type:事件类型,比如task.completed、payment.synced。source:来源系统标识,签名校验时用到。target_uri:投递目标,出站时拼接最终 URL。payload:业务数据,统一用 JSON 表示。created_at、updated_at:时间戳,用于追踪和超时判断。trace_id:链路追踪 ID,方便把 Madeira 内部日志和上下游日志串起来。retry_count:已重试次数。status:枚举值,包括pending、delivering、succeeded、dead、cancelled。
数据库表我设计得很克制,两张表就够用。event表保存事件的完整快照和当前状态,delivery_log表记录每一次投递尝试的时间、状态码、返回摘要、耗时。查询一条事件时,先用event表看总状态,再用delivery_log表看它经历过几次尝试、每次发生了什么。这个设计看起来很平铺直叙,但排查问题的时候特别好用,因为所有信息都是按时间线排列的。
Redis 里只放三类数据:Stream 里流动的待消费事件、幂等标识、限流计数。终态数据一律不依赖 Redis,这是后来线上事故逼出来的铁律。
3. 核心实现:把异步任务做正确
3.1 签名校验与幂等控制
签名这块我踩过一次很尴尬的坑:上游服务端用 SHA256 对整个 JSON 串签名,转发时不小心调整了 JSON 字段顺序,结果签名一直对不上。从那以后我统一约定,签名内容是“时间戳 + 换行 + 原始请求体字节”,参与签名的请求体用原始字节,不允许序列化两次。下面是示例代码:
import hashlib import hmac import time def sign(secret: str, payload: bytes) -> tuple[str, str]: timestamp = str(int(time.time())) msg = (timestamp + "\n").encode("utf-8") + payload digest = hmac.new(secret.encode("utf-8"), msg, hashlib.sha256).hexdigest() return timestamp, digest def verify(secret: str, payload: bytes, timestamp: str, digest: str) -> bool: msg = (timestamp + "\n").encode("utf-8") + payload expected = hmac.new(secret.encode("utf-8"), msg, hashlib.sha256).hexdigest() return hmac.compare_digest(expected, digest)时间戳本身的宽限窗口我设在 300 秒。太短容易误伤不同机器间的时钟偏差,太长又会给重放攻击留下空间。校验通过后,下一步就是幂等。
幂等控制我用 Redis 的SET NX实现。同一个event_id第二次进来时直接返回“重复”,不重复入队。这里有个细节:入队和幂等标记必须放在同一个事务边界里。我的做法是先用 Lua 脚本把幂等标记和XADD一起执行,脚本在 Redis 侧保证原子性;如果中间崩溃,要么都成功,要么都失败,不会出现“标记留下了但消息没进去”的脏状态。
-- KEYS[1]: idempotent key -- KEYS[2]: stream key -- ARGV[1]: event_id -- ARGV[2]: message payload local ok = redis.call("SET", KEYS[1], "1", "NX", "EX", "86400") if not ok then return 0 end redis.call("XADD", KEYS[2], "*", "data", ARGV[2]) return 13.2 延迟队列、重试与退避策略
异步系统里,重试策略设计得不好比不设计还危险。最典型的问题是固定间隔重试:下游恢复瞬间,几百个定时重试请求同时涌进去,再次把下游打崩,形成恶性循环。我给 Madeira 配置了指数退避加抖动。初始延迟 1 分钟,每次翻倍,上限 30 分钟,同时加 20% 以内的随机抖动,避免重试请求扎堆。
延迟队列用 Redis Stream 做比较别扭,因为 Stream 本身不支持“到时间才能消费”。我最后采用的办法是:为每种延迟级别建一个独立的 Stream,例如madeira:delay:60、madeira:delay:300。消费者只负责把到期的消息转移到主工作流 Stream,转移前检查一下当前时间是否已经超过消息里记录的due_at。如果没到,就稍等再查,或者重新把消息放回同一个延迟 Stream。这个实现不复杂,但能很好地缓解“所有任务都挤在一个队列里等待”的问题。
消费者代码的核心是消费组机制。每次消费一批消息,处理完用XACK确认,中间如果消费者宕掉,消息会留在 Pending 列表里。恢复后用XAUTOCLAIM把超时未确认的消息重新领走:
import redis import json r = redis.Redis.from_url("redis://127.0.0.1:6379") def consume_once(): resp = r.xreadgroup( groupname="madeira-workers", consumername="worker-1", streams={"madeira:delivery": ">"}, count=10, block=2000, ) for stream, messages in resp or []: for msg_id, fields in messages: data = json.loads(fields[b"data"]) ok = send_webhook(data) if ok: r.xack("madeira:delivery", "madeira-workers", msg_id) else: r.xack("madeira:delivery", "madeira-workers", msg_id) r.xadd("madeira:delay:300", {"data": fields[b"data"]})这段代码有一个很容易错的地方:失败后必须先把当前消息XACK掉,再把它放进延迟 Stream,否则当前消息会被同一个消费者反复拿到,产生“卡死循环”。很多人忽略这一点,结果是任务表面上看在重试,实际上同一批消费根本推不进去。
3.3 出站推送与失败补偿
出站部分最重要的经验是:别把“HTTP 状态码 200”当成成功的唯一标准。有些下游接口即使返回 200,响应体里也可能带着业务失败码;有些接口在 5xx 和超时时表现完全不一样。我在出站判断里把结果分成了四类:
| 结果分类 | 判定条件 | 处理动作 |
|---|---|---|
| 成功 | HTTP 2xx 且响应体可解析、无业务错误码 | 标记完成,记录日志 |
| 可重试失败 | 5xx、超时、连接被重置 | 进入延迟队列,按退避策略重试 |
| 不可重试失败 | 4xx 参数错误、鉴权失败 | 进入死信队列,并保留失败原因 |
| 未知失败 | 响应体解析失败或状态码异常 | 多试几次,直到达到上限 |
失败补偿依赖之前说的延迟再投递机制。出站消费者失败后,把消息连同本次失败原因写到delivery_log,然后投递到延迟 Stream;延迟 Stream 的消费者到期后,把消息重新丢回出站 Stream。每次重试都会更新消息头里的retry_count,达到上限后进入deadStream。
死信队列的人工处理我也做了产品化:提供一个管理后台,可以查看死信详情、手动重放、批量取消。手动重放的场景是第三方接口修复后,把积压的死信按原顺序重新推送;批量取消的场景是确认某批数据已经没人要了,直接终结掉,避免它们一直占用后续资源。这个“能人工干预”的设计,比全自动补偿更能兜底,也是项目上线后运维同事最依赖的功能之一。
4. 部署与排错:从脚本到可观测
4.1 单机部署与容器编排
服务刚成型时我用 systemd 加裸 Python 进程跑,部署一次要手工 pull 代码、装依赖、重启,很累。后来切到 Docker Compose,整个组件的启动命令收敛成一个文件,本地开发体验好了几个量级。线上我也用 Compose,一套跑 API,一套跑 Worker,Redis 单独跑在云厂商的托管实例上。
Compose 文件的核心部分长这样:
services: madeira-api: build: . command: uvicorn app.main:app --host 0.0.0.0 --port 8000 --workers 4 environment: - REDIS_URL=redis://${REDIS_HOST}:6379/0 - DATABASE_URL=postgresql://${DB_USER}:${DB_PASSWORD}@${DB_HOST}:5432/madeira ports: - "8000:8000" restart: unless-stopped madeira-worker: build: . command: python app/worker.py environment: - REDIS_URL=redis://${REDIS_HOST}:6379/0 - DATABASE_URL=postgresql://${DB_USER}:${DB_PASSWORD}@${DB_HOST}:5432/madeira restart: unless-stopped有一点我要特别强调:API 和 Worker 必须分开部署,不能把消费循环放进 API 进程里。我曾经图省事,在 FastAPI 启动事件里挂了一个后台消费线程,结果线上接入量一上来,消费线程把事件循环占满,HTTP 接口的响应时间从几十毫秒涨到十几秒,整条链路直接受影响。现在 API 只负责接收入站请求、写队列;消费和出站全部放在单独的 Worker 进程里,相互之间只通过 Redis 通信。进程多了就水平扩 Worker 数量,架构简单清晰。
4.2 监控指标与日志设计
异步系统最怕的就是“看起来没报错,但任务就是没结果”。所以我从第一天起就规定,Madeira 必须暴露一套可以横向对比的指标,至少包含以下几项:
- 入站 QPS:单位时间收到的请求数量。
- 入站失败率:签名校验失败和格式校验失败的比例。
- 队列积压:当前 Stream 中未消费的消息数量。
- 消费速率:单位时间成功从队列中取走并处理的消息数量。
- 出站成功率:投递成功的消息占总投递次数的比例。
- 重试分布:重试 1 次、2 次、5 次以上的消息各占多少。
- 死信增长速度:单位时间进入死信队列的消息数量。
这些指标会输出为 Prometheus 标准格式,用 Grafana 展示。其中“重试分布”特别重要,它是判断下游是否健康的重要信号。如果某个下游的重试次数开始集中上涨,说明它已经接近瓶颈,运维应该提前通知对方扩容,而不是等到 5xx 一片的时候才被动响应。
日志方面,所有模块共享同一个trace_id。入站时生成trace_id,写进消息头,出站时再把trace_id带在 HTTP 请求的X-Trace-Id上。这样一旦下游投诉说收到了重复请求或者没收到请求,我可以只用trace_id就把整条链路的日志拉出来。排查异步问题最大的障碍是“上下文断裂”,trace_id就是用来缝合断点的。
4.3 现场实录:三个常见故障
第一个故障是 Redis 连接池耗尽。现象是 API 响应偶尔变慢,日志里出现一堆Timeout reading from socket。原因并不复杂:Worker 消费线程里每条消息都新建连接,请求高峰时连接数飙升,把 Redis 连接池塞满。后来我把 Redis 客户端实例改成全局复用,把连接池上限调大,同时给消费循环加了信号量限制并发,问题就消失了。教训是:不要在高频路径里反复创建 Redis 连接,连接一定要复用,而且消费并发必须有显式上限。
第二个故障很有意思。某个下游回调接口一直返回 200,但业务方反馈数据没到。我查了事件轨迹,发现每次请求都拿到了 200,再仔细看响应体内容,里面是success=false的业务错误。原先我的出站判断只看 HTTP 状态码,没有解析响应体。后来我把成功判断改成“HTTP 状态码为 2xx 且响应体可解析、无业务错误码”,并为这种“半成功”场景单独记录日志。这个坑再次验证了:HTTP 层的成功和业务层的成功是两回事。
第三个故障是消息重复。下游收到的重复请求变多,排查后发现是 Worker 在处理消息时阻塞时间太长,XACK没有发出去,消息留在 Pending 列表,被另一个新 Worker 用XAUTOCLAIM领走,于是同一个事件被处理了两次。这个问题的根本原因是下游响应时间不可控,Worker 的等待时长超过了我预设的确认窗口。对策是增加“消费中”标记:出站前先在 Redis 里写一个带 TTL 的投递中标记,处理完成后再检查标记是否能删除;如果标记已过期但事件还在 Pending 里,说明上一个消费实例可能已经失联,这时才允许重新领取。虽然不能做到 100% 杜绝重复,但能把业务上的重复率降到可以接受的程度。
5. 给想抄作业的人三句话
项目做完之后,我从这些经验里提炼出三条比较实用的点,放在最后聊。
第一,异步中间层一定要把“可观测性”当成第一需求来设计,而不是后补的功能。事件从入站开始就带上trace_id,日志带、落库带、出站请求头也带;指标从第一天就接好。这能让你在线上出问题时用最短时间定位,省下的运维成本远超搭建监控的那点工作量。
第二,凡是涉及跨系统写入的地方,都要先想好“重复了怎么办”。哪怕上游拍着胸脯说不会重复推送,也要在接入层做幂等校验。幂等不是业务方的需求,是系统的自我保护,代价通常只是几个字节的 Redis key,收益却是避免一次灾难性的重复投递。
第三,设计重试的时候永远要站在下游的角度想问题。下游接口能力有限,我们就应该用指数退避加抖动,避免重试风暴;下游接口未恢复,我们就应该把流量收进延迟队列,而不是一遍遍冲击。真正的可靠不是“永远成功”,而是在失败发生时仍然有秩序、可追踪、可恢复。
如果你也在搭类似的事件分发或者任务调度系统,希望这篇记录能帮你少踩几个坑。我回看这套东西,最大的感受是:复杂的分布式技术不一定酷,但把一件普通事情做正确、做细致,往往才是线上体验最值得投入的部分。