Skip to content

DDK Event Starter ​

发布聚合根登记的领域事件;标了 @IntegrationEvent 的事件交给 Spring Modulith,在业务事务里登记、提交后可靠地投递到进程外;消费端用 IdempotentConsumer 去重。

text
aggregate.registerEvent(e) ──► 仓储写入 ──► DomainEventPublisher ──► Spring ApplicationEventPublisher
                                                                        │
          进程内   @TransactionalEventListener(AFTER_COMMIT) ◄───────────┤
          进程外   @IntegrationEvent + Spring Modulith     ◄───────────┘
                   在业务事务中写入 event_publication,提交后投递到 Kafka / AMQP / JMS / MessageChannel

为什么复用 Spring Modulith ​

事务性 Outbox 的发布侧,Spring Modulith 已经做好了:事件发布记录在业务事务里写入,提交后投递,失败的记录停在 FAILED 可以重新提交,还有过期监控和重启重投。验证确认它的 JDBC 实现与 MyBatis-Plus 共用同一个 DataSource 和事务:提交后有记录,回滚后没有。DDK 按「不重复实现 Spring Modulith 已有能力」的原则直接复用,只补两处缺口:

  1. 领域层怎样声明对外事件。Modulith 用自己的 @Externalized 或 jMolecules 注解选择事件,领域层不能依赖它们。DDK 在 com.ddk.core.domain 提供纯 JDK 的 @IntegrationEvent,由本 starter 翻译给 Modulith。
  2. 消费端去重。投递是「至少一次」,Modulith 只管发布侧。

集成事件 ​

java
@IntegrationEvent(value = "user-events", key = "userId", id = "eventId")
public record UserRegisteredEvent(UUID eventId, UserId userId, String username, Instant occurredOn) implements DomainEvent {
}
属性含义
value投递目标:Kafka topic、AMQP exchange、JMS 目的地或 MessageChannel Bean 名,取决于引入的外发模块
key作为消息 key 的访问方法名。类型化标识取原始值;同一个 key 的消息在 Kafka 里保持顺序
id事件自带标识的访问方法名,取值写进消息头 ddk-event-id。必须是事件自身的数据,重投前后才保持不变
type / version契约名称与版本,见下一节

应用再引入 Modulith 的 JDBC 发布记录和一个外发模块,版本由 DDK 的 BOM 管理:

xml
<dependency>
    <groupId>org.springframework.modulith</groupId>
    <artifactId>spring-modulith-starter-jdbc</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.modulith</groupId>
    <artifactId>spring-modulith-events-kafka</artifactId>
</dependency>

消息体里的类型化标识写成原始值("userId":42),由本 starter 注册的 IdentifierJacksonModule 处理。应用自己声明了 EventExternalizationConfiguration 时,DDK 的选择规则让位。

投递到 RocketMQ ​

Spring Modulith 没有 RocketMQ 模块,由本 starter 补上。引入 RocketMQ 客户端(版本由 DDK 的 BOM 管理)代替 Modulith 的外发模块,再配置 NameServer:

xml
<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
</dependency>
yaml
ddk:
  event:
    rocketmq:
      name-server: localhost:9876
mermaid
sequenceDiagram
    participant App as 应用服务
    participant DB as event_publication
    participant L as DDK RocketMQ 监听器
    participant MQ as RocketMQ
    App->>DB: 在业务事务里登记事件
    Note over App,DB: 提交
    DB-->>L: 提交之后
    L->>MQ: 发到 topic:tag,按 key 哈希选队列
    MQ-->>L: SEND_OK
    L->>DB: 标记发布完成
方面行为
目标value 写成 topic 或 topic:tag,与 rocketmq-spring 的约定一致
key有 key 的消息按 key 哈希选队列,同一个聚合的事件保持顺序;key 同时写进消息的 keys,方便在控制台按业务标识查消息
消息头ddk-event-type、ddk-event-version、ddk-event-id 写成用户属性,消费方用 MessageExt.getUserProperty(...) 读取
消息体用应用的 JsonMapper 序列化成 JSON;String 和 byte[] 原样发送
失败发送结果不是 SEND_OK 就算失败,发布记录保持未完成,可以重投
producer应用已有 DefaultMQProducer(例如 rocketmq-spring 注册的那个)时直接用它,否则按 ddk.event.rocketmq.* 创建

