relay 按行领取,不用游标
上级:events
状态: 已在 v0.1.76 发布(2026-09-27)—— 设计与落地记录见 StandMeet 仓库的 docs/design/event-bus-outbox-webhooks.md。
relay 用 SELECT … WHERE fanned_out_at IS NULL AND poisoned_at IS NULL ORDER BY seq FOR UPDATE SKIP LOCKED 领取还没扇出的行,在同一事务里入队并打上 fanned_out_at。它从不按序号游标读,因为游标会丢事件。
relay 怎么跑
- relay 是每个进程里的一个循环(
internal/infra/events/run.go),不是 River 任务。SKIP LOCKED已经防住重复扇出,所以不需要选主。高频的周期任务只会把river_job塞满。 - 触发器和
Record发的NOTIFY standmeet_events唤醒它。一个 1 分钟的周期任务events relay sweep只负责戳它一下,兜住监听重连期间丢掉的 NOTIFY。 - 每轮最多领 200 行。一次 2000 篇的导入不会变成一个巨大的事务。
- 对每个(事件,订阅方)匹配插入一个 River job,并在同一事务里打上
fanned_out_at和fanout(每个订阅方拿到了哪个 job)。 - 结果:扇出恰好一次,处理至少一次。
- 某轮失败就退避:2 秒起,翻倍,封顶 1 分钟。之后失败的那批逐行重试,一行出错不阻塞其它行。同一行失败 5 次就标为 poison(
poisoned_at),放到一边并告警;owner 在任务面板的事件详情里把它放回队列(events.requeue)。投递失败交给各自的任务重试(retry-has-one-owner)。 - 扇出提交后,relay 发
NOTIFY standmeet_events_fanned,带上事件 id。写回执就等这个(AwaitFanout,async-response-contract)。
同样的 SKIP LOCKED 领取已经在跑微站构建队列:builder-claim-skip-locked。
合并由每个订阅自己选择
同一批里,设了 Coalesce: true 的订阅对每个主体只拿一个任务,对应最新那条事件。只有 corpus.index 设了它:它重读笔记的当前状态,最新一次变更就覆盖了前面的。webhook 和邮件从不合并:同一主体的两条事件是两件事实。
- 没用 River 的唯一任务来做这件事。River 要求唯一任务的
ByState包含running,于是索引任务运行期间到来的变更会并进正在运行的任务里而丢掉。 - 早期实现对所有订阅都合并,结果同一批里
supplier.activated跟在supplier.connected后面时,后者被吞了。UTTestEverySubjectEventReachesANonCoalescingSubscriber守住这个修复。
游标导致的丢信 bug
最初的 relay 按 seq > cursor 读。序号在事务开始时分配,提交顺序却可能相反。设计审查中发现了丢信:
修法:不用游标。每行自带 fanned_out_at。没打标记的行一定还会被领到,没有“被跳过”这一说。一条 UT 跑两个交错提交的事务,两条事件都最终被扇出。
端到端时序:一次语料写入
relay 崩了怎么办
领取、入队、打标记在同一事务里。批处理中途崩溃,三者一起回滚,重启后重新领取。见 message-loss-guarantees。
测试
events 的 UT(relay、触发器、recorder 合计约 26 条):按行领取;交错提交不丢;批大小;poison 标记;多个 relay 并发不重复扇出;崩溃回滚后重领;同一主体的每条事件都到达不合并的订阅方。见 events-test-plan。