JIMMY
返回项目列表

Forge

分布式任务队列
GoRedis Stream分布式

单队列 8k 任务/秒 · 延迟任务误差 <100ms · 死信可回放

forge — queue
 

背景与动机

很多活不该在请求线程里干:发邮件、跑构建、生成报表。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,主请求路径轻装上阵。