Skip to content

hello-mq消息队列、事件流与可靠消息模式

用统一实验场景解释消息系统核心语义,用可运行的容器化 Demo 验证关键结论,用横向矩阵说明能力边界与选型依据

产品覆盖

产品状态代表性分卷实验
RabbitMQ✅ 已落地传统消息队列与灵活路由8 页5 个(basic / routing / consumer-crash / retry-dlq / backlog-recovery)
Apache Kafka✅ 已落地分区式持久日志与事件流8 页4 个(basic / consumer-group / ordering-replay / idempotence-transaction)
Apache RocketMQ✅ 已落地面向业务消息的分布式中间件8 页4 个(basic / fifo-delay / transaction / retry-dlq)
Apache Pulsar✅ 已落地存储计算分离、云原生多租户8 页3 个(basic / subscriptions / redelivery-replay)
Redis Streams⏸ 不在范围Redis 内的追加日志与消费组
NATS + JetStream⏸ 不在范围低延迟 Core NATS 与持久化 JetStream

从同步调用到事件驱动

切换下面的模式,观察调用关系如何从「点对点强耦合」演化为「经 Broker 解耦」;点「播放消息流」可高亮一条消息的流转路径。

交互拓扑:同步调用 → 事件驱动

订单服务逐个直调下游:任一环节慢或挂,整条链路一起慢、一起挂。

订单服务库存服务积分服务通知服务
  • 订单服务 → 库存服务 (HTTP,阻塞等待)
  • 订单服务 → 积分服务 (HTTP,阻塞等待)
  • 订单服务 → 通知服务 (HTTP,阻塞等待)

模型细节与产品对照见消息模型

验证快照示例:消费者崩溃与幂等拦截

下面是一次真实运行的 consumer-crash 实验快照:消费者在业务提交后、ACK 前崩溃(exit 137),重启后重投被幂等表拦截为 duplicate_skipped,业务表最终恰好 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

完整解读见消费者崩溃与重投

横向矩阵速览

四个产品在关键能力上的支持方式不同——「原生」不等于「免费」,「业务实现」也不等于「不可行」。完整矩阵与证据链接见 横向矩阵

能力RabbitMQKafkaRocketMQPulsar
端到端 exactly-once业务实现仅集群内 EOS业务实现业务实现
顺序消息单队列内分区内MessageGroup 内分区 + 订阅类型相关
内置重试与 DLQ组合配置(TTL+DLX)业务实现原生(Broker 内置)组合配置(DeadLetterPolicy)
延迟消息组合配置(TTL+DLX)业务实现原生(定时投递)业务实现
消息回放不适用(ACK 即删)原生(位点/时间戳)原生(位点重置)原生(reset-cursor)

选型没有万能冠军:按输入维度筛选候选,见 选型指南

本地实验 Quick Start

bash
git clone <repo-url> && cd hello-mq
npm install

npm run lab -- list                # 查看产品与实验
npm run lab -- rabbitmq basic      # 启动 RabbitMQ,发 3 条、收 3 条并校验
npm run lab -- rabbitmq clean      # 仅清理 hello-mq 实验资源

环境要求与完整说明见快速开始

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