Forge
分布式任务队列 GoRedis Stream分布式
单队列 8k 任务/秒 · 延迟任务误差 <100ms · 死信可回放
背景与动机
很多活不该在请求线程里干:发邮件、跑构建、生成报表。Forge 是一个基于 Redis Stream 的任务队列——提交方只管扔,消费方各取所需,延迟任务、指数退避、死信回放一应俱全。
为什么不用现成的 Celery?它绑定 Python 生态,而我的消费方有 Go 有 TS。Forge 的协议语言无关:消息体是 JSON,谁都能消费。
架构与取舍
Redis Stream 天生带消费者组(consumer group)语义:待处理列表(PEL)天然解决“任务领取后消费者崩溃”的归还问题,XAUTOLOAD 一调,任务自动易主。
延迟任务不用轮询数据库:提交时算好执行时间戳,写入按时间分桶的 ZSET,调度协程每 100ms 扫一次到期桶,误差控制在百毫秒级——对构建类任务完全够用。
难点与解决
重试退避的难点在公平:一个一直失败的任务不能堵死队列。Forge 的退避时间从 2s 指数涨到 15 分钟封顶,并且重试任务和新鲜任务分 Stream 存放,互不占道。
死信不是终点:DLQ 里的任务带完整失败上下文,一条命令即可批量回放——修复 bug 之后,把损失找回来。
数据与结果
单队列实测 8k 任务/秒吞吐,300 并发消费者下消息零丢失(PEL 机制保底)。它现在支撑着 Aegis 的异步用量结算——网关记账走 Forge,主请求路径轻装上阵。