ARTICLE DETAIL

资讯详情

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

Madeira:异步任务调度与事件回执分发服务的设计与实践

Madeira:异步任务调度与事件回执分发服务的设计与实践 当我第一次把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导致一批在途任务失去了来源上下文排查了很久。从那时起我规定所有已经进入终态的事件都要落 PostgreSQLRedis 只承担“短时流转”和“消费结算”不承担“永久存储”。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( groupnamemadeira-workers, consumernameworker-1, streams{madeira:delivery: }, count10, block2000, ) for stream, messages in resp or []: for msg_id, fields in messages: data json.loads(fields[bdata]) 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[bdata]})这段代码有一个很容易错的地方失败后必须先把当前消息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一套跑 WorkerRedis 单独跑在云厂商的托管实例上。Compose 文件的核心部分长这样services: madeira-api: build: . command: uvicorn app.main:app --host 0.0.0.0 --port 8000 --workers 4 environment: - REDIS_URLredis://${REDIS_HOST}:6379/0 - DATABASE_URLpostgresql://${DB_USER}:${DB_PASSWORD}${DB_HOST}:5432/madeira ports: - 8000:8000 restart: unless-stopped madeira-worker: build: . command: python app/worker.py environment: - REDIS_URLredis://${REDIS_HOST}:6379/0 - DATABASE_URLpostgresql://${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再仔细看响应体内容里面是successfalse的业务错误。原先我的出站判断只看 HTTP 状态码没有解析响应体。后来我把成功判断改成“HTTP 状态码为 2xx 且响应体可解析、无业务错误码”并为这种“半成功”场景单独记录日志。这个坑再次验证了HTTP 层的成功和业务层的成功是两回事。第三个故障是消息重复。下游收到的重复请求变多排查后发现是 Worker 在处理消息时阻塞时间太长XACK没有发出去消息留在 Pending 列表被另一个新 Worker 用XAUTOCLAIM领走于是同一个事件被处理了两次。这个问题的根本原因是下游响应时间不可控Worker 的等待时长超过了我预设的确认窗口。对策是增加“消费中”标记出站前先在 Redis 里写一个带 TTL 的投递中标记处理完成后再检查标记是否能删除如果标记已过期但事件还在 Pending 里说明上一个消费实例可能已经失联这时才允许重新领取。虽然不能做到 100% 杜绝重复但能把业务上的重复率降到可以接受的程度。5. 给想抄作业的人三句话项目做完之后我从这些经验里提炼出三条比较实用的点放在最后聊。第一异步中间层一定要把“可观测性”当成第一需求来设计而不是后补的功能。事件从入站开始就带上trace_id日志带、落库带、出站请求头也带指标从第一天就接好。这能让你在线上出问题时用最短时间定位省下的运维成本远超搭建监控的那点工作量。第二凡是涉及跨系统写入的地方都要先想好“重复了怎么办”。哪怕上游拍着胸脯说不会重复推送也要在接入层做幂等校验。幂等不是业务方的需求是系统的自我保护代价通常只是几个字节的 Redis key收益却是避免一次灾难性的重复投递。第三设计重试的时候永远要站在下游的角度想问题。下游接口能力有限我们就应该用指数退避加抖动避免重试风暴下游接口未恢复我们就应该把流量收进延迟队列而不是一遍遍冲击。真正的可靠不是“永远成功”而是在失败发生时仍然有秩序、可追踪、可恢复。如果你也在搭类似的事件分发或者任务调度系统希望这篇记录能帮你少踩几个坑。我回看这套东西最大的感受是复杂的分布式技术不一定酷但把一件普通事情做正确、做细致往往才是线上体验最值得投入的部分。
返回列表