Skip to content

消费者崩溃与重投(consumer-crash)

本页结论:消费者在「业务已提交、ACK 未发出」之间崩溃,RabbitMQ 会重投同一条消息;幂等表让这次重投只产生一次 duplicate_skipped,而不会重复落库——这正是 at-least-once 下业务必须预期重复的含义。

适用场景

  • 理解 at-least-once 投递的具体作用范围:Broker 保证「至少投一次」,不保证「只投一次」。
  • 验证规格 §5.4 的幂等消费基准实现:DB 提交后才 ACK,崩溃窗口由幂等表兜底。
  • 这是 Phase 1 的退出条件实验:从零启动 RabbitMQ,复现「崩溃 → 重投 → 幂等拦截」。

崩溃窗口在哪里

手动 ACK 的消费流程存在一个无法消除的窗口:

如果第 2 步之后、第 3 步之前进程终止,Broker 认为消息从未被确认,必然重投。反过来,若先 ACK 后写库,则存在「ACK 已发、业务未提交」的丢消息窗口。顺序只能是:业务事务提交 → ACK,重复交给幂等表处理。

实验步骤

bash
npm run lab -- rabbitmq consumer-crash

编排分两轮:

  1. 第一轮:Producer 发送 3 条消息;Consumer 处理完第 1 条(business_committed)后,在 ACK 前执行崩溃注入(Runtime.halt(137)),进程退出码 137。lab 断言崩溃确实发生——失败注入必须被观察到,而不是「恰好没触发」。
  2. 第二轮:lab 重启 Consumer。第 1 条消息带 redelivered=true 再次到达,幂等表命中,status=duplicate_skipped 并安全 ACK;随后正常处理第 2、3 条。

两轮共用同一个 SQLite 幂等库(.lab/ 目录),模拟同一业务实例重启。

断言

断言期望说明
crashExitCode137崩溃注入确实发生
receivedTotal43 条唯一消息 + 1 条重投
redeliveredCount1仅 order-1001 被重投
duplicatesObserved1幂等表观察到 1 次重复
duplicatesApplied0重复未产生任何业务写入
business_rows3orders 表恰好 3 行

提交快照

verifiedrabbitmq / consumer-crashbroker 4.1.4 · java-amqp-client-5.34.0
镜像rabbitmq:4.1.4-management@sha256:294b01e1796a8acede4619f32a1c394fae1f8021e57986ea01aad38dc2a4f502
捕获时间2026-08-19T07:55:38.494Z
耗时 / 退出码14776 ms / exit 0
断言
crashExitCode137
crashAfterBusinessCommit1
receivedTotal4
redeliveredCount1
duplicatesObserved1
duplicatesApplied0
uniqueMessageIds3
business_rows3
queueDepthAfter0
consumerExitCode0
归一化日志
[producer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash routingKey=orders.crash durationMs=<ms> status=confirmed
[producer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.crash routingKey=orders.crash durationMs=<ms> status=confirmed
[producer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.crash routingKey=orders.crash durationMs=<ms> status=confirmed
[producer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash destination=orders.crash confirmed=3 status=done
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash status=crash_injected
[assert] crashExitCode=137 PASS
[assert] crashAfterBusinessCommit=1 PASS
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash consumer=consumer-1 attempt=1 redelivered=true status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.crash attempt=1 status=duplicate_skipped
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.crash consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.crash attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.crash consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.crash attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash queue=orders.crash received=3 status=done
[inspect] timestamp=<ts> level=INFO service=order-service product=rabbitmq lab=consumer-crash business_rows=3 processed_rows=3 status=snapshot
[assert] receivedTotal=4 PASS
[assert] redeliveredCount=1 PASS
[assert] duplicatesObserved=1 PASS
[assert] duplicatesApplied=0 PASS
[assert] uniqueMessageIds=3 PASS
[assert] business_rows=3 PASS
[assert] queueDepthAfter=0 PASS
[assert] consumerExitCode=0 PASS
npm run lab -- rabbitmq consumer-crash

快照中关键三行:

text
[consumer] ... messageId=mid-1 ... status=business_committed   # 第一轮:业务已提交
[consumer] ... messageId=mid-1 ... status=crash_injected        # ACK 前崩溃(exit 137)
[consumer] ... messageId=mid-1 ... redelivered=true ... status=duplicate_skipped  # 第二轮:幂等拦截

常见误区

  • 「ACK 之后才落库」——顺序反了,崩溃窗口变成丢消息窗口。
  • 「重投是小概率,可以不管」——at-least-once 意味着业务必须预期重复,而不是「偶尔可能重复」。
  • 「幂等表只是优化」——它是崩溃窗口的正确性组件,不可省略。
  • redelivered=true 就跳过」——重投消息可能第一次就没处理完,必须走完整幂等流程。

官方资料与版本说明

以统一实验验证消息系统语义边界