一句话总结
同步 RPC 像打电话--必须等对方接了才能说事.
异步消息像发快递--发出去就不用等,收件方按自己的节奏处理.
微服务之间大量的交互本质上是"通知"而非"询问",用消息队列解耦是最自然的选择.
为什么需要异步通信
上一篇讲了 gRPC 同步调用--A 发请求,等 B 响应.这个模型有三个隐含假设:
- B 必须在线(时间耦合)
- B 必须在 deadline 内响应(性能耦合)
- A 必须知道 B 的地址(空间耦合)
当这些假设不成立时,同步调用就是错误的工具:
| 场景 | 同步 RPC 的问题 | 异步消息的优势 |
|---|---|---|
| 订单创建后发通知邮件 | 邮件服务慢/挂了拖垮订单接口 | 发个事件,通知服务自己消费 |
| 支付完成后更新积分 + 物流 + 统计 | 串行调用 3 个下游,延迟叠加 | 一个事件广播给所有消费者 |
| 流量洪峰(秒杀) | 下游直接被打垮 | 消息队列削峰填谷 |
| 跨团队集成 | A 团队要知道 B/C/D 的接口细节 | A 只发布事件,谁感兴趣谁订阅 |
同步调用 = 老板站在你工位后面等你处理完才走.
异步消息 = 老板把任务放进你的收件箱,你按优先级处理.
前者让老板被阻塞,后者让双方都以自己的节奏运行.
消息队列基础模型
核心概念
| 概念 | 含义 |
|---|---|
| Producer(生产者) | 发送消息的服务 |
| Consumer(消费者) | 接收并处理消息的服务 |
| Broker(代理) | 存储和转发消息的中间件(Kafka / RabbitMQ / NATS) |
| Topic / Queue | 消息的逻辑分类通道 |
| Consumer Group | 一组消费者共同消费一个 Topic,实现负载均衡 |
| Offset | 消费者在 Topic 中的读取位置(Kafka 特有) |
两种基本模式
|
| 点对点(Queue) | 发布/订阅(Topic) | |
|---|---|---|
| 消费关系 | 一条消息只被一个消费者处理 | 一条消息被所有订阅者各处理一次 |
| 典型场景 | 任务分发,工作队列 | 事件广播,多下游联动 |
| 扩缩容 | 加消费者 = 提高吞吐 | 加订阅者 = 新增下游能力 |
Kafka 深入
Apache Kafka 是微服务领域最主流的消息系统--高吞吐,持久化,支持回溯消费.
架构概览
|
核心设计决策:
- Partition:Topic 被分割成多个分区,每个分区是一个有序的,不可变的消息序列
- 消息顺序:只保证分区内有序,不保证跨分区有序
- Consumer Group:同一 Group 内的消费者瓜分 Partition,实现负载均衡;不同 Group 各自独立消费全量数据
- 持久化:消息写入磁盘,支持按 offset 回溯重新消费
Partition Key 决定消息顺序
相同 key 的消息一定落在同一个 Partition,因此保序.
典型做法:用 user_id 或 order_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 示例
|
|
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:全局唯一,用于幂等去重
- 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') |
幂等表实现示例
|
幂等表和业务操作必须在同一个事务中
如果先插幂等表,再执行业务逻辑--业务失败了怎么办?幂等表已经标记"已处理",重试时就跳过了,消息永远丢失.
正确做法:把幂等表插入和业务操作放在同一个数据库事务中,要么一起成功,要么一起回滚.
发件箱模式(Transactional Outbox)
一个经典问题:服务需要"写数据库 + 发消息"两个操作.它们不在同一个事务里--如何保证一致性?
双写问题
|
无论先做哪个,都可能出现不一致.这就是双写问题(Dual Write Problem).
发件箱方案
核心思想:不直接发消息,而是把待发送的消息写入本地数据库的 outbox 表(与业务操作在同一个事务中).一个独立的进程轮询 outbox 表,将消息投递到消息队列.
|
|
发件箱模式 = 本地消息表
国内常说的"本地消息表"和这里的发件箱模式(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)
传统方式存储实体的"当前状态"(一行记录).事件溯源存储实体经历的"所有事件序列"--当前状态通过回放事件推导得出.
核心概念
|
适用场景与代价
| 优势 | 代价 |
|---|---|
| 完整审计轨迹(金融,合规) | 查询复杂(需要物化视图/CQRS) |
| 可回溯到任意时间点 | 事件 schema 演进困难 |
| 天然支持事件驱动 | 开发心智模型转变大 |
| 易于调试(回放重现问题) | 存储量大(事件只增不删) |
不要默认使用事件溯源
事件溯源是一个高成本,高收益的架构选择.只在以下场景考虑:
- 强合规要求(金融交易必须可审计)
- 业务本身就是事件序列(订单状态机,工作流引擎)
- 需要时间旅行能力(回到过去某一刻的状态)
大多数 CRUD 业务用传统方式 + 事件通知即可.
死信队列(Dead Letter Queue)
消费者处理消息失败后怎么办?无限重试会卡住整个队列(Head-of-Line Blocking).标准做法:重试 N 次后转入死信队列(DLQ).
|
DLQ 消息的处理:
- 人工排查原因后修复 bug,再从 DLQ 重新投递
- 配置自动重试策略(如 1 小时后再试一次)
- 持续失败的消息标记为永久失败,记录到监控报表
消息顺序性
并非所有场景都需要严格有序.分清楚:
| 场景 | 是否需要有序 | 实现方式 |
|---|---|---|
| 同一订单的状态变更 | 需要(创建→支付→发货) | 用 order_id 做 Partition Key |
| 不同用户的订单 | 不需要 | 不同分区并行消费 |
| 日志收集 | 尽量有序但允许乱序 | 按时间戳在消费端排序 |
顺序性 vs 吞吐量是对立的
要严格有序就只能单分区单消费者--吞吐量极低.
实践中的策略:按业务实体(order_id / user_id)分区,保证同一实体内有序,不同实体之间并行处理.这在绝大多数场景下是足够的.
背压与流控
当消费者处理速度跟不上生产者时:
- Kafka 天然支持背压:消费者拉模型(pull),处理完一批再拉下一批.消息堆积在 Broker 的磁盘上(Kafka 的持久化日志设计让堆积是安全的)
- RabbitMQ:推模型需要手动设置 prefetch count 限制同时处理的消息数
生产环境的应对:
- 监控消费延迟(Consumer Lag):当前 offset 与最新 offset 的差值,是最重要的 Kafka 监控指标
- 水平扩容消费者:增加 Consumer Group 中的实例数(不超过 Partition 数)
- 限流上游:如果堆积持续增长,考虑对 Producer 限流
事件驱动与 Saga 模式(预告)
异步消息天然适合编排跨服务的业务流程.例如"创建订单"涉及库存扣减,支付,物流--每一步都是一个事件触发下一步:
|
这就是 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,同实体内有序即可
动手练习
- 启动 Kafka 集群:用 Docker Compose 启动一个 Kafka 集群(单 Broker + ZooKeeper),创建一个 3 分区的 Topic,用 kafka-console-producer/consumer 验证消息收发
- Consumer Group 分配:写一个 Go Producer 发送 100 条带 key 的消息,启动 3 个 Consumer(同一 Group),观察消息如何在消费者间分配
- 实现 DLQ 重试:模拟消费失败--让消费者在处理第 5 条消息时返回错误,实现"重试 3 次后发送到 DLQ"的逻辑
- 实现发件箱模式:在创建订单的同一事务中写入 outbox 表,写一个 goroutine 轮询 outbox 并发送到 Kafka
- 实现幂等消费:用 event_id 去重表确保同一事件不会被重复处理