事件驱动触发与消息集成

本文系统讲解工作流的事件驱动触发与消息集成,回答事件源如何映射到工作流实例、Kafka 与 Pulsar 集成时顺序与去重怎么保证、背压与触发风暴如何治理。覆盖推拉两种触发模型、分区键与消费组语义、幂等消费与去重表、事件时间与乱序处理、限流与合并窗口、事件驱动与定时调度的边界、死信兜底与观测告警,并给出可落地的配置片段与代码。

引言

定时轮询是最容易实现、也最容易失控的触发方式。一个「每 5 分钟扫一遍订单表,把待发货的挑出来」的任务,在数据量小的时候毫无问题,一旦表里有几百万行待处理记录,扫描本身就变成了负担;更麻烦的是延迟被调度周期锁死,业务方要求「支付成功后 3 秒内开始发货」,轮询无论如何都做不到。

事件驱动把「什么时候触发」这个问题交给了事件源:支付系统在支付成功的那一刻发一条消息,工作流消费到消息就启动实例。延迟从分钟级降到毫秒级,扫描成本归零,而且天然携带了业务上下文(订单号、金额、渠道),不需要再回查数据库。

但事件驱动的复杂度不在「发消息」这一步,而在四个隐蔽的地方:同一件事被重复投递怎么办、跨分区的乱序怎么处理、上游洪峰把下游打爆怎么办、以及一个业务键的事件来得太频繁导致触发风暴怎么办。这四个问题在轮询模型里几乎不存在(因为轮询天然去重、天然限流),在事件驱动里却必须显式设计。

本文按「事件源 → 触发映射 → 顺序与去重 → 背压与风暴 → 与定时调度的边界」的顺序展开,重点讲工程上会踩的具体坑。想先看触发机制在整个引擎体系里的位置,可以从 工作流引擎全景与选型 开始;重试与幂等的通用设计在 重试幂等与补偿设计 里有更完整的讨论。

目录

  1. 事件驱动触发要解决什么
  2. 事件源的类型与触发模型
  3. 触发器与工作流实例的映射
  4. Kafka 集成:消费组与位移提交
  5. 用消息键保证分区内有序
  6. Pulsar 的订阅模式与差异
  7. 事件去重与幂等消费
  8. 顺序保证与乱序处理
  9. 事件时间与水位线
  10. 背压与限流
  11. 触发风暴与合并窗口
  12. 事件驱动与定时调度的边界
  13. 事件网关与 CloudEvents 信封
  14. 回调、Webhook 与轮询的取舍
  15. 与 outbox 模式的配合
  16. 死信与人工兜底
  17. 观测与告警
  18. 落地路线图
  19. 权衡取舍
  20. 常见坑清单
  21. 小结

1. 事件驱动触发要解决什么

轮询模型有三个绕不开的代价。第一是延迟:*/5 * * * * 意味着平均 2.5 分钟的等待,紧急流程等不起。第二是空转:绝大多数轮询周期里没有新数据,但查询照跑,数据库连接与 CPU 被白白消耗。第三是状态无关:轮询只能看到「当前表里有什么」,看不到「刚刚发生了什么」,因此无法区分「这条记录是新产生的」还是「这条记录被改过又被改回来」。

事件驱动把触发信号与业务状态解耦:事件源负责在状态变化的瞬间广播一个事实,工作流负责消费这个事实并推进流程。它带来三个直接收益:延迟从分钟级压到毫秒级;触发成本与数据量脱钩(没有事件就没有消耗);事件本身携带了发生时刻的完整上下文,工作流不需要回查就能判断该走哪条分支。

代价同样明确:投递语义从「精确一次查询」退化为「至少一次投递」,去重责任转移到消费端;事件流是有序性有限的(分区内有序、跨分区无序),需要额外机制兜底;上游的流量特征直接冲击下游,没有轮询天然具备的削峰效果。本文后续所有小节,本质上都在处理这三个代价。

2. 事件源的类型与触发模型

按投递机制,事件源可以分成四类,工程上的处理方式差别很大:

事件源典型实现语义主要风险
消息队列Kafka、Pulsar、RabbitMQ、RocketMQ至少一次重复投递、乱序
变更数据捕获Debezium 读 binlog / WAL至少一次表结构变更、快照与增量衔接
推送回调Webhook、S3 事件通知、K8s Event至多一次丢失、伪造、无序
主动拉取数据库轮询、API 分页拉取精确一次(按查询结果)延迟、空转、水位管理

触发模型则分三种:**推送(push)**由事件源主动调用接收端,延迟最低但要求接收端有稳定地址与容量;**拉取(pull)**由消费者主动取,消费者完全控制速率,天然背压,代价是空转;**长轮询(long poll)**介于两者之间,消费者发起请求后服务端挂起,有数据才返回,Kafka 的 fetch.max.wait.ms 与 Pulsar 的 receiverQueueSize 都是这个思路。

选型的判断标准是「谁控制速率」。上游是稳定的内部系统、下游有余量时用推送;上游可能突发、下游容量有限时用拉取,让消费者自己决定什么时候取多少。多数生产系统最终落在「消息队列拉取 + 关键回调推送」的混合形态上。

长轮询是两种模型的折中,实现上只是给拉取加了一个「挂起等待」的语义:

# 长轮询:没有数据时挂起 20 秒再返回,避免高频空转
resp = requests.get(
    EVENTS_API,                              # 内部事件接口地址
    params={"cursor": cursor, "wait": 20, "limit": 100},
    timeout=25,          # 必须大于服务端 wait,否则客户端先超时
)
# 返回空列表说明窗口内没有新事件,立即发起下一次长轮询

长轮询的关键是「客户端超时 > 服务端挂起时间」,否则每次请求都会在服务端返回前被客户端掐断,退化成高频短轮询。它适合无法部署消息队列、又必须近实时响应的外部系统对接场景。

3. 触发器与工作流实例的映射

事件与工作流实例之间有三种映射关系,选错了会导致实例数爆炸或状态丢失:

1:1   一个事件 -> 一个工作流实例
      例:每笔支付启动一个发货流程。用业务键做 WorkflowId,天然幂等。

N:1   多个事件 -> 一个工作流实例
      例:一个订单的「创建/支付/取消」事件都投给同一个订单实例。
      用长驻实例 + 信号(Signal)实现,实例生命周期与订单一致。

1:N   一个事件 -> 多个工作流实例
      例:一次大促开始,触发所有仓库的备货流程。
      用事件里的维度键展开,注意扇出数量要有上限。

N:1 是最容易被低估的模式。它的实现是「用业务键查实例是否已存在,存在就发信号,不存在才启动」,也就是 Temporal 里的 WorkflowExecutionAlreadyStarted 捕获后改发信号。这个模式的好处是订单的完整状态集中在一个实例里,查询与补偿都简单;风险是实例变成「长驻」,历史会持续增长,必须配合历史截断策略。

// N:1:先尝试启动,已存在则改发信号
try {
    client.newUntypedWorkflowStub("OrderWorkflow", options).start(input);
} catch (WorkflowExecutionAlreadyStarted e) {
    client.newUntypedWorkflowStub("order-" + orderId).signal("onEvent", event);
}

这里有一个必须提前决策的点:实例结束后再来事件怎么办。订单流程走完进入终态,此时又收到一条迟到的「取消」事件,直接丢弃会造成业务不一致。常见做法是让实例在终态保留一个短窗口(比如用 Workflow.await 多等 24 小时),或者把迟到事件路由到一个专门的「异常处理流程」。参见 状态机引擎与状态流转 里关于非法流转防御的讨论。

4. Kafka 集成:消费组与位移提交

Kafka 是最常见的事件源。工作流侧通常有一个消费者进程,从主题拉消息、解析事件、启动或推进工作流实例。三个参数决定语义:

# 关闭自动提交,由代码在处理完成后手工提交
enable.auto.commit=false
# 第一次消费没有位移时的行为
auto.offset.reset=earliest
# 单次拉取的最大记录数,控制单批处理规模
max.poll.records=200

