引言
定时轮询是最容易实现、也最容易失控的触发方式。一个「每 5 分钟扫一遍订单表,把待发货的挑出来」的任务,在数据量小的时候毫无问题,一旦表里有几百万行待处理记录,扫描本身就变成了负担;更麻烦的是延迟被调度周期锁死,业务方要求「支付成功后 3 秒内开始发货」,轮询无论如何都做不到。
事件驱动把「什么时候触发」这个问题交给了事件源:支付系统在支付成功的那一刻发一条消息,工作流消费到消息就启动实例。延迟从分钟级降到毫秒级,扫描成本归零,而且天然携带了业务上下文(订单号、金额、渠道),不需要再回查数据库。
但事件驱动的复杂度不在「发消息」这一步,而在四个隐蔽的地方:同一件事被重复投递怎么办、跨分区的乱序怎么处理、上游洪峰把下游打爆怎么办、以及一个业务键的事件来得太频繁导致触发风暴怎么办。这四个问题在轮询模型里几乎不存在(因为轮询天然去重、天然限流),在事件驱动里却必须显式设计。
本文按「事件源 → 触发映射 → 顺序与去重 → 背压与风暴 → 与定时调度的边界」的顺序展开,重点讲工程上会踩的具体坑。想先看触发机制在整个引擎体系里的位置,可以从 工作流引擎全景与选型 开始;重试与幂等的通用设计在 重试幂等与补偿设计 里有更完整的讨论。
目录
- 事件驱动触发要解决什么
- 事件源的类型与触发模型
- 触发器与工作流实例的映射
- Kafka 集成:消费组与位移提交
- 用消息键保证分区内有序
- Pulsar 的订阅模式与差异
- 事件去重与幂等消费
- 顺序保证与乱序处理
- 事件时间与水位线
- 背压与限流
- 触发风暴与合并窗口
- 事件驱动与定时调度的边界
- 事件网关与 CloudEvents 信封
- 回调、Webhook 与轮询的取舍
- 与 outbox 模式的配合
- 死信与人工兜底
- 观测与告警
- 落地路线图
- 权衡取舍
- 常见坑清单
- 小结
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. 常见坑清单
- 用
enable.auto.commit=true,消息还没处理完位移就提交,崩溃后事件永久丢失。 max.poll.interval.ms保持默认 5 分钟而单批处理超过它,触发 rebalance 后大范围重复消费。- 生产端未开幂等生产者,重试导致同一分区内消息乱序,破坏「分区内有序」的假设。
- 消费端把同一分区的消息丢进线程池并发处理,分区内有序保证失效。
- 用消息时间戳或消费时刻代替业务事件时间,上游积压时业务判断全错。
- 去重表没有 TTL,几个月后膨胀到上亿行,去重查询本身成为瓶颈。
- 先写去重标记再执行业务,业务失败后事件被永久标记为已消费。
- Webhook 接收端同步处理业务逻辑,对方超时重试导致重复率成倍放大。
- 不校验 Webhook 签名,任何人都能伪造「支付成功」事件启动工作流。
- 同一业务键的高频事件直接触发工作流,地址改 20 次就启动 20 个实例。
- 死信队列建了但没有告警与重放工具,坏消息静默堆积无人知晓。
- 只监控 lag 不看端到端延迟,出现「lag 正常但业务响应变慢」的隐蔽故障。
21. 小结
事件驱动触发的核心矛盾是「延迟」与「确定性」的交换:轮询用固定的延迟换来了天然的去重与限流,事件驱动用毫秒级响应换来了必须自己处理重复、乱序与背压的责任。理解了这笔交易,剩下的都是具体机制的选择。
工程上的三条底线是:所有消费端幂等(去重表或业务幂等键,二选一或都用)、所有业务判断基于事件时间(不用消费时刻)、所有下游都有背压(有界队列 + 限流)。这三条任何一条缺失,都会在上游异常时以「重复扣款」「状态错乱」「雪崩」的形式暴露。
如果流程以「按时间周期批量处理」为主,事件驱动不是合适的模型,应该看 Airflow DAG 调度体系 ;如果事件驱动的目的是编排跨服务的分布式事务,那么补偿的设计比触发本身更重要,参见 Saga 与分布式事务补偿 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。