DDK Event Starter
发布聚合根登记的领域事件;标了
@IntegrationEvent的事件交给 Spring Modulith,在业务事务里登记、提交后可靠地投递到进程外;消费端用IdempotentConsumer去重。
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 已有能力」的原则直接复用,只补两处缺口:
- 领域层怎样声明对外事件。Modulith 用自己的
@Externalized或 jMolecules 注解选择事件,领域层不能依赖它们。DDK 在com.ddk.core.domain提供纯 JDK 的@IntegrationEvent,由本 starter 翻译给 Modulith。 - 消费端去重。投递是「至少一次」,Modulith 只管发布侧。
集成事件
@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 管理:
<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:
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
</dependency>ddk:
event:
rocketmq:
name-server: localhost:9876sequenceDiagram
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-type | type(),默认取类的简单名 | 消费方按这个稳定名称分发,而不是按 Java 类;移动包、改类名不影响下游。改类名时把旧名称写进 type() |
ddk-event-version | version(),默认 1 | 消息体是哪个版本的结构 |
ddk-event-id | id() 指向的访问方法 | 去重,见下一节 |
演进规则:
同一版本内 只新增可选字段;消费方必须忽略不认识的字段
破坏性修改 删除、改名字段,或改变字段的类型、含义
→ 新建一个 type() 相同、version = N + 1 的事件类
→ 两个版本并行发布,直到所有消费方都切换过去
→ 再删除旧版本启动时,本 starter 扫描应用包里所有 @IntegrationEvent,声明有误就直接启动失败:key() 或 id() 不是无参访问方法、版本号小于 1、目标为空、两个类声明了相同的类型和版本。否则这些错误要到事务提交后的投递环节才暴露,那时发布记录已经停在 FAILED。
消费端幂等
@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")。
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.enabled | true | 注册基于 Spring 的 DomainEventPublisher |
ddk.event.inbox.enabled | false | 注册 IdempotentConsumer |
ddk.event.inbox.table | ddk_processed_message | 已处理消息表 |
ddk.event.inbox.initialize-schema | true | 启动时建表 |
ddk.event.rocketmq.name-server | NameServer 地址,多个用 ; 分隔;设置后由 DDK 创建 producer | |
ddk.event.rocketmq.producer-group | ddk-event-producer | 该 producer 的生产者组 |
ddk.event.rocketmq.send-timeout | 3s | 单次发送超时 |
spring.modulith.events.* | Spring Modulith | 发布记录表、重启重投、完成模式、过期监控 |
不适用的场景 / 已知问题
- 乱序:key 只保证同一分区(RocketMQ 是同一队列)内有序;跨分区、队列数变化或重投导致的乱序,需要消费方按聚合版本丢弃过期事件。
- 死信处理:失败的发布记录停在
FAILED,需要自己接告警,并通过 Modulith 的FailedEventPublications重新提交。