阶段一 · 微服务架构与通信

异步通信与事件驱动

一句话总结

同步 RPC 像打电话--必须等对方接了才能说事.
异步消息像发快递--发出去就不用等,收件方按自己的节奏处理.
微服务之间大量的交互本质上是"通知"而非"询问",用消息队列解耦是最自然的选择.

为什么需要异步通信

上一篇讲了 gRPC 同步调用--A 发请求,等 B 响应.这个模型有三个隐含假设:

  1. B 必须在线(时间耦合)
  2. B 必须在 deadline 内响应(性能耦合)
  3. A 必须知道 B 的地址(空间耦合)

当这些假设不成立时,同步调用就是错误的工具:

场景 同步 RPC 的问题 异步消息的优势
订单创建后发通知邮件 邮件服务慢/挂了拖垮订单接口 发个事件,通知服务自己消费
支付完成后更新积分 + 物流 + 统计 串行调用 3 个下游,延迟叠加 一个事件广播给所有消费者
流量洪峰(秒杀) 下游直接被打垮 消息队列削峰填谷
跨团队集成 A 团队要知道 B/C/D 的接口细节 A 只发布事件,谁感兴趣谁订阅

同步调用 = 老板站在你工位后面等你处理完才走.
异步消息 = 老板把任务放进你的收件箱,你按优先级处理.
前者让老板被阻塞,后者让双方都以自己的节奏运行.

消息队列基础模型

核心概念

概念 含义
Producer(生产者) 发送消息的服务
Consumer(消费者) 接收并处理消息的服务
Broker(代理) 存储和转发消息的中间件(Kafka / RabbitMQ / NATS)
Topic / Queue 消息的逻辑分类通道
Consumer Group 一组消费者共同消费一个 Topic,实现负载均衡
Offset 消费者在 Topic 中的读取位置(Kafka 特有)

两种基本模式

┌──────────────┐
│ 点对点(Queue) │
├──────────────┤
│ • Producer │
│ • Consumer A │
│ • Consumer B │
└──────────────┘
``` ```text
┌────────────────┐
│ 发布/订阅(Topic) │
├────────────────┤
│ • Producer │
│ • 订阅者 A(通知服务) │
│ • 订阅者 B(统计服务) │
│ • 订阅者 C(积分服务) │
└────────────────┘
点对点(Queue) 发布/订阅(Topic)
消费关系 一条消息只被一个消费者处理 一条消息被所有订阅者各处理一次
典型场景 任务分发,工作队列 事件广播,多下游联动
扩缩容 加消费者 = 提高吞吐 加订阅者 = 新增下游能力

Kafka 深入

Apache Kafka 是微服务领域最主流的消息系统--高吞吐,持久化,支持回溯消费.

架构概览

┌───────────────────────┐
│ Topic: order-events │
├───────────────────────┤
│ • Partition 0 │
│ • Partition 1 │
│ • Partition 2 │
└───────────────────────┘
┌────────────────────────────────┐
│ Consumer Group: notify-group │
├────────────────────────────────┤
│ • Consumer 0 │
│ • Consumer 1 │
│ • Consumer 2 │
└────────────────────────────────┘

核心设计决策:

  • Partition:Topic 被分割成多个分区,每个分区是一个有序的,不可变的消息序列
  • 消息顺序:只保证分区内有序,不保证跨分区有序
  • Consumer Group:同一 Group 内的消费者瓜分 Partition,实现负载均衡;不同 Group 各自独立消费全量数据
  • 持久化:消息写入磁盘,支持按 offset 回溯重新消费

Partition Key 决定消息顺序

相同 key 的消息一定落在同一个 Partition,因此保序.
典型做法:用 user_idorder_id 做 key,保证同一用户/订单的事件有序.没有 key 则 round-robin 分布到各分区.

投递保证

语义 含义 实现方式
At Most Once 最多一次(可能丢) 消费者先 commit offset 再处理
At Least Once 至少一次(可能重复) 消费者处理完再 commit offset(默认)
Exactly Once 恰好一次(不丢不重) Kafka 事务 + 幂等 Producer + 幂等消费者