位移提交时机决定了投递语义。处理完再提交是「至少一次」,崩溃时最多重复处理一批;拉取后立刻提交是「至多一次」,崩溃会丢消息。工作流场景必须选前者,重复由幂等消费兜住。

max.poll.interval.ms(默认 5 分钟)是另一个高频事故点:如果单批处理时间超过这个值,消费者会被判定为死亡并触发 rebalance,分区被重新分配,下一轮从上次提交的位移重新消费,造成大范围重复。调优方向是减小 max.poll.records 或把耗时处理异步化,而不是盲目调大 max.poll.interval.ms。

消费者组的分区分配也有讲究。工作流的消费者通常需要「按业务键亲和」——同一个订单的事件必须落到同一个实例上,这就要求消息的生产端用订单号做 key,消费端不做额外的重分区。更多 Kafka 的底层语义可以对照 Kafka 消息队列专题 阅读。

5. 用消息键保证分区内有序

Kafka 只保证分区内有序,跨分区无序。要让「同一个订单的事件按顺序处理」,必须让这些事件进入同一个分区,也就是用同一个 key:

// 生产端:用业务键做 key,默认分区器对 key 做 murmur2 哈希后取模
ProducerRecord<String, String> record =
    new ProducerRecord<>("order-events", orderId, payload);

// 消费端:单分区单线程顺序处理,或按 key 分派到同一处理槽

这里有一个直接的容量约束:分区数就是该主题的最大并行度。分区太少,消费者数量超过分区数时有消费者空转;分区太多,单个分区吞吐下降且 rebalance 变慢。常见的做法是按「日均事件量 / 单分区目标吞吐」估算,比如单分区 10 MB/s、日均 500 GB 则需要 50 个以上分区。

顺序保证还有两个隐含前提。第一,生产端要开启 max.in.flight.requests.per.connection=1 或使用幂等生产者(enable.idempotence=true),否则重试可能导致同一分区内乱序。第二,消费端不能为了加速而把同一分区的消息丢进线程池并发处理——一旦并发,分区内有序的保证就失效了。如果确实需要并发,就按 key 做二级分派(同一个 key 永远进同一个处理槽),而不是简单轮转。

6. Pulsar 的订阅模式与差异

Pulsar 在订阅模型上比 Kafka 更细,四种订阅模式直接对应不同的编排需求:

订阅模式语义适用场景
Exclusive单消费者独占严格全局有序
Failover主备切换有序 + 高可用
Shared多消费者共享,轮询投递高吞吐、不要求顺序
Key_Shared按 key 分派,同 key 单消费者按业务键有序 + 并行

Key_Shared 是 Pulsar 相对 Kafka 的明显优势:它把「按 key 有序」和「多消费者并行」同时做到了,不需要用分区数硬性限制并行度。工作流场景里,如果一个主题承载多种业务键(订单、退款、库存),用 Key_Shared 可以避免为每种键单独建主题。

Pulsar 的另一个差异是计算存储分离:BookKeeper 存消息、Broker 无状态,因此扩容与堆积处理更平滑;同时分层存储(tiered storage)可以把老消息卸载到对象存储,适合「事件要保留 30 天供重放」的场景。代价是运维组件更多,小规模团队的成本不划算。

7. 事件去重与幂等消费

至少一次投递意味着同一事件可能被处理两次。去重的第一道防线是事件自带的唯一 ID(event_id / message_id),而不是业务内容:

CREATE TABLE consumed_event (
    event_id    VARCHAR(64) PRIMARY KEY,
    consumer    VARCHAR(64) NOT NULL,
    consumed_at TIMESTAMP   NOT NULL DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_consumed_at (consumed_at)
);

INSERT INTO consumed_event (event_id, consumer) VALUES (?, ?)
ON CONFLICT DO NOTHING;
-- 影响行数为 0 表示已消费过,直接 ack 跳过

