Skip to content

RocketMQ 集成

本文档面向使用 pragmatic-ddd 框架、需要把领域事件投递到 RocketMQ 的开发者,说明 pragmatic-ddd-rocketmq 提供的两种事件管理器(Remoting / gRPC)、统一配置、订阅者注册、死信与可靠投递的用法。

1. 概述

1.1 核心定位

pragmatic-ddd-rocketmq 在 core 的 IEventManager 端口上提供两种 RocketMQ 实现,作为领域事件的可靠传输通道:本地发布的领域事件事实经订阅者顺序/条件机制处理后,由本模块异步投递到 RocketMQ,并在消费失败时重试/进死信。

跨进程/跨服务的"谁来响应"不归本模块管:订阅者通过 core 的 @ExternalDependency + IDependency 端口(防腐适配器做 HTTP/RPC)调用外部系统;RocketMQ 仅负责把事件可靠送达消息中间件。

实现类协议适用版本特点
RocketMqEventManagerRemoting(rocketmq-client)RocketMQ 4.x / 5.x兼容性好,社区成熟
RocketMqGrpcEventManagergRPC(rocketmq-client-java)RocketMQ 5.x + gRPC Proxy弹性伸缩更好,5.x 推荐

两者均实现 core 的 IEventManager,构造与使用方式一致,可无缝替换。

它解决的典型场景问题

  • 领域事件可靠投递到消息中间件:core 的本地 ThreadPoolEventManager 仅在进程内内存执行,进程退出即丢;本模块把事件异步、持久化地投递到 RocketMQ,作为可靠的传输通道(对比进程内实现),IEventManager 端口屏蔽了中间件差异。
  • 消费失败重试与死信兜底,保证最终一致性:网络抖动或下游故障导致消费失败时,框架按 maxReconsumeTimes 自动重试,耗尽后投递死信队列,避免领域事件丢失,保障本地发布与下游处理的最终一致。
  • 双协议并存,业务不感知协议:既有 4.x 集群只能走 Remoting,新建 5.x 集群(开启 gRPC Proxy)推荐 gRPC;框架用同一套 IEventManager API 适配两者,接入层无需改动业务代码。

1.2 模块依赖与类型关系

text
pragmatic-ddd-rocketmq
  ├── RocketMqConfig            (统一配置; bind 前缀 rocketmq)
  ├── RocketMqConfiguration     (聚合配置门面, 基于 IConfigurationContext)
  ├── Fastjson2EventSerializer  (实现 core IEventSerializer; 默认序列化器)
  ├── RocketMqEventManager      (Remoting; extends AbstractMQEventManager)
  │     └── builder().config().topicResolver().producer().orderManager().serializer().metrics().build()
  └── RocketMqGrpcEventManager  (gRPC; extends AbstractMQEventManager)
        └── builder() 同构, Producer 类型为 gRPC Producer

依赖 core 端口(io.pragmatic.ddd.event.spi):
  IEventManager / IEventSerializer / ITopicResolver / IEventMetrics / ISubscriberOrderManager

本文档聚焦事件管理器用法(第 2~4 节);与 Outbox 可靠投递的配合见第 5 节。

1.3 前置概念

阅读第 2 节前,需认识以下来自 core 的基础术语:

术语来源含义
IDomainEventcore领域事件接口;发布的事件需实现它,提供 entityId 等标识
IEventManagercore事件管理器端口;publish 发布、registerSubscriber 注册订阅者
ITopicResolvercore把事件类型解析为 RocketMQ topic 的组件;两个管理器构造时必填
DeliveryPolicycore投递策略枚举,DELAYED 表示延迟投递
IEventSerializercore事件序列化器;缺省使用 Fastjson2EventSerializer
IEventMetricscore指标采集端口;缺省 NoOpEventMetrics(无操作)

2. 核心概念详解

2.1 统一配置 RocketMqConfig

RocketMqConfig 是两种管理器的统一配置入口。两种构建方式:

