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.created、user.registered。载荷带足够下游用的字段,别让下游反查你。版本化事件,字段只加不删。
5. 幂等消费
消息可能重复投递(网络重试),消费方必须幂等:
func handle(e OrderCreated) error { if dedup.Seen(e.ID) { return nil // 已处理过 } // 处理业务 }
用事件 id 或业务 id 去重,保证重复消费不产生重复积分/重复扣库存。
6. 最终一致
事件驱动下数据不立即可见(积分晚几秒到账),这是最终一致。前端要接受这个延迟,或用本地消息表先落库再发,保证至少一次投递。
7. 上线清单
事件名用「域.过去式」,载荷自包含。消费必须幂等,靠事件 id 去重。重要事件用 JetStream/Kafka 持久化,别用会丢的纯内存队列。发布失败要重试,本地消息表保证不丢。监控消费延迟和死信,堆积说明下游处理不过来。
下一篇讲 DDD,把业务逻辑从 handler 里长出来。