三个细节决定这个方案是否可靠。一是键的粒度:consumer 字段不能省,否则同一事件被两个不同流程消费时,第二个会被误判为重复。二是表的清理:去重表必须有 TTL(比如保留 7 天),否则会无限增长;清理窗口要大于「消息最大重投间隔」,否则老事件重投时会被当成新事件。三是插入与处理的原子性:如果先去重表插入、再执行业务逻辑,业务失败后事件已被标记为「已消费」,就永久丢了。正确顺序是「业务处理与去重标记在同一事务里提交」,或者用「先插入占位、失败后删除占位」的补偿写法。

对于无法做到事务的跨系统场景,第二道防线是业务幂等键:业务单号 + 步骤名,由下游系统按此去重。这与重试幂等设计里讲的幂等键是同一套方法。

当去重压力很高(每秒数万事件)时,数据库去重表会成为瓶颈,可以用 Redis 做前置拦截:

def is_duplicate(event_id, ttl_seconds=7 * 86400):
    # SET NX 原子占位,成功说明是新事件
    ok = redis.set(f"dedup:{event_id}", "1", nx=True, ex=ttl_seconds)
    return not ok

Redis 方案的代价是不持久:Redis 故障或数据淘汰后,去重能力丢失,会出现重复处理。因此它只适合作为「前置过滤器」——Redis 判重后仍要在业务侧保留幂等键兜底,形成「Redis 挡量 + 数据库兜底」的两级结构。

8. 顺序保证与乱序处理

即使同一 key 走了同一分区,乱序仍可能从三个地方进来:生产端多线程并发发送(不同线程的发送完成顺序不定)、上游系统本身的事件时间无序(支付回调可能晚于发货事件到达)、以及重投递(失败的事件重试后落在新事件之后)。

处理乱序的标准做法是给事件带上单调递增的版本号或序列号,由消费端做状态比较:

def handle(event):
    current = load_state(event.aggregate_id)
    if event.version <= current.version:
        log.warning("stale event dropped id=%s v=%s", event.id, event.version)
        return  # 旧事件,丢弃
    if event.version > current.version + 1:
        # 出现空洞,说明中间事件还没到
        buffer_event(event)          # 暂存,等补齐
        schedule_gap_timeout(event.aggregate_id, seconds=30)
        return
    apply(event)
    drain_buffer(event.aggregate_id)  # 补齐后尝试放行暂存事件

这段逻辑有三个必须显式设计的点:版本号从哪里来(由产生事件的聚合根维护,不能用消息时间戳代替);空洞等待多久(等待超时会永久卡住,必须有超时后「跳过并告警」的路径);缓冲区放哪(内存缓冲在重启后丢失,重要流程要落到外部存储)。乱序窗口的大小取决于上游的最大延迟,用生产环境的 P99 延迟加上余量来定,而不是拍脑袋。

9. 事件时间与水位线

处理时间(processing time)与事件时间(event time)的区分,是事件驱动工作流最容易出错的地方。 处理时间是消费者拿到消息的墙上时钟,事件时间是事件在业务上真实发生的时刻。两者在上游积压时可能相差几小时。

后果很直观:如果工作流用处理时间判断「订单是否超时未支付」,当上游积压了 3 小时,一批本已超时的订单会被判为「刚创建」,走上错误的正常分支。正确做法是把业务判断全部建立在事件时间上,并用事件自带的 occurred_at 字段:

// 错误:用消费时刻做业务判断
if (Instant.now().isAfter(order.createdAt().plus(Duration.ofMinutes(30)))) { ... }

// 正确:用事件时间与业务时钟
if (event.occurredAt().isAfter(order.deadline())) { ... }

水位线(watermark)是「系统认为不会再有更早事件到达」的时间点,用于决定何时可以安全地触发窗口计算。工作流引擎里没有内建的 watermark 概念,需要自己实现:维护一个「已见最大事件时间」,当它减去允许延迟超过某个阈值时,就认为窗口可以关闭。这条规则要和乱序缓冲配合使用,否则会出现「窗口已关闭但迟到事件又到」的矛盾。

10. 背压与限流

事件驱动的消费者是「上游给多少就吃多少」,如果没有背压,一次上游故障恢复后的补偿性投递就能把下游打爆。背压的实现分三层:

