领域事件应在 domain 层发布,由聚合根或领域服务在状态变更后立即生成并调用内存态 publish 方法,不涉及序列化或网络调用;domain 层仅定义事件结构和 publisher 接口,具体实现由 application 层注入。

领域事件该在哪一层发布:domain 包内触发,不依赖 infrastructure
领域事件不是“发给消息队列的通知”,而是“业务事实已发生”的声明。它必须由 domain 层的聚合根或领域服务在状态变更后立即生成,比如 Order.Confirm() 成功执行后,调用 o.publish(OrderConfirmed{ID: o.ID, ...}) —— 这个 publish 方法只是把事件塞进一个内存通道或切片,**不涉及序列化、网络调用或 Kafka 客户端**。
常见错误是把 publish 写成直接调用 kafka.Producer.SendMessage,这会让 domain 层强依赖外部基础设施,破坏分层边界,也导致单元测试无法隔离。
- domain 层只定义事件 struct(如
OrderPaidEvent)和Publisher接口(如type Publisher interface { Publish(Event) }) - application 层实现该接口,注入具体的消息发送器(
KafkaPublisher或InMemoryPublisher) - 测试时用内存实现替换,断言事件类型和字段即可,无需启动 Kafka
消费端如何保证幂等:用 event.ID + 处理状态做唯一键
跨服务消费事件时,“一次处理”不等于“一次送达”。Kafka 可能重复投递,NATS JetStream 也可能重传,所以消费者必须自己判重。最可靠的方式是:以 event.ID(全局唯一,建议用 ULID 或 UUIDv7)为 key,在 DB 或 Redis 中记录“该事件是否已成功处理”。
不要用业务字段(如 order_id)单独做幂等键——同一订单可能触发多个事件(OrderCreated、OrderPaid),混在一起会误判。
- Redis 方案:用
SETNX event_id:xxx 1 EX 3600,成功则处理,失败则跳过 - DB 方案:建表
event_consumed(event_id TEXT PRIMARY KEY, processed_at TIMESTAMPTZ),INSERT ON CONFLICT DO NOTHING - 避免在事务中查 DB 判重后再发 DB 更新——若事务中途失败,状态未写入,下次又会重试;应先 INSERT 幂等记录,再执行业务逻辑
为什么不能在事务里直接 publish:DB 提交与事件发送无法原子提交
如果你在数据库事务内调用 kafka.Producer.SendMessage,会出现两种失败情况:DB 提交成功但消息发送失败(事件丢失),或 消息发送成功但 DB 回滚(事件幽灵)。Go 没有两阶段提交(2PC)支持,Kafka 事务也只能保证“写入 Kafka 和写入自身 topic 原子性”,无法涵盖你的 MySQL/PostgreSQL。
在 Go 中使用 google/wire 实现编译时依赖注入——wire.NewSet、wire.Build、wire.Bind(接口→实现)、wire.Struct、wire.Value、wire.Interface
正确做法是“本地消息表 + 定时扫描”或“Kafka 事务 + 业务表同属一个 Kafka cluster”。前者更通用:
- 在业务 DB 同一事务中,插入一条
outbox_events记录(含事件 payload、topic、status=‘pending’) - 另起一个独立 goroutine,定时 SELECT WHERE status = 'pending',逐条发送并更新 status = 'sent'
- 失败时标记为 'failed',人工介入或自动告警
用 github.com/ThreeDotsLabs/watermill 时 consumer group 必须显式设置
Watermill 默认不设 ConsumerGroup,若你启动多个服务实例监听同一 topic,所有实例都会收到每条消息——这不是负载均衡,是广播风暴。Kafka 的 partition-level 负载均衡依赖 consumer group 协调,而 RabbitMQ 的 queue 绑定也需明确指定 queueName 作为 group 标识。
典型错误配置:
sub, _ := kafka.NewSubscriber(kafka.SubscriberConfig{
Brokers: []string{"localhost:9092"},
Unmarshaler: kafka.DefaultMarshaler{},
})
// ❌ 缺少 ConsumerGroup 字段,所有实例竞争消费
正确写法:
sub, _ := kafka.NewSubscriber(kafka.SubscriberConfig{
Brokers: []string{"localhost:9092"},
ConsumerGroup: "inventory-service-group", // ✅ 必须设
Unmarshaler: kafka.DefaultMarshaler{},
})
另外,Ack 必须在业务逻辑彻底完成(DB 写入、缓存更新、下游通知全部成功)后才调用,否则消息被提前确认,失败后无法重试。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










