底层三块,我们只写了一块

一共三块,只有一块是我们自己写的。这是故意的:分布式队列里最容易写错的那一层,我们没有动手。

本页目录

底层实现

投递交给 asynq

入队、延时投递、指数退避、把半路挂掉的 worker 手里的活收回来,这些全都来自 asynq,很多 Go 项目已经在生产里跑着它。它的投递语义是至少一次,同一个事件可能到两遍。

状态是每次调用两个 Redis 哈希

<callId>_status 每个步骤一个字段,<callId>_meta 放你传进来的 metadata、Memo 的结果和已投递的信号。至少一次能安全,靠的就是这份状态:重复投递会发现这一步已经做完,直接跳过,不会再做一遍。

归属靠 CAS、epoch 和心跳

所有状态写入都走一段 Redis Lua 脚本做原子 CAS。正在执行的步骤带着 epoch 和心跳,一个卡住又活过来的进程覆盖不了别人已经做完的活。每次调用还持有一把重放租约,30 秒过期、每 10 秒续一次,保证两个 worker 不会同时重放同一次调用。进程被 kill 掉,租约自己就到期。

github.com/hibiken/asynq

队列和命名空间

所有 flow 共用一对队列:gotick 和 gotick_critical。消费端先取 gotick_critical,让已经开始的流程尽快跑完,而不是所有流程一起慢慢挪。事件里带着 flowId,消费端靠它分发。

队列数不随 flow 数量增长,这一点很关键。如果一个 flow 配一对队列,轮询就要遍历所有队列,50 个 flow 就是 100 个队列,系统完全空闲时也在不停刷 Redis。

用 Config.Queue 换命名空间,默认是 gotick。同一个命名空间上的 worker 应该注册同一组 flow。worker 收到自己不认识的 flow 时,会把事件原样转投出去,不会丢掉,因为丢掉会让别人的流程永远停住。滚动发布期间会短暂走到这条路径,几秒内就恢复正常。

但两个不相干的服务共用默认命名空间就麻烦了:每个事件都要转投好几次。8 个 worker、命中率 1/8 时,一个两步的流程要几秒才跑完,而不是几十毫秒。这种情况给它们各自分配命名空间。

gotick.Config{ Queue: "myapp", // default "gotick"}

Render diagnostics