第一层是拉取节奏。消费者滞后(consumer lag)是天然的背压信号:当 lag 超过阈值时,主动 pause() 分区,等处理速率追上再 resume()。Kafka 的 pause/resume 与 Pulsar 的 receiverQueueSize 都能做到这一点,不需要丢弃消息。

第二层是有界队列。消费者与工作流引擎之间通常有一个内存队列,队列必须设上限,满了就停止拉取(阻塞式背压),而不是无限堆积后 OOM。

第三层是限流。当下游是外部系统(支付网关、短信通道)时,必须按对方的配额限流:

// 令牌桶:每秒 50 个请求,允许突发 100
RateLimiter limiter = RateLimiter.create(50.0);
// 每个任务执行前获取令牌,超时则抛出让重试机制处理
if (!limiter.tryAcquire(2, TimeUnit.SECONDS)) {
    throw new RetryableException("rate limited");
}

限流的粒度要和下游配额对齐:如果外部系统按商户维度限流,那么限流器也必须是每商户一个桶,全局共享一个桶会导致某个大商户把额度吃光。这与 重试幂等与补偿设计 里「重试预算」的思路一致——限流与重试是同一件事的两面。

11. 触发风暴与合并窗口

幂等启动能挡住重复事件,但挡不住「同一业务键的高频不同事件」。典型的触发风暴场景:一个用户在一分钟内修改了 20 次收货地址,每次修改都发一条事件;如果每条事件都触发一次「重新计算运费」的工作流,系统会被同一用户的重复计算淹没。

三种治理手段,按介入位置排序:

去抖(debounce)  收到事件后等 N 秒,期间有新事件则重新计时,只处理最后一次
节流(throttle)  每 N 秒最多处理一次,丢弃窗口内的中间事件
合并(coalesce)  把窗口内的多个事件合并成一个批次事件,全部处理但只触发一次

去抖适合「只关心最终状态」的场景(地址修改),节流适合「限速」的场景(日志上报),合并适合「不能丢事件但可以批量」的场景(库存变更)。

去抖的落地有两种方式。轻量做法是在消费者侧维护一个「最近事件时间」表,距上次事件不足窗口就只更新标记不发信号;重量做法是用一个长驻的工作流实例做聚合器,事件以信号形式投递给它,实例内部用 Workflow.await 做窗口等待。后者更容易保证「窗口内的最后一次事件一定被处理」,因为信号与状态都在引擎里持久化了。

// 聚合器实例:收到信号后重新计时,静默 30 秒才真正处理
@Override
public void run() {
    while (true) {
        Workflow.await(() -> this.dirty);   // 等到有事件
        this.dirty = false;
        boolean quiet = Workflow.await(Duration.ofSeconds(30), () -> this.dirty);
        if (quiet) {                        // 窗口内无新事件,执行合并处理
            activities.recompute(this.aggregateId);
        }                                   // 否则回到循环,重新计时
    }
}

聚合器的实例数等于活跃业务键的数量,因此需要给业务键设上限(比如只对 VIP 用户做实时聚合,其余走批量)。长驻实例要配合历史截断,否则窗口反复触发的实例会很快撞上事件数上限。

12. 事件驱动与定时调度的边界

不是所有触发都该事件驱动。判断标准是「这个动作有没有一个明确的、值得即时响应的前置事实」:

场景推荐触发理由
支付成功即发货事件驱动有明确前置事实,延迟敏感
每日凌晨生成报表定时调度无前置事实,按业务周期
文件到达后解析事件驱动文件到达是事实,且时间不定
超时未支付则取消定时兜底本质是「时间流逝」,没有事件
库存低于阈值补货事件驱动 + 定时兜底事件为主,定时防漏
上游数据同步完成后跑下游事件驱动依赖是数据就绪,不是时间

最实用的模式是**「事件为主、定时兜底」**:正常路径由事件触发,同时挂一个低频定时任务扫描「应该已触发但状态没变」的记录。这样即使事件丢失(回调模型下是至多一次),系统也能在兜底周期内自愈。兜底任务的查询条件是「状态未推进且超过预期时间」,而不是全表扫描。