Exactly Once 是个谎言--在端到端层面

Kafka 自身支持事务性的 Exactly Once(EOS),但这只保证 Kafka 内部不重复.
你的消费者处理完消息后写数据库,调 API--这段逻辑 Kafka 管不了.
实践中的标准答案:At Least Once + 幂等消费.后面会详细讲幂等设计.

Go Producer / Consumer 示例

// Producer(使用 confluent-kafka-go)
producer, _ := kafka.NewProducer(&kafka.ConfigMap{
"bootstrap.servers": "localhost:9092",
"acks": "all", // 等待所有副本确认
})

event := OrderCreatedEvent{OrderID: "ord-123", UserID: "u-1"}
payload, _ := json.Marshal(event)

producer.Produce(&kafka.Message{
TopicPartition: kafka.TopicPartition{Topic: &topic},
Key: []byte(event.OrderID), // 保证同一订单的事件有序
Value: payload,
Headers: []kafka.Header{
{Key: "event-type", Value: []byte("order.created")},
{Key: "trace-id", Value: []byte(traceID)},
},
}, nil)
// Consumer
consumer, _ := kafka.NewConsumer(&kafka.ConfigMap{
"bootstrap.servers": "localhost:9092",
"group.id": "notification-service",
"auto.offset.reset": "earliest",
"enable.auto.commit": false, // 手动 commit,保证 at-least-once
})
consumer.Subscribe("order-events", nil)

for {
msg, err := consumer.ReadMessage(time.Second)
if err != nil { continue }

// 处理消息
if err := handleMessage(msg); err != nil {
// 处理失败:记录日志,后续重试或发送到死信队列
sendToDeadLetterQueue(msg, err)
}

// 处理成功后手动提交 offset
consumer.CommitMessage(msg)
}

RabbitMQ 对比

RabbitMQ 和 Kafka 面向不同的设计目标:

维度 Kafka RabbitMQ
定位 分布式事件流平台 传统消息队列(AMQP)
消息模型 持久化日志,消费者主动拉取 推模型,Broker 主动推给消费者
消息保留 按时间/大小保留,可回溯 消费后即删(可配置持久化)
顺序保证 分区内有序 单队列内有序
吞吐量 百万级/秒 万级/秒
路由能力 简单(Topic + Partition Key) 强大(Exchange + Binding + Routing Key)
适用场景 事件流,日志收集,大数据管道 任务队列,RPC,复杂路由,延时队列

选择建议

Kafka:需要高吞吐,事件回溯,多消费者独立消费同一份数据时选择.
RabbitMQ:需要灵活路由,延时队列,轻量级任务分发时选择.
大多数微服务场景中 Kafka 是默认选择--它的 Consumer Group + 持久化 + 回溯能力天然适合事件驱动架构.

事件驱动架构(EDA)

事件的三种类型

类型 内容 示例 消费者行为
Event Notification 只通知"发生了什么",不携带详情 {type: "order.created", order_id: "123"} 消费者需要回调查询详情
Event-Carried State Transfer 携带变更后的完整/增量数据 {type: "order.created", order: {id, items, total, ...}} 消费者本地保存副本,不需要回调
Domain Event 描述领域中发生的业务事实 {type: "payment.completed", amount: 9900, method: "wechat"} 触发后续业务流程

Event Notification vs State Transfer 的取舍

Notification 轻量但引入了回调耦合(消费者还是要知道生产者的 API).
State Transfer 消息大但彻底解耦(消费者完全不依赖生产者在线).
实践中常见折中:关键字段内嵌 + 非关键字段按需回查.

事件 Schema 设计

一个生产级事件的标准结构:

{
"event_id": "evt-a1b2c3",
"event_type": "order.created",
"source": "order-service",
"timestamp": "2024-03-15T10:30:00Z",
"trace_id": "trace-xyz",
"version": "1.0",
"data": {
"order_id": "ord-123",
"user_id": "u-456",
"items": [...],
"total_cents": 29900
}
}