它工作在 Modulith 默认的监听器模式下;spring.modulith.events.externalization.enabled=false 或 mode=outbox 时不注册。测试用 Testcontainers 起真实的 RocketMQ broker。

事件契约 ​

每条投递出去的集成事件都带契约消息头:

消息头来源用途
ddk-event-typetype(),默认取类的简单名消费方按这个稳定名称分发,而不是按 Java 类;移动包、改类名不影响下游。改类名时把旧名称写进 type()
ddk-event-versionversion(),默认 1消息体是哪个版本的结构
ddk-event-idid() 指向的访问方法去重,见下一节

演进规则:

text
同一版本内     只新增可选字段;消费方必须忽略不认识的字段
破坏性修改     删除、改名字段,或改变字段的类型、含义
              → 新建一个 type() 相同、version = N + 1 的事件类
              → 两个版本并行发布,直到所有消费方都切换过去
              → 再删除旧版本

启动时,本 starter 扫描应用包里所有 @IntegrationEvent,声明有误就直接启动失败:key() 或 id() 不是无参访问方法、版本号小于 1、目标为空、两个类声明了相同的类型和版本。否则这些错误要到事务提交后的投递环节才暴露,那时发布记录已经停在 FAILED。

消费端幂等 ​

java
@KafkaListener(topics = "user-events")
void on(UserRegisteredEvent event, @Header("ddk-event-id") String eventId) {
    idempotentConsumer.handle("welcome-mail", eventId, () -> mailService.sendWelcome(event.userId()));
}

用 RocketMQ 时,事件 ID 从用户属性里取:message.getUserProperty("ddk-event-id")。

text
handle(consumer, messageId, handler)            加入调用方的事务,没有则新开
  保存点:INSERT (consumer, message_id)          主键冲突 → 回滚到保存点,返回 false
  handler.run()                                  抛异常 → 整个事务连同登记一起回滚
  return true
  • 登记与处理逻辑的写入在同一个事务里:处理失败时消息保持未处理,重投后会再次执行。
  • 登记放在保存点之后:PostgreSQL 在语句失败后会中止整个事务,重复消息只回滚到保存点,调用方的其他写入照常提交。这一点有基于真实 PostgreSQL 的测试。
  • purgeOlderThan(Duration) 清理旧记录,保留时长要超过中间件可能重投的最长时间。

开启 ddk.event.inbox.enabled=true 后,启动时用 CREATE TABLE IF NOT EXISTS 创建 ddk_processed_message,适用于 MySQL、PostgreSQL、H2;其他数据库或由迁移工具管理表结构时,关闭 initialize-schema 并自行建表。

配置项 ​

配置项默认值说明
ddk.event.enabledtrue注册基于 Spring 的 DomainEventPublisher
ddk.event.inbox.enabledfalse注册 IdempotentConsumer
ddk.event.inbox.tableddk_processed_message已处理消息表
ddk.event.inbox.initialize-schematrue启动时建表
ddk.event.rocketmq.name-serverNameServer 地址,多个用 ; 分隔;设置后由 DDK 创建 producer
ddk.event.rocketmq.producer-groupddk-event-producer该 producer 的生产者组
ddk.event.rocketmq.send-timeout3s单次发送超时
spring.modulith.events.*Spring Modulith发布记录表、重启重投、完成模式、过期监控

不适用的场景 / 已知问题 ​

  1. 乱序:key 只保证同一分区(RocketMQ 是同一队列)内有序;跨分区、队列数变化或重投导致的乱序,需要消费方按聚合版本丢弃过期事件。
  2. 死信处理:失败的发布记录停在 FAILED,需要自己接告警,并通过 Modulith 的 FailedEventPublications 重新提交。

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权