反过来,定时任务也有它不可替代的位置:任何「按时间流逝推进」的语义(超时、SLA、宽限期结束)本质上都是时间触发,硬用事件模拟会引入大量不必要的状态。调度与回填的细节在 Airflow DAG 调度体系 里展开。

兜底扫描的查询条件必须能被索引命中,否则每次全表扫描会把数据库拖垮:

-- 兜底任务:找出「应已推进但状态未变」的记录
SELECT id, business_key, status, updated_at
FROM order_flow
WHERE status IN ('WAIT_PAY', 'WAIT_SHIP')          -- 未终态
  AND updated_at < NOW() - INTERVAL '15 minutes'   -- 超过预期推进时间
  AND updated_at > NOW() - INTERVAL '24 hours'     -- 只看近 24 小时,避免历史堆积
ORDER BY updated_at
LIMIT 1000;
-- 需要 (status, updated_at) 复合索引,否则退化为全表扫描

兜底任务的频率要低于事件到达的典型间隔,否则会与正常路径重复触发;同时它必须有幂等保护,因为兜底触发的事件可能与正常事件同时到达。

13. 事件网关与 CloudEvents 信封

当系统接入多个事件源时,每个源的事件格式都不一样,工作流侧会被迫写一堆适配代码。解决办法是在事件源与工作流之间加一层事件网关,统一事件信封:

{
  "specversion": "1.0",
  "id": "b1c2d3e4-...",
  "source": "/payment/order-service",
  "type": "com.example.order.paid.v1",
  "subject": "order-1001",
  "time": "2026-10-07T10:30:00+08:00",
  "datacontenttype": "application/json",
  "data": { "orderId": "1001", "amount": "199.00", "channel": "wechat" }
}

这是 CloudEvents 1.0 的标准字段,价值在于把「路由所需的元数据」与「业务载荷」分离:网关可以只依赖 type 与 subject 做路由,完全不解析 data,因此新增业务类型不需要改网关代码。id 直接用作去重键,time 用作事件时间,subject 用作业务键。

网关承担四项职责:协议适配(把 HTTP 回调、Kafka 消息、CDC 事件统一成同一信封)、签名与鉴权校验、schema 校验(拒绝格式非法的事件,避免脏数据进入工作流)、以及路由(按 type 分派到不同主题或不同工作流定义)。注意 type 字段必须带版本后缀(.v1),这样同一事件的结构升级可以并行存在两套,与 Temporal 与持久化执行 里的事件历史版本兼容策略呼应。

14. 回调、Webhook 与轮询的取舍

Webhook 是最不可靠的事件源,因为它本质上是「别人调你的 HTTP 接口」,而 HTTP 调用的成功与否取决于你的服务当时是否可用。三种失败模式都要处理:对方没发(对方系统故障或配置错误)、发了你没收到(网络抖动、服务重启)、收到了但处理失败(返回 5xx 后对方重试,你又处理了一次)。

@app.post("/webhook/payment")
def receive(request):
    body = request.get_data()
    if not verify_signature(body, request.headers.get("X-Signature"), secret):
        return "invalid signature", 401        # 拒绝伪造
    event = parse(body)
    if is_duplicate(event["id"]):              # 去重表兜底
        return "ok", 200                       # 已处理过,返回成功避免对方重试
    enqueue(event)                             # 先落盘/入队,再返回
    return "ok", 200                           # 快速返回,处理异步化

三个要点:签名校验必须做,否则任何人都能伪造「支付成功」;去重必须在返回 200 之前,否则对方的正常重试会产生重复;处理必须异步化,同步处理慢会导致对方超时并重试,把重复率放大。此外签名要带时间戳并校验时效,防止重放攻击。

无论回调多可靠,都建议配一个轮询兜底:定时查询对方系统的「最近交易列表」,与本地记录比对,找出漏掉的事件。这是唯一能发现「对方根本没发」的手段。

签名校验的实现有一个常见错误:直接对原始 body 做 HMAC 却在校验前先做了一次 JSON 反序列化再序列化,导致字节不一致而校验永远失败。正确做法是在解析之前对原始字节做校验:

def verify_signature(raw_body: bytes, header: str, secret: str) -> bool:
    ts, sig = header.split(",", 1)                 # 形如 t=1696000000,v1=abcdef
    ts = ts.split("=", 1)[1]
    if abs(time.time() - int(ts)) > 300:           # 5 分钟时效,防重放
        return False
    expected = hmac.new(
        secret.encode(), f"{ts}.".encode() + raw_body, hashlib.sha256
    ).hexdigest()
    return hmac.compare_digest(expected, sig.split("=", 1)[1])

时间戳与 compare_digest(恒定时间比较)都不能省:前者防重放攻击,后者防时序侧信道。

15. 与 outbox 模式的配合

事件驱动的一个根本难题是「业务数据写库」与「事件发送」这两件事必须同时成功或同时失败。先写库后发消息,发消息失败就丢了事件;先发消息后写库,写库失败就产生了「幽灵事件」。

**事务性发件箱(outbox)**是标准解法:把要发的事件与业务数据写在同一个数据库事务里,落到一张 outbox 表,再由独立的投递进程读表发送:

BEGIN;
INSERT INTO orders (id, status, amount) VALUES (?, 'PAID', ?);
INSERT INTO outbox (id, aggregate_id, event_type, payload, created_at)
VALUES (?, ?, 'order.paid.v1', ?, NOW());
COMMIT;
-- 投递进程:轮询或 CDC 读 outbox,发送成功后标记 sent_at

投递进程有两种实现:轮询(简单,延迟秒级,需要处理并发投递的重复)与 CDC(Debezium 读 binlog,延迟毫秒级,不占用业务库查询资源)。无论哪种,投递都是至少一次,消费端仍要去重——outbox 解决的是「不丢」,去重表解决的是「不重」,两者缺一不可。outbox 与 Saga 协同模式的配合在 Saga 与分布式事务补偿 里有更完整的讨论。

16. 死信与人工兜底

总有一类事件无法自动处理:格式解析失败、引用的业务实体不存在、触发了业务规则里没覆盖的分支。这些事件不能无限重试,也不能静默丢弃。

# 消费者重试策略:3 次快速重试 + 指数退避,之后进死信
retry:
  max_attempts: 5
  backoff:
    initial: 1s
    multiplier: 2.0
    max: 5m
dead_letter:
  topic: workflow-events-dlq
  include_headers: true      # 保留原始 topic/partition/offset,便于定位
  include_stacktrace: true

死信队列的设计有三个要点。一是保留原始元数据(来源主题、分区、位移、失败原因、重试次数),否则重放时无从下手。二是必须配监控与告警,死信深度是「有事件正在悄悄失败」的唯一信号,没有告警的死信队列等于黑洞。三是提供重放工具,修复 bug 后要能把死信按原顺序重新投递,而不是手工造消息。重放时要特别注意幂等——这些事件之前可能已经部分执行过。

对于无法自动修复的事件(比如引用的订单确实不存在),需要一条人工处理路径:把它转成一个待办任务,交由运营在后台界面处理。这与 人工任务与审批流表单 里的人工节点是同一类设计。

17. 观测与告警

事件驱动系统有四个必须监控的维度:

指标含义告警阈值参考
consumer lag分区积压的消息数持续 5 分钟超过 10 万
消费速率每秒处理事件数环比下降 50%
触发成功率事件成功启动/推进实例的比例低于 99.5%
去重命中率被判为重复的事件比例突增 5 倍
死信深度进入死信队列的事件数大于 0 即告警
端到端延迟事件发生到实例推进的耗时P99 超过 30 秒

consumer lag 是最重要的单一指标,它同时反映上游流量与下游容量。但要注意 lag 的绝对值没有意义,必须结合流量基线看:一个大促期间 lag 涨到 100 万可能是正常的(消费速率也在涨),而平时 lag 涨到 1 万就说明消费者挂了。

端到端延迟是唯一能反映「业务是否真的及时响应」的指标,它的测量方式是让事件带上产生时间戳,在工作流推进时计算差值。这个指标会暴露那些「lag 正常但处理慢」的隐蔽问题。观测体系的完整设计见 工作流可观测与调试 。

