← 返回博客
go2026-09-12 22:13:021 分钟 · 439 0

事件驱动架构:用 NATS/Kafka 解耦服务

用消息队列(NATS/Kafka)让服务通过事件通信而非直接调用,讲清发布订阅、事件建模、幂等消费和最终一致。

1. 为什么事件驱动

订单服务创建完直接调库存、通知、积分三个服务,任一个挂订单就失败。改成发一个 OrderCreated 事件,三个服务各自订阅自己关心的事,互不知道对方存在,挂了也不影响下单。

2. 发布事件

func (h *OrderHandler) Create(c *gin.Context) {
    order := h.svc.Create(c)
    evt, _ := json.Marshal(OrderCreated{ID: order.ID, UID: order.UID, Amount: order.Amount})
    js.Publish("order.created", evt) // NATS JetStream
    OK(c, order)
}

下单只管落库和发事件,剩下的交给订阅方。

3. 订阅消费

js.Subscribe("order.created", func(m *nats.Msg) {
    var e OrderCreated
    json.Unmarshal(m.Data, &e)
    pointsSvc.Add(e.UID, e.Amount) // 加积分
    m.Ack()
})

库存、通知、积分各起一个 consumer group,互不影响。Kafka 用 consumer group 做并行消费和重放。

4. 事件建模

事件名用「领域.发生的事」过去式:order.createduser.registered。载荷带足够下游用的字段,别让下游反查你。版本化事件,字段只加不删。

5. 幂等消费

消息可能重复投递(网络重试),消费方必须幂等:

func handle(e OrderCreated) error {
    if dedup.Seen(e.ID) {
        return nil // 已处理过
    }
    // 处理业务
}

用事件 id 或业务 id 去重,保证重复消费不产生重复积分/重复扣库存。

6. 最终一致

事件驱动下数据不立即可见(积分晚几秒到账),这是最终一致。前端要接受这个延迟,或用本地消息表先落库再发,保证至少一次投递。

7. 上线清单

事件名用「域.过去式」,载荷自包含。消费必须幂等,靠事件 id 去重。重要事件用 JetStream/Kafka 持久化,别用会丢的纯内存队列。发布失败要重试,本地消息表保证不丢。监控消费延迟和死信,堆积说明下游处理不过来。

下一篇讲 DDD,把业务逻辑从 handler 里长出来。

相关推荐

本文为原创文章,采用CC BY-NC-SA 4.0协议授权,转载请保留署名与原文链接。原文链接:https://www.wxbuluo.com/article/201