
组里做 DDD 架构改造代码评审的时候几乎每个项目都会有人问同一个问题MQ 消息到底放到哪一层处理问法各不相同——消费 Listener 算 Interface 层还是 Infrastructure 层消费到的业务逻辑能不能写在领域服务里消息体直接往聚合根传行不行坦白讲这个问题很多人背过标准答案但一到动手写就乱。原因在于一条 MQ 消息从进来到处理完本身就跨越了好几个层次网络连接是基础设施字节反序列化是基础设施里面承载的业务语义要进入领域而编排这些业务动作的又应该是应用层。所以放哪一层这个问法本身就不完整更准确的问题是消息处理的每个阶段分别该属于哪一层的职责。这篇文章就把这件事彻底拆开分别讲清消费端、生产端、双写一致性、幂等重试最后给一套可以直接抄走的模块结构和代码示例。1. 这个问题真正在问什么不是包结构是依赖边界很多人看到MQ 放哪一层第一反应是去看包的目录结构想把 MQ 相关类塞进某个包里就算完事。但 DDD 的分层本质上是一组依赖规则外层可以依赖内层内层绝不能反向依赖外层。这个规则和你的代码放在infrastructure还是adapter包里关系不大重要的是类与类之间的依赖箭头指向哪边。1.1 四层架构职责速览先把经典的 DDD 四层职责对一下表避免后面越聊越偏。层核心职责与 MQ 的关系Interface 层协议翻译REST、RPC、WebSocket、消息入口可以作为消息入口适配器但不写业务Application 层用例编排、事务边界、权限校验、事件分发定义消息业务处理的入口定义发布/消费所需的端口接口不感知具体中间件Domain 层业务规则、状态流转、聚合、领域事件完全不感知 MQ只表达业务事实Infrastructure 层数据库实现、MQ 客户端、序列化、定时任务、文件存储持有真实 MQ 连接、Listener、Codec、Outbox 轮询实现这里最容易被忽略的是Application 层不感知具体中间件这句。应用层需要的是一个把事件发出去的能力它应该依赖一个自己定义的接口比如IntegrationEventPublisher而不是直接依赖RabbitTemplate或者 Kafka 的Producer。至于这个接口背后是 RabbitMQ、RocketMQ 还是 Kafka那是基础设施层该操心的事。1.2 MQ 在架构里的真实身份MQ 和数据库一样属于基础设施服务。消息本身是一个运输工具就像快递包裹包裹箱、快递单、运输车辆是基础设施包裹里装的商品才是业务关心的东西。你不能因为商品装了箱就把整个快递体系搬进业务代码里。有了这个认知再回头看MQ 放哪一层结论就出来了消息的接收动作监听、连接、确认、反序列化放在基础设施层。消息触发的业务处理必须回到应用层由应用服务调用领域模型完成。领域层只负责产生业务事实通过领域事件表达出来不关心这个事实是通过什么渠道通知出去。同样的依赖必须向内思想在嵌入式系统的模块划分里也能看到只不过服务端落地时具体表现为外层适配器 内层端口的模式。接下来分别拆消费端和生产端这两个方向的处理逻辑差异很大不能混在一起谈。2. 消费端MQ 接收是技术活业务处理必须回到应用层消费端是最容易写错的地方。很多人图省事在RabbitListener或者KafkaListener标注的方法里直接写业务代码一个方法干完所有事。短时间看是快了后面改起来想哭都来不及。2.1 一条消息从网卡到业务方法的完整路径一条消息到达消费进程后至少要经过这几个阶段连接与投递MQ broker 把消息投递给消费进程涉及消费组、确认机制、重投策略。反序列化把字节流解析成结构化 DTO做格式校验。语义转换把外部系统定义的协议转换成当前限界上下文能理解的结构这一步在跨系统对接时尤其重要。业务编排调用应用服务执行具体用例操作领域模型。确认与重试消费成功后确认消息失败按策略重试或者进入死信队列。这五个阶段里1、2、5 是纯基础设施工作4 是应用层的活3 属于防腐层Anti-Corruption Layer的范畴一般放在基础设施层的适配器里实现。这条路径拆清楚之后消费逻辑放哪一层其实就变成了你打算让哪个阶段留在 Listener 里。2.2 推荐的消费端结构Listener 适配 AppService 编排先看一段我推荐的代码骨架。消息监听类放在基础设施层只做解码和转发Component public class OrderMessageConsumer { private final OrderAppService orderAppService; private final OrderCreatedMessageCodec codec; RabbitListener(queues ${mq.order-created.queue}) public void onOrderCreated(byte[] body) { // 基础设施层负责把字节流解码成应用层能理解的 DTO OrderCreatedMessage message codec.decode(body); // 业务处理全权交给应用层 orderAppService.handleOrderCreated(message); } }应用服务里才是真正的业务编排Service public class OrderAppService { private final OrderRepository orderRepository; Transactional public void handleOrderCreated(OrderCreatedMessage message) { // 幂等判断详见第 5 节 if (orderRepository.existsByOrderNo(message.orderNo())) { return; } Order order Order.create(message.orderNo(), message.customerId()); orderRepository.save(order); } }注意这里OrderAppService接收的OrderCreatedMessage是应用层 DTO不是基础设施层的消息对象。如果团队要求严格更规范的做法是让 Listener 把消息翻译成语义完整的命令对象比如CreateOrderCommand再调应用层。这样应用层完全不知道消息的传输格式将来 MQ 从 RabbitMQ 换成 Kafka应用层一行不用改。还有一种常见做法是把消息监听类放在 Interface 层因为它本质上是外部输入适配器。这属于团队规范的选择两种都行关键判断标准只有一条Listener 里不能出现业务逻辑。一旦 Listener 开始直接操作仓储、直接改状态机就说明业务逻辑正在漏出应用层。2.3 为什么不能把业务逻辑写进 Listener我在评审里经常看到这样的写法KafkaListener(topics order-paid) public void onMessage(String json) { Order order JsonUtil.parse(json, Order.class); order.setStatus(PAID); orderRepository.save(order); }这段代码的问题非常典型消息中间件绑架了业务。哪天消费组要调整、topic 要换名、或者中间件本身要替换业务逻辑就要跟着搬家。无法直接测试。业务逻辑必须依赖 MQ 环境才能跑而它本身和 MQ 没有任何关系。复用性为零。如果这个业务同时能被 REST 接口触发同样的逻辑就要在 Controller 里再写一遍两份代码迟早出现偏差。事务边界失控。Listener 方法里的Transactional和业务事务纠缠在一起领域事件收集、聚合持久化、幂等记录的顺序很容易搞乱。正确的心智模型是Listener 只是适配器不是业务体。业务入口永远收敛在应用层消息、REST、RPC、定时任务都只是不同入口通道进到应用层之后走同一条业务路径。这样才能保证业务规则单一、可测、可复用。3. 生产端领域事件和集成事件要分开触发放应用层聊完消费端再看出厂方向。消息生产放哪一层是另一个高频争论点。这个问题的核心不是在哪发消息而是什么时候、由谁决定要发消息。3.1 两种事件的本质区别生产端必须区分两类事件混用是很多架构腐化的起点。维度领域事件集成事件表达语言领域通用语言业务术语跨系统协议稳定契约可见范围限界上下文内部跨上下文、跨服务使用方领域层内部逻辑、应用层编排外部系统消费者载体领域对象内部的事件对象独立的版本化 DTO例子OrderConfirmedDomainEventOrderConfirmedIntegrationEvent领域事件描述的是聚合里发生了一件业务事实它应该产生在领域层由聚合根在状态变更时记录。集成事件描述的是我要告诉别的系统一个消息它属于应用层的对外契约由应用服务负责翻译和发布。3.2 领域层只产生事实应用层负责发出通知聚合根在状态变化时只负责产生领域事件不负责发送public class Order extends AggregateRoot { private OrderStatus status; public void confirm() { if (this.status ! OrderStatus.PAID) { throw new OrderStateException(只有已支付订单才能确认); } this.status OrderStatus.CONFIRMED; // 领域事件只表达业务结果不关心通知方式 this.addDomainEvent(new OrderConfirmedDomainEvent(this.id, this.customerId)); } }应用层在事务边界内把这个领域事件翻译成集成事件再通过端口发布Service public class OrderAppService { private final OrderRepository orderRepository; private final IntegrationEventPublisher eventPublisher; Transactional public void confirmOrder(ConfirmOrderCommand cmd) { Order order orderRepository.find(cmd.orderId()); order.confirm(); orderRepository.save(order); // 从聚合根取出领域事件翻译成集成事件并发布 for (DomainEvent de : order.pullDomainEvents()) { if (de instanceof OrderConfirmedDomainEvent oce) { eventPublisher.publish(OrderConfirmedIntegrationEvent.of(oce)); } } } }这里有一个关键动作领域对象被保存后应用层统一收集领域事件并转换。很多团队的漏洞恰恰在转换这一步他们喜欢直接在聚合根里拼一个消息对象发给 MQ结果领域层被硬生生塞进了序列化逻辑。记住领域层的聚合根只告诉世界我发生了什么至于世界怎么知道不是它的事。3.3 发布器接口放哪用端口反解应用层需要一个发布事件的接口这个接口应该由应用层自己定义称为端口Portpublic interface IntegrationEventPublisher { void publish(IntegrationEvent event); }基础设施层提供适配器实现Component public class RabbitIntegrationEventPublisher implements IntegrationEventPublisher { private final RabbitTemplate rabbitTemplate; Override public void publish(IntegrationEvent event) { rabbitTemplate.convertAndSend(event.topic(), event.serialize()); } }这样一来应用层依赖的是接口不是 RabbitMQ将来的实现换成RocketIntegrationEventPublisher、KafkaIntegrationEventPublisher对应用层毫无影响。这就是端口-适配器模式也是 DDD 依赖倒置在消息场景下的具体落地。有人会问端口能不能定义在领域层我的建议是看语义。集成事件属于跨上下文协作协议不属于领域模型的语言放进领域层会污染领域概念。领域层最多定义一个事件记录器之类的存储端口用于把领域事件持久化到 Outbox而对外发布集成事件的端口留在应用层更干净。4. 先写库还是先发消息Outbox 解决双写不一致生产端绕不开一个经典问题业务数据写库和 MQ 发送是两个系统怎么保证一致性这就是先写数据库、先发布 MQ、补发消息这些讨论背后的原始诉求。4.1 双写问题的本质假设订单确认时要更新数据库状态同时发一条消息给下游仓储系统。数据库事务和 MQ 发送无法放在同一个原子操作里于是出现两种选择先写库后发 MQ数据库提交成功之后MQ 发送失败下游永远不知道订单已确认。先发 MQ后写库MQ 发出去了数据库写入失败下游收到一条不存在订单的虚假消息。有人尝试用分布式事务XA把数据库和 MQ 包进同一个全局事务但实现成本高、性能差很多中间件还不原生支持。更实用的思路是放弃跨系统原子性接受最终一致用 Outbox 模式保证消息不会丢。4.2 Transactional Outbox 实现细节Transactional Outbox 的核心思想是把要发消息当成一条业务数据和普通业务写入放进同一个本地事务。业务提交之后异步进程把待发送事件捞出来发给 MQ发成功之后再标记状态。先看表结构CREATE TABLE outbox_event ( id BIGINT AUTO_INCREMENT PRIMARY KEY, aggregate_type VARCHAR(64) NOT NULL, aggregate_id VARCHAR(64) NOT NULL, event_type VARCHAR(128) NOT NULL, payload JSON NOT NULL, status TINYINT NOT NULL DEFAULT 0, -- 0 待发送 1 已发送 2 发送失败 retry_count INT NOT NULL DEFAULT 0, created_at DATETIME NOT NULL, sent_at DATETIME NULL, INDEX idx_status_created (status, created_at) ) ENGINEInnoDB;业务方法里做两件事更新业务状态同时插入一条 outbox 记录。这两个动作在同一个事务里Transactional public void confirmOrder(ConfirmOrderCommand cmd) { Order order orderRepository.find(cmd.orderId()); order.confirm(); orderRepository.save(order); outboxRepository.save(new OutboxRecord( OrderConfirmedIntegrationEvent.of(order))); }事务一旦提交业务数据和待发消息就同时落库了不会再出现库更新了但消息没记录的情况。至于线程中某个进程挂了消息还留在 outbox 表里由另一个轮询器负责补偿发送。轮询发送的过程要特别注意并发。多个实例同时扫表会重复发送最好用数据库的行锁能力限定每次只能捞走一部分SELECT * FROM outbox_event WHERE status 0 ORDER BY id LIMIT 50 FOR UPDATE SKIP LOCKED;拿到记录后发送 MQ成功后把status更新为 1或者直接删除。发送失败则累加retry_count超过阈值标记为失败并告警。4.3 轻量替代AFTER_COMMIT 事件发布的适用场景如果 Outbox 表不想建还有一种轻量做法利用 Spring 的事务提交回调。事务成功提交后触发监听方法进行发布TransactionalEventListener(phase TransactionPhase.AFTER_COMMIT, fallbackExecution true) public void onOrderConfirmed(OrderConfirmedDomainEvent event) { integrationEventPublisher.publish(OrderConfirmedIntegrationEvent.of(event)); }优点是零额外存储、延迟几乎为零。代价是事务提交之后、回调执行之前应用进程如果崩溃事件仍然会丢而且回调里发布失败也没有事务可以回滚。严格来说它省掉了消息落库这层保险可靠性低于 Outbox。我的建议是内部系统间的非关键通知比如用户行为统计、异步刷新缓存用 AFTER_COMMIT 完全可以接受但凡是下游要做库存扣减、支付确认、状态对账的老老实实上 Outbox。4.4 方案对比与选型建议方案一致性消息丢失概率复杂度适用场景先写库后发 MQ不一致窗口明显较高低内部通知可接受丢失先发 MQ 后写库存在虚假消息中低几乎不用仅可补偿场景AFTER_COMMIT 回调基本一致低低非关键事件、低延迟要求本地消息表 定时扫描最终一致极低中中小规模秒级延迟可接受Outbox 轮询最终一致极低中大多数业务系统的默认选择Outbox CDCBinlog最终一致极低较高高吞吐、毫秒级延迟要求如果你只是问先写数据库还是先发 MQ我的答案是两个单独做都不对把要发的消息写进和业务同一个数据库事务再用异步进程补发才是可靠路线。这也是补发消息热搜词背后真正要解决的问题。5. 消费幂等和重试陷阱消息进业务前的最后一道关消息能可靠发出不等于消费端能正确收到一次。MQ 的投递语义绝大多数是 at-least-once同一消息可能因为消费超时、进程重启、ACK 丢失等原因被投递多次。所以消费端必须幂等这是比放哪一层更优先的基础问题。5.1 at-least-once 交付与幂等的对应关系Kafka 默认配置下消息可能重复RabbitMQ 手动确认时消费端在处理完成但确认前崩溃消息也会被重新投递重试机制更会主动造成重复。这些和代码写得好不好无关是分布式的常态。幂等设计的目标是同一个业务请求被处理多次最终状态和只处理一次完全一致。5.2 幂等方案的落地代码最可靠的幂等做法是唯一键 唯一索引 同事务占位。先看代码Transactional public void handlePaymentSucceeded(PaymentSucceededMessage message) { try { processedRecordMapper.insertIfAbsent( new ProcessedMessage(message.getMessageId())); } catch (DuplicateKeyException e) { // 消息已处理过直接跳过 return; } Order order orderRepository.findByOrderNo(message.orderNo()); order.markPaid(message.getPaidAmount()); orderRepository.save(order); }processed_record表对message_id建唯一索引重复消息插入时直接抛唯一键冲突业务逻辑不用执行第二次。这里有个很关键的细节占位记录和业务更新必须在同一个事务里。如果先插入占位记录再执行业务业务失败时整个事务回滚占位记录也跟着回滚下次消息还能重新进来处理不会出现占位了但不处理的死结。除了消息 ID 去重业务侧的状态机校验同样重要。比如订单已经处于PAID状态再来一条支付成功消息领域方法直接抛状态异常或者在应用层前置判断只有待支付订单才能进入已支付状态。两层幂等叠加覆盖的场景更全面。5.3 重试次数、死信和手工补偿消费失败时不要无限重试。对于临时性错误数据库锁冲突、下游接口超时可以用指数退避重试几次对于永久性错误消息格式损坏、业务状态非法重试多少次都不会成功反而会阻塞队列。常规设计是约定最大重试次数超过次数进入死信队列并触发告警由人工或补偿脚本处理。这里要强调一点重试机制属于基础设施层的可靠性设计但重试后执行的业务逻辑仍然要回到应用层的统一入口。比如 Kafka 的重试子 topicorder-paid-retry它的 Listener 最终还是要调用OrderAppService.handlePaymentSucceeded而不是在重试处理器里重新拼一段业务代码。这样一致性入口不变重试只是换了个触发时机。还有一个小建议给每条集成事件带上messageId、occurredAt、source等公共字段统一打印日志。排查问题时能顺着 messageId 把生产端、MQ、消费端串起来否则多服务联调的时候消息到底在哪一环丢的根本定位不到。6. 一个可直接参考的订单模块落地方案说了这么多原则给出一套可以直接抄走的模块结构。以订单确认、发集成事件、下游消费为例。6.1 包结构与职责对照orderservice ├── interfaces │ └── rest │ └── OrderController ├── application │ ├── service │ │ └── OrderAppService │ ├── dto │ │ └── OrderCreatedMessage │ ├── command │ │ └── ConfirmOrderCommand │ └── port │ ├── IntegrationEventPublisher │ └── OutboxRepository ├── domain │ ├── model │ │ └── order │ │ ├── Order │ │ ├── OrderStatus │ │ └── OrderConfirmedDomainEvent │ └── repository │ └── OrderRepository └── infrastructure ├── mq │ ├── RabbitIntegrationEventPublisher │ ├── OrderCreatedMessageCodec │ └── OrderMessageConsumer ├── repository │ ├── OrderRepositoryImpl │ └── OutboxRepositoryImpl ├── outbox │ ├── OutboxEvent │ └── OutboxPoller └── config这个结构里interfaces只放 REST 入口MQ 消费入口放在infrastructure.mq因为消费本质上是一个基础技术适配动作。如果团队更愿意把所有外部输入都叫驱动适配器把 MQ 消费类挪到interfaces下也不冲突只要确保它们只依赖应用层接口即可。6.2 关键代码串起来跑一遍完整链路如下用户调用POST /orders/{id}/confirmOrderController把请求转成ConfirmOrderCommand调用OrderAppService.confirmOrder。应用服务加载聚合根执行order.confirm()聚合根状态变更并产生OrderConfirmedDomainEvent。事务提交前应用服务保存订单同时把集成事件写入 outbox 表。事务提交后OutboxPoller定时扫描待发送记录通过RabbitIntegrationEventPublisher发布到 MQ。下游仓储服务消费消息由它的OrderMessageConsumer解码消息再调用它自己的应用服务处理入库。这一条链路里每个阶段都清楚对应一个层次入口归接口层编排归应用层业务规则归领域层技术细节归基础设施层。没有任何一段业务逻辑被丢在 Listener 或者 Outbox 轮询器里。6.3 演进中需要守住的边界项目演进过程中最容易破坏边界的是消息协议升级。集成事件加字段时一定要带version字段并在消费端做兼容检查不能因为下游要一个新字段就直接改领域模型来迎合消息内容。如果消息里的某个字段和领域概念对不上说明语义转换这层防腐没做好问题出在基础设施层的翻译不怪领域模型。另一种常见的腐化是为了消息消费方便在应用服务里暴露了好几个几乎一样的处理方法。这时候停下来想想是不是应用层正在被不同的输入协议牵着走。应用层应该面向用例设计而不是面向消息设计。消息只是触发用例的一种方式应用层的方法签名应该是用例语义而不是消息字段。7. 实际评审中我常拦下的几种写法最后聊聊我在代码评审里见过最多的几种问题写法以及为什么拦下来。7.1 在 Listener 里直接写业务逻辑最典型的一段 DDD 改造前的代码我几乎在每个项目里都见过KafkaListener(topics order-paid) public void onMessage(String json) { Order order JsonUtil.parse(json, Order.class); order.setStatus(PAID); orderRepository.save(order); }这段代码的每一个动作都有问题直接在监听器里解析领域模型、直接改状态、直接操作仓储。它完全绕过了应用层导致幂等逻辑、领域事件、事务边界全部缺失。我通常要求改成调应用服务解析只负责提取必要字段最终的状态变更必须回到领域方法里去。7.2 在 Domain 层注入 MQ 客户端比在 Listener 里写业务更隐蔽的是在聚合根里直接注入 MQ 客户端public class Order { Autowired private RabbitTemplate rabbitTemplate; // 反面教材 public void confirm() { // 业务规则... rabbitTemplate.convertAndSend(order-confirmed, this); } }这种写法的直接后果是领域对象变成了技术类无法脱离 Spring 和 MQ 做单元测试领域模型对基础设施的依赖让整个架构的依赖方向彻底反转。真正的领域方法应该只记录事件发布动作放到应用层。我在评审里的判断标准很简单单测一下这个聚合根如果测试代码里需要 mock RabbitTemplate这个设计基本就是错的。7.3 事务注解和 ACK 模式搭配错位还有一个不太起眼但线上事故频发的地方消费监听方法和事务注解的组合。如果消费端使用自动 ACK同时监听方法里加了Transactional那么方法返回后事务可能还没提交消息却已经确认了。此时进程一旦崩溃消息就永久丢失。正确做法是使用手动确认模式确认逻辑放在事务成功提交之后执行或者干脆用第 4 节的 Outbox 兜底让消息确认和业务状态解耦。这类问题平时测不出来只有流量高峰或者故障演练时才会暴露。我处理的办法是给团队的消费端代码定一条约束消费监听方法里不写业务、不控制事务只做解码 投递应用层事务边界统一交给应用服务控制。这条约束写进开发规范比评审时一次次口头提醒有效得多。我自己在实际项目里见过太多次因为消息分层不清导致的线上问题从Listener 改了个字段导致线上订单状态错乱到领域对象里注入 MQ 导致单测全挂。回过头看所有问题都能用一条最基本的依赖规则解释清楚外层适配器负责所有技术细节应用层统一编排领域层保持纯净。下次再有人问你 MQ 放哪一层你可以反问一句你是说 MQ 客户端放哪还是说消息里那部分业务逻辑放哪把两个问题拆开答案自然就出来了。