18. 落地路线图

  • 第 1 周:选一个延迟敏感、事件源稳定的流程(比如支付成功通知),用 Kafka 主题接进来,实现「消费 → 幂等启动工作流」。重点验证重复投递下的行为。
  • 第 2 周:加入去重表与事件 ID,做一次人为重投测试(把同一批消息重新投递),确认不会产生重复实例。
  • 第 3 周:接入背压与限流,用压测工具制造 10 倍洪峰,观察 lag 是否可控、下游是否被打爆。
  • 第 4 周:补上死信队列、告警与重放工具,做一次「制造坏消息 → 进死信 → 修复 → 重放」的完整演练。
  • 第 5 周:为关键流程加定时兜底扫描,验证「事件丢失时系统能自愈」。

试点流程要选「事件量中等、业务影响可控」的,不要一上来就接支付核心链路。验证幂等最简单的方法是在消费者里人为注入一次「处理后抛异常」,观察重投时是否重复执行业务动作。

19. 权衡取舍

选择收益代价
事件驱动触发延迟毫秒级、无空转需处理重复、乱序、背压
定时轮询实现简单、天然去重限流延迟受周期限制、空转成本
推送模型延迟最低要求接收端稳定与有容量
拉取模型消费者控速、天然背压空转、需要位移管理
分区键保序分区内严格有序并行度受分区数限制
Pulsar Key_Shared保序与并行兼得运维组件多、生态较小
去重表可靠、可审计需要 TTL 与清理、写放大
去抖合并抑制触发风暴引入固定延迟、可能丢中间态
事件为主 + 定时兜底兼顾延迟与可靠性两套逻辑需保持一致
outbox 模式业务与事件不丢多一个投递组件与延迟
死信队列不丢异常事件需要重放工具与人工流程

20. 常见坑清单

  1. 用 enable.auto.commit=true,消息还没处理完位移就提交,崩溃后事件永久丢失。
  2. max.poll.interval.ms 保持默认 5 分钟而单批处理超过它,触发 rebalance 后大范围重复消费。
  3. 生产端未开幂等生产者,重试导致同一分区内消息乱序,破坏「分区内有序」的假设。
  4. 消费端把同一分区的消息丢进线程池并发处理,分区内有序保证失效。
  5. 用消息时间戳或消费时刻代替业务事件时间,上游积压时业务判断全错。
  6. 去重表没有 TTL,几个月后膨胀到上亿行,去重查询本身成为瓶颈。
  7. 先写去重标记再执行业务,业务失败后事件被永久标记为已消费。
  8. Webhook 接收端同步处理业务逻辑,对方超时重试导致重复率成倍放大。
  9. 不校验 Webhook 签名,任何人都能伪造「支付成功」事件启动工作流。
  10. 同一业务键的高频事件直接触发工作流,地址改 20 次就启动 20 个实例。
  11. 死信队列建了但没有告警与重放工具,坏消息静默堆积无人知晓。
  12. 只监控 lag 不看端到端延迟,出现「lag 正常但业务响应变慢」的隐蔽故障。

21. 小结

事件驱动触发的核心矛盾是「延迟」与「确定性」的交换:轮询用固定的延迟换来了天然的去重与限流,事件驱动用毫秒级响应换来了必须自己处理重复、乱序与背压的责任。理解了这笔交易,剩下的都是具体机制的选择。

工程上的三条底线是:所有消费端幂等(去重表或业务幂等键,二选一或都用)、所有业务判断基于事件时间(不用消费时刻)、所有下游都有背压(有界队列 + 限流)。这三条任何一条缺失,都会在上游异常时以「重复扣款」「状态错乱」「雪崩」的形式暴露。

如果流程以「按时间周期批量处理」为主,事件驱动不是合适的模型,应该看 Airflow DAG 调度体系 ;如果事件驱动的目的是编排跨服务的分布式事务,那么补偿的设计比触发本身更重要,参见 Saga 与分布式事务补偿 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 调度、回填与补数