设计原则:

  • event_id:全局唯一,用于幂等去重
  • event_type:使用 领域.动作 命名(过去时),如 order.created,payment.completed
  • source:产生事件的服务名
  • version:schema 版本,支持消费者区分处理
  • trace_id:串联分布式追踪链路

幂等消费

At Least Once 意味着同一条消息可能被投递多次(网络抖动,消费者重启,rebalance).消费者必须做到"处理 N 次和处理 1 次效果相同"--这就是幂等性.

幂等实现策略

策略 原理 适用场景
唯一约束 数据库 UNIQUE KEY 防止重复插入 创建类操作(创建订单,创建用户)
幂等表 / 消息去重表 记录已处理的 event_id,重复则跳过 通用方案,适合所有场景
条件更新 UPDATE ... WHERE version = N 乐观锁 更新类操作
天然幂等 操作本身重复执行无副作用 设置状态(SET status='paid')

幂等表实现示例

func handleMessage(ctx context.Context, msg *Event) error {
// 1. 检查是否已处理(幂等表)
exists, err := db.ExecContext(ctx,
"INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW()) ON CONFLICT DO NOTHING",
msg.EventID,
)
if err != nil { return err }
if exists.RowsAffected() == 0 {
// 已处理过,跳过
log.Info("duplicate event, skipping", "event_id", msg.EventID)
return nil
}

// 2. 执行业务逻辑(在同一个事务中!)
return processOrderCreated(ctx, msg.Data)
}

幂等表和业务操作必须在同一个事务中

如果先插幂等表,再执行业务逻辑--业务失败了怎么办?幂等表已经标记"已处理",重试时就跳过了,消息永远丢失.
正确做法:把幂等表插入和业务操作放在同一个数据库事务中,要么一起成功,要么一起回滚.

发件箱模式(Transactional Outbox)

一个经典问题:服务需要"写数据库 + 发消息"两个操作.它们不在同一个事务里--如何保证一致性?

双写问题

场景:订单服务创建订单后发送事件

方案 A:先写 DB 再发消息
1. INSERT INTO orders ... ✓
2. kafka.Produce("order.created") ✗ ← 发送失败!
结果:订单已创建,但事件丢失,下游永远不知道

方案 B:先发消息再写 DB
1. kafka.Produce("order.created") ✓
2. INSERT INTO orders ... ✗ ← 写入失败!
结果:事件已发出,但订单不存在,下游处理一个幽灵事件

无论先做哪个,都可能出现不一致.这就是双写问题(Dual Write Problem).

发件箱方案

核心思想:不直接发消息,而是把待发送的消息写入本地数据库的 outbox 表(与业务操作在同一个事务中).一个独立的进程轮询 outbox 表,将消息投递到消息队列.

sequenceDiagram
participant APP as 订单服务
participant DB as 数据库
participant RELAY as Outbox Relay
participant MQ as Kafka

APP->>DB: BEGIN TRANSACTION
APP->>DB: INSERT INTO orders (...)
APP->>DB: INSERT INTO outbox (event_type, payload)
APP->>DB: COMMIT
Note over DB: 数据库事务保证原子性

loop 轮询 / CDC
RELAY->>DB: SELECT * FROM outbox WHERE status='pending'
RELAY->>MQ: Produce message
RELAY->>DB: UPDATE outbox SET status='sent'
end
-- outbox 表结构
CREATE TABLE outbox (
id BIGSERIAL PRIMARY KEY,
event_id UUID NOT NULL UNIQUE,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
status VARCHAR(20) DEFAULT 'pending',
created_at TIMESTAMP DEFAULT NOW(),
sent_at TIMESTAMP
);

-- 业务代码:在同一个事务中写业务表 + outbox
BEGIN;
INSERT INTO orders (id, user_id, total) VALUES ('ord-123', 'u-1', 299);
INSERT INTO outbox (event_id, event_type, payload)
VALUES ('evt-abc', 'order.created', '{"order_id":"ord-123","user_id":"u-1","total":299}');
COMMIT;

发件箱模式 = 本地消息表

国内常说的"本地消息表"和这里的发件箱模式(Transactional Outbox)是同一套机制--把"发消息"降级为"写一行本地记录",塞进业务的同一个本地事务,从而规避双写问题.
两个名字可视为同义词,前者源自国内(阿里系)技术圈,后者是微服务社区的英文术语.

发件箱模式 = 先放自家发件箱再由邮差取走:
直接投递到邮局(MQ)可能因为路上出岔子而和你家的账本(DB)对不上.
发件箱模式是先把信放进自家门口的发件箱(outbox 表,和记账同一个动作完成),再由邮差(Relay)按部就班取走投递--账本和发件箱是同时落定的,绝不会一个有一个没有.

投递语义:at-least-once

发件箱模式保证消息不丢,但代价是可能重:Relay "已经投递到 MQ,但还没来得及把 outbox 标记为 sent" 时崩溃,重启后会重发同一条消息.因此它的投递语义是 at-least-once(至少一次),而非 exactly-once.

完整方案 = 发送方消息表 + 接收方去重表

这正好呼应前面的幂等消费:发件箱保证"不丢",去重表(processed_events)消化"可能重".
狭义的 Outbox 只管发送方可靠发出消息;而国内讲"本地消息表"做最终一致性事务时,往往指端到端闭环--发送方 outbox 表 + 接收方去重表配对出现,两者合起来才得到可靠的最终一致性.

CDC 替代方案

轮询 outbox 表有延迟和性能问题.更优方案是使用 CDC(Change Data Capture)--监听数据库的 binlog / WAL,实时将变更发送到 Kafka:

  • Debezium:最流行的 CDC 工具,支持 MySQL / PostgreSQL / MongoDB
  • 直接监听 outbox 表的 INSERT 事件,零轮询,低延迟
  • 也可以监听业务表本身的变更--连 outbox 表都不用了

生产推荐

小规模:outbox + 定时轮询(简单可靠,延迟秒级).
中大规模:outbox + Debezium CDC(实时,不增加数据库查询压力).
超大规模:直接 CDC 监听业务表变更(省掉 outbox 表,但 schema 变更管理更复杂).

事件溯源(Event Sourcing)

传统方式存储实体的"当前状态"(一行记录).事件溯源存储实体经历的"所有事件序列"--当前状态通过回放事件推导得出.

核心概念

传统 CRUD:
orders 表:{id: "ord-1", status: "shipped", total: 299, ...}
→ 只看到最终状态,中间发生了什么?不知道

事件溯源:
Event 1: OrderCreated {order_id: "ord-1", items: [...], total: 299}
Event 2: PaymentReceived {order_id: "ord-1", amount: 299, method: "wechat"}
Event 3: OrderShipped {order_id: "ord-1", tracking: "SF123456"}
→ 当前状态 = 回放所有事件得出
→ 完整审计轨迹,可以回到任意历史时间点

适用场景与代价

优势 代价
完整审计轨迹(金融,合规) 查询复杂(需要物化视图/CQRS)
可回溯到任意时间点 事件 schema 演进困难
天然支持事件驱动 开发心智模型转变大
易于调试(回放重现问题) 存储量大(事件只增不删)

不要默认使用事件溯源

事件溯源是一个高成本,高收益的架构选择.只在以下场景考虑:

  1. 强合规要求(金融交易必须可审计)
  2. 业务本身就是事件序列(订单状态机,工作流引擎)
  3. 需要时间旅行能力(回到过去某一刻的状态)

大多数 CRUD 业务用传统方式 + 事件通知即可.

死信队列(Dead Letter Queue)

消费者处理消息失败后怎么办?无限重试会卡住整个队列(Head-of-Line Blocking).标准做法:重试 N 次后转入死信队列(DLQ).

[RETRY] → [主队列] → [Consumer]
[DLQ] → [告警 + 人工处理]

DLQ 消息的处理:

  • 人工排查原因后修复 bug,再从 DLQ 重新投递
  • 配置自动重试策略(如 1 小时后再试一次)
  • 持续失败的消息标记为永久失败,记录到监控报表

消息顺序性

并非所有场景都需要严格有序.分清楚:

场景 是否需要有序 实现方式
同一订单的状态变更 需要(创建→支付→发货) 用 order_id 做 Partition Key
不同用户的订单 不需要 不同分区并行消费
日志收集 尽量有序但允许乱序 按时间戳在消费端排序

顺序性 vs 吞吐量是对立的

要严格有序就只能单分区单消费者--吞吐量极低.
实践中的策略:按业务实体(order_id / user_id)分区,保证同一实体内有序,不同实体之间并行处理.这在绝大多数场景下是足够的.

背压与流控

当消费者处理速度跟不上生产者时:

  • Kafka 天然支持背压:消费者拉模型(pull),处理完一批再拉下一批.消息堆积在 Broker 的磁盘上(Kafka 的持久化日志设计让堆积是安全的)
  • RabbitMQ:推模型需要手动设置 prefetch count 限制同时处理的消息数

生产环境的应对:

  1. 监控消费延迟(Consumer Lag):当前 offset 与最新 offset 的差值,是最重要的 Kafka 监控指标
  2. 水平扩容消费者:增加 Consumer Group 中的实例数(不超过 Partition 数)
  3. 限流上游:如果堆积持续增长,考虑对 Producer 限流

事件驱动与 Saga 模式(预告)

异步消息天然适合编排跨服务的业务流程.例如"创建订单"涉及库存扣减,支付,物流--每一步都是一个事件触发下一步:

订单服务: OrderCreated →
库存服务: InventoryReserved →
支付服务: PaymentCompleted →
物流服务: ShipmentCreated

任何一步失败 → 发送补偿事件 → 回滚前序操作

这就是 Saga 模式的核心思想--第 07 篇会深入讲解.

生产最佳实践

  • 消息不要太大:单条消息建议 < 1MB.大文件放对象存储,消息里只带 URL
  • Schema 注册表:用 Confluent Schema Registry 管理 Avro / Protobuf schema,防止生产者随意改格式
  • 消费者要快:消息处理逻辑应尽量简单,耗时操作异步化(消费者内部再用 goroutine pool)
  • 监控三大指标:Consumer Lag(消费延迟),消息处理延时 P99,死信队列积压量
  • 消息要带 trace_id:否则跨服务的链路追踪在异步环节断裂
  • 版本兼容:事件 schema 只做加法,新增字段不影响旧消费者(和 Protobuf 同理)

快速回顾

  • 异步 vs 同步:通知类交互(不需要即时响应)用消息队列解耦
  • Kafka:分区有序,持久化日志,Consumer Group 负载均衡,可回溯
  • RabbitMQ:灵活路由,延时队列,轻量任务分发
  • At Least Once + 幂等消费是生产标准,不要追求端到端 Exactly Once
  • 发件箱模式:解决"写 DB + 发消息"的双写一致性问题
  • 事件溯源:高价值高成本,仅在合规审计,状态机场景使用
  • 死信队列:兜底方案,防止毒消息卡住正常消费
  • 顺序性:按业务实体做 Partition Key,同实体内有序即可

动手练习

  1. 启动 Kafka 集群:用 Docker Compose 启动一个 Kafka 集群(单 Broker + ZooKeeper),创建一个 3 分区的 Topic,用 kafka-console-producer/consumer 验证消息收发
  2. Consumer Group 分配:写一个 Go Producer 发送 100 条带 key 的消息,启动 3 个 Consumer(同一 Group),观察消息如何在消费者间分配
  3. 实现 DLQ 重试:模拟消费失败--让消费者在处理第 5 条消息时返回错误,实现"重试 3 次后发送到 DLQ"的逻辑
  4. 实现发件箱模式:在创建订单的同一事务中写入 outbox 表,写一个 goroutine 轮询 outbox 并发送到 Kafka
  5. 实现幂等消费:用 event_id 去重表确保同一事件不会被重复处理