relay 按行领取,不用游标

更新于 · 在 sijie.xyz 查看原条目 ↗

状态: 已在 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 后面时,后者被吞了。UT TestEverySubjectEventReachesANonCoalescingSubscriber 守住这个修复。

游标导致的丢信 bug

最初的 relay 按 seq > cursor 读。序号在事务开始时分配,提交顺序却可能相反。设计审查中发现了丢信:

sequenceDiagram
  participant A as Transaction A
  participant B as Transaction B
  participant O as events
  participant R as relay (cursor version)
  A->>O: INSERT, gets seq 101 (uncommitted)
  B->>O: INSERT, gets seq 102
  B->>O: COMMIT
  R->>O: read seq > 100, sees only 102
  R->>R: cursor moves to 102
  A->>O: COMMIT, 101 only now visible
  Note over R,O: 101 < cursor, never read → lost

修法:不用游标。每行自带 fanned_out_at。没打标记的行一定还会被领到,没有“被跳过”这一说。一条 UT 跑两个交错提交的事务,两条事件都最终被扇出。

端到端时序:一次语料写入

sequenceDiagram
  autonumber
  participant UC as Use case (any write path)
  participant DB as Postgres
  participant T as corpus_notes trigger
  participant R as Relay loop (every process)
  participant Q as River
  participant IX as Subscriber corpus.index
  participant WF as webhook.fanout
  participant WD as webhook.deliver
  UC->>DB: BEGIN, UPDATE corpus_notes …
  DB->>T: AFTER UPDATE (watched column changed)
  T->>DB: INSERT events(corpus.note.changed), pg_notify
  UC->>DB: COMMIT (change and event land or vanish together)
  DB-->>R: NOTIFY wake-up (fallback: 1-minute sweep)
  R->>DB: BEGIN, SELECT events … FOR UPDATE SKIP LOCKED
  R->>Q: InsertTx (one job per matching subscriber, coalesced where opted in)
  R->>DB: UPDATE events SET fanned_out_at, fanout, COMMIT (fan-out exactly once)
  R-->>DB: NOTIFY standmeet_events_fanned
  par in process
    Q->>IX: Handle(event)
    IX-->>Q: ok / error → River backoff retry
  and leaves the instance
    Q->>WF: Handle(event)
    WF->>WF: type globs and scope per endpoint
    WF->>Q: InsertTx webhook.deliver per endpoint (unique by args)
    Q->>WD: Work(endpoint, event)
    WD->>WD: take lease → sign → POST
    WD-->>Q: 2xx / non-2xx → retry on schedule
  end

relay 崩了怎么办

领取、入队、打标记在同一事务里。批处理中途崩溃,三者一起回滚,重启后重新领取。见 message-loss-guarantees。

测试

events 的 UT(relay、触发器、recorder 合计约 26 条):按行领取;交错提交不丢;批大小;poison 标记;多个 relay 并发不重复扇出;崩溃回滚后重领;同一主体的每条事件都到达不合并的订阅方。见 events-test-plan。