java
// 方式一:链式 setter
RocketMqConfig config = new RocketMqConfig()
        .setNameServer("127.0.0.1:9876")          // Remoting 必填
        .setProducerGroup("ORDER_PRODUCER_GROUP")
        .setConsumerGroup("ORDER_CONSUMER_GROUP")
        .setRetryTimesWhenSendFailed(3)
        .setSendMsgTimeout(3000)
        .setMaxReconsumeTimes(16)
        .setDefaultDelayLevel(3);

// 方式二:从配置源按 rocketmq 前缀绑定
RocketMqConfig config = RocketMqConfig.bind(configurationSource);

RocketMqConfiguration 是聚合配置门面,基于统一配置上下文按语义取数(无需感知裸 key):

java
RocketMqConfiguration rmqCfg = new RocketMqConfiguration(configurationContext);
RocketMqConfig config = rmqCfg.config();          // 等价于 RocketMqConfig.bind(...)
String nameServer = rmqCfg.nameServer();          // rocketmq.name-server
String proxyAddr = rmqCfg.proxyAddr();            // rocketmq.proxy-addr
参数默认值适用协议说明
nameServer-RemotingNameServer 地址;框架自建 Producer/Consumer 时必填
proxyAddr-gRPCgRPC Proxy 地址;gRPC 实现必填(可选依赖 rocketmq-client-java
producerGroupDEFAULT_PRODUCER_GROUP通用Producer 组名
consumerGroupPRAGMATIC_DDD_RMQ_CONSUMER通用Consumer 组名(全局唯一)
retryTimesWhenSendFailed3Remoting发送失败重试次数
sendMsgTimeout3000Remoting发送超时(毫秒)
compressMsgBodyOverHowmuch4096Remoting消息体压缩阈值(字节),超过触发压缩
maxReconsumeTimes16通用消费最大重试次数
defaultDelayLevel3通用默认延迟级别(DELAYED 策略使用)

bind 键约定(前缀 rocketmq):name-server / proxy-addr / retry-times-when-send-failed / send-msg-timeout / compress-msg-body-over-howmuch / producer-group / default-delay-level / max-reconsume-timesnameServer 仅在框架自建 Producer/Consumer 时需要;外部注入 Producer 时可不配。

2.2 Remoting 事件管理器

RocketMqEventManager 基于 rocketmq-client 的 DefaultMQProducer / DefaultMQPushConsumer,兼容 4.x / 5.x Broker。

java
RocketMqEventManager eventManager = RocketMqEventManager.builder()
        .config(new RocketMqConfig().setNameServer("127.0.0.1:9876")
                .setProducerGroup("ORDER_PRODUCER_GROUP"))
        .topicResolver(myTopicResolver)               // 必填
        .serializer(new Fastjson2EventSerializer())   // 可选,缺省即此
        .build();

eventManager.start();                                 // 受控启动,待依赖就绪后调用

// 发布事件
eventManager.publish(new OrderCancelledEvent("order-001"));

// 注册订阅者
eventManager.registerSubscriber("notify-customer", OrderCancelledEvent.class,
        event -> sendNotification(event.getEntityId()));

// 关闭
eventManager.shutdown();

组件能力:

成员类型说明
builder()静态唯一构造入口,返回 Builder
.config(...)必填RocketMqConfig
.topicResolver(...)必填ITopicResolver,缺则 build()NPE
.producer(...)可选外部注入 Remoting MQProducer,与 Spring 容器共享;未注入则框架自建(单实例复用)
.serializer(...) / .metrics(...) / .orderManager(...)可选缺省 Fastjson2EventSerializer / NoOpEventMetrics / SubscriberOrderManager
start()方法真正拉起 Producer/Consumer 收发;init(),先 buildstart
shutdown()方法释放 Consumer 与自建 Producer(外部注入的 Producer 不关闭)

2.3 gRPC 事件管理器

RocketMqGrpcEventManager 基于 rocketmq-client-java 的 Producer / PushConsumer,仅支持 5.x Broker(需开启 gRPC Proxy),4.x 不可用。

java
RocketMqGrpcEventManager eventManager = RocketMqGrpcEventManager.builder()
        .config(new RocketMqConfig().setProxyAddr("127.0.0.1:8081")
                .setProducerGroup("ORDER_PRODUCER_GROUP"))
        .topicResolver(myTopicResolver)               // 必填
        .build();

eventManager.start();
// 发布 / 订阅 / 关闭与 Remoting 完全一致
eventManager.shutdown();

与 Remoting 的差异(来自实现):

维度RemotinggRPC
依赖rocketmq-client(必引)rocketmq-client-java(optional)
地址nameServerproxyAddr
Producer/Consumer 创建new DefaultMQProducer/DefaultMQPushConsumerClientServiceProvider 工厂 build
延迟消息setDelayTimeLevel(级别)setDeliveryTimestamp(绝对时间戳),内部按级别映射到毫秒
消费确认ConsumeConcurrentlyStatusConsumeResult.SUCCESS/FAILURE
关闭shutdown()close()(封装在 shutdown() 内)

gRPC 默认延迟级别映射(与 RocketMQ 等级一致):1s / 5s / 10s / 30s / 1m / 2m / ... / 1h / 2h,共 18 级;defaultDelayLevel 越界时自动夹取到有效区间。

2.4 订阅者注册

订阅者注册方式与 core 本地管理器一致(基于 AbstractMQEventManager)。每个 topic 独立 Consumer,消费隔离。

java
// 基本注册
eventManager.registerSubscriber("notify-customer", OrderCancelledEvent.class,
        event -> sendNotification(event.getEntityId()));

// 带执行条件
eventManager.registerSubscriber("refund", OrderCancelledEvent.class,
        event -> doRefund(event.getEntityId()),
        event -> event.getCancelReason() != null);

// 带投递策略(延迟)
eventManager.registerSubscriber("delayed-notification", OrderCancelledEvent.class,
        event -> sendNotification(event.getEntityId()),
        DeliveryPolicy.DELAYED);

// 带依赖顺序(依赖 notify-customer 先执行)
eventManager.registerSubscriber("update-read-model", OrderCancelledEvent.class,
        event -> updateReadModel(event.getEntityId()),
        new DefaultExecuteCondition<>(),
        "notify-customer");

2.5 端到端示例

把配置、构造、发布、订阅拼成完整流程(Remoting 为例,gRPC 仅替换管理器与地址字段):

java
// 1. 配置
RocketMqConfig config = new RocketMqConfig()
        .setNameServer("127.0.0.1:9876")
        .setProducerGroup("ORDER_PRODUCER_GROUP")
        .setConsumerGroup("ORDER_CONSUMER_GROUP");

// 2. 构造(topicResolver 必填)
RocketMqEventManager eventManager = RocketMqEventManager.builder()
        .config(config)
        .topicResolver(myTopicResolver)
        .build();

// 3. 注册订阅者(应在 start 前完成)
eventManager.registerSubscriber("notify-customer", OrderCancelledEvent.class,
        event -> sendNotification(event.getEntityId()));

// 4. 应用依赖就绪后启动
eventManager.start();

// 5. 发布
eventManager.publish(new OrderCancelledEvent("order-001"));

// 6. 关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(eventManager::shutdown));

2.6 订阅者执行顺序与依赖传播

RocketMqEventManager 的订阅者执行顺序由 core 的 ISubscriberOrderManager依赖边图驱动,注册时以虚拟根 _root_ 为起点建立 事件 → 订阅者 → 后继订阅者 的边。其能力可概括为:同层无依赖的订阅者并行,跨层有依赖的订阅者顺序

注册时通过 registerSubscriber(..., dependsOn) 声明依赖(见 2.4);SubscriberOrderManager 在注册期即检测环依赖,成环则抛 IllegalStateException

能力示意

并行(同层多根,互不依赖,并发消费)OrderCancelledEvent 挂在 _root_ 下的两个订阅者,各自独立 topic,消息同时发出、并行执行。

text
        ┌─────────────┐
_root_ ─┤ notify-customer │  (topic: OrderCancelledEvent, 独立消息)
        ├─────────────┤
        │ refund        │  (topic: OrderCancelledEvent, 独立消息)
        └─────────────┘
        → notify-customer 与 refund 并行,无先后顺序

顺序(依赖链,前者完成才触发后者)update-read-model 声明 dependsOn("notify-customer"),形成 _root_ → notify-customer → update-read-model 的链。

text
_root_ → notify-customer ──(完成后发新消息)──▶ update-read-model
        (先执行)                               (后执行)

顺序链如何在 MQ 上串联

顺序不是进程内串行调用,而是跨 MQ 的消息重投:消费侧执行完当前订阅者 handleEvent 后,AbstractMQEventManager 通过 orderManager.findNextSubscribers(event, name) 取出直接后继,对每个后继递归 publish 一条新 MQ 消息,由后者的 Consumer 执行。即"前一个订阅者完成 → 后继作为新消息被投递 → 后继 Consumer 执行",逐层推进。

因此依赖顺序的语义是"最终顺序一致",而非强实时串行;每一跳都是独立的 MQ 投递,享受各自的重试/死信保障(见 3.6)。

3. 关键机制与避坑指南

3.1 Consumer Group 唯一性

⚠️ 重要约束consumerGroup 必须全局唯一,且禁止与 topic 同名,否则会导致 rebalance 抢队列、消息被错误消费。框架默认值为 PRAGMATIC_DDD_RMQ_CONSUMER,多实例部署时务必显式配置为各自唯一值。

3.2 外部注入 Producer 的生命周期

⚠️ 重要约束:通过 builder().producer(...) 注入的 Producer(Remoting 为 MQProducer,gRPC 为 Producer)由调用方持有,shutdown() 不会关闭它;仅框架自建的 Producer 才在 shutdown() 中被释放。与 Spring 容器共享 Producer 时,需自行管理其生命周期。

3.3 受控启动与 init 误区

⚠️ 重要约束:管理器构造(build)后不会建立任何网络连接;真正的收发由 start() 触发(gRPC 的 build 即连接,故推迟到 start)。不存在 init() 方法,调用会编译失败。应用应待全部下游依赖就绪后再调 start(),避免 Consumer 提前拉消息而下游未准备好。

3.4 死信队列格式

⚠️ 重要约束:消费重试超过 maxReconsumeTimes 后,框架将消息投递到死信队列,其 topic 格式为 原topic%DLQ%(代码实现:topic + "%DLQ%"),而非 %DLQ%{consumerGroup}。例如 topic 为 OrderCancelledEvent,死信 topic 为 OrderCancelledEvent%DLQ%。注意与原生 RocketMQ %DLQ%{consumerGroup} 约定不同,运维查死信时需按此格式。

3.5 可选依赖与协议选择

⚠️ 重要约束rocketmq-client-java(gRPC 客户端)标记为 optional=true;仅使用 Remoting 时不会引入,也不会触发 gRPC 类加载。RocketMqGrpcEventManager 仅在 classpath 存在 gRPC 依赖时可实例化。协议选择:4.x → 只能 Remoting;5.x 无 gRPC Proxy → Remoting;5.x + gRPC Proxy → 推荐 gRPC。

3.6 顺序链的幂等与循环依赖

⚠️ 重要约束dependsOn 声明的顺序链是跨 MQ 的消息重投实现(见 2.6),每一跳都是独立投递。因此:

  • 订阅逻辑必须幂等:同一事件可能因重试、重投被多次执行,非幂等操作(如重复扣款)需用业务键去重。
  • 禁止依赖成环:注册期 SubscriberOrderManager 会检测环依赖,成环立即抛 IllegalStateException(fail-fast),须在开发期修正依赖声明。

4. 异常与错误处理体系

本模块复用 core 的事件异常类型,不做独立异常体系:

阶段触发条件异常 / 行为
构造configtopicResolver 为 nullNullPointerExceptionbuild()requireNonNull
初始化 Consumersubscribe 失败RegisterDomainEventException(topic, cause)
发送Producer 发送失败PublishEventException(entityId, cause);同时 metrics.recordPublish(..., false, ...)
启动 Producer/Consumerstart() 失败RegisterDomainEventException(Remoting)/ RuntimeException(gRPC)
消费单条失败返回 RECONSUME_LATER / FAILURE框架重试;耗尽后 handleDeadLetter 投死信,metrics.recordDlq(...)

最佳实践:发布失败会向上抛 PublishEventException,业务代码应捕获并处理(或交由 Outbox 兜底,见第 5 节);消费失败由框架自动重试,订阅逻辑需保证幂等(同一事件可能重复投递)。

5. 与 Outbox 配合(可靠投递)

RocketMQ 事件管理器(即时推送通道)可与 core 的 OutboxRelay 配合,进一步加固"事件不丢":事件先随业务事务写入 Outbox 表,由 Relay 异步轮询补推,当即时投递(Producer 发送)失败时由 Outbox 兜底,避免事件因发送异常而丢失,保障最终一致性。

java
// 1. RocketMQ 事件管理器
RocketMqEventManager eventManager = RocketMqEventManager.builder()
        .config(config).topicResolver(myTopicResolver).build();
eventManager.start();

// 2. Outbox 兜底轮询
OutboxRelay relay = new OutboxRelay(
        outboxStore,
        eventManager,
        new Fastjson2EventSerializer(),
        Executors.newScheduledThreadPool(1),
        new OutboxRelayConfig(Duration.ofSeconds(5), Duration.ofSeconds(30),100, 5));
relay.start();

流程(事务内落库 + 主动推送,失败由 Relay 兜底):

text
业务事务提交 → Outbox 落库(PENDING) + EagerPublisher 主动推送
                                    ↓ 推送失败
                          OutboxRelay 兜底轮询 → 重新推送
                                    ↓ 重试耗尽
                                markFailed(死信)

6. 事件指标

实现 core 的 IEventMetrics 接口可采集发布/消费/死信指标,构造时通过 .metrics(...) 注入(缺省 NoOpEventMetrics):

java
public class MyEventMetrics implements IEventMetrics {
    @Override
    public void recordPublish(String topic, String eventType, boolean success, long latencyMs) {
        // 记录发布指标
    }

    @Override
    public void recordConsume(String topic, String eventType, boolean success, long reconsumeTimes) {
        // 记录消费指标
    }

    @Override
    public void recordDlq(String topic, String cause) {
        // 记录死信
    }
}

注意签名与 core IEventMetrics 一致:recordPublish(topic, eventType, success, latencyMs)recordConsume(topic, eventType, success, reconsumeTimes)recordDlq(topic, cause),均为 4 / 3 参数,与 NoOpEventMetrics 默认实现对齐。

7. 总结速查

概念关键事实最关键约束
RocketMqConfig链式 setter 或 bind(source)RocketMqConfiguration 为门面nameServer(Remoting)/ proxyAddr(gRPC)二选一必填
RocketMqEventManagerRemoting,4.x/5.xbuilder().config().topicResolver().build(),无 init(),用 start()
RocketMqGrpcEventManagergRPC,仅 5.x + Proxy同 builder 结构;Producer 类型为 gRPC Producer
topicResolverITopicResolver,两个管理器构造必填缺失 build() 抛 NPE
serializer / metrics缺省 Fastjson2EventSerializer / NoOpEventMetrics可注入自定义实现
外部 Producerbuilder().producer(...)shutdown() 不关闭外部注入的 Producer
consumerGroup默认 PRAGMATIC_DDD_RMQ_CONSUMER全局唯一,禁止与 topic 同名
死信 topic代码格式 topic%DLQ%%DLQ%{consumerGroup}
延迟消息DELAYED + defaultDelayLevelRemoting 用级别,gRPC 映射为绝对时间戳(18 级 1s~2h)
可靠投递OutboxRelay + outboxStore解决发送失败丢事件;订阅逻辑需幂等