消息驱动架构
#消息队列 · #Kafka · #RocketMQ · #事务消息 · #最终一致性 · #幂等
消息驱动架构通过异步解耦实现系统间松耦合通信,是构建高可用、高吞吐分布式系统的核心模式。本专题聚焦消息系统的设计模式、可靠性保障与工程实践。
为什么需要消息驱动
- 解耦:上下游系统通过消息通信,互不感知对方实现
- 削峰填谷:突发流量写入消息队列,消费端按能力消费
- 异步处理:非核心链路异步化,降低主链路延迟
- 事件溯源:消息作为事件日志,支持回放和审计
消息投递语义
| 语义 | 含义 | 实现难度 |
|---|---|---|
| At Most Once | 最多投递一次,可能丢失 | 低 |
| At Least Once | 至少投递一次,可能重复 | 中 |
| Exactly Once | 精确一次,不丢不重 | 高 |
工程实践:大多数系统选择 At Least Once + 幂等消费 来近似实现 Exactly Once。
消息可靠性保障
生产端可靠发送
java
// Kafka 可靠发送配置
Properties props = new Properties();
props.put("acks", "all"); // 所有 ISR 副本确认
props.put("retries", 3); // 失败重试
props.put("enable.idempotence", true); // 幂等生产者
props.put("max.in.flight.requests.per.connection", 5);Broker 端可靠存储
- 多副本:Kafka ISR 机制、RocketMQ DLedger
- 刷盘策略:同步刷盘(可靠)vs 异步刷盘(高性能)
- 持久化:消息写入磁盘后才确认
消费端可靠消费
python
# 手动提交 offset,确保消费成功后才确认
def consume_message(msg):
try:
process(msg) # 业务处理
consumer.commit() # 处理成功后提交
except Exception as e:
log.error(f"消费失败: {e}")
# 不提交 offset,消息会被重新投递
raise幂等消费
为什么需要幂等
消息重复投递不可避免(网络抖动、消费者重启等),消费端必须保证:同一消息处理多次,结果与处理一次相同。
幂等方案
| 方案 | 原理 | 适用场景 |
|---|---|---|
| 唯一 ID + 去重表 | 消费前查询是否已处理 | 通用 |
| 数据库唯一约束 | 利用 UNIQUE KEY 天然去重 | 插入操作 |
| 乐观锁/版本号 | 带版本号更新,重复更新不生效 | 更新操作 |
| Redis SETNX | 利用 Redis 原子性判断是否已消费 | 高并发场景 |
| 状态机 | 业务状态只能单向流转 | 订单等有状态业务 |
python
# 唯一 ID + Redis 去重
def idempotent_consume(msg):
msg_id = msg.get("msg_id")
# 利用 Redis SETNX 判断是否已消费
if not redis.set(f"consumed:{msg_id}", "1", nx=True, ex=86400):
log.info(f"消息已消费,跳过: {msg_id}")
return
try:
process(msg)
except Exception:
redis.delete(f"consumed:{msg_id}") # 处理失败,删除标记允许重试
raise事务消息
问题场景
下单流程:
1. 创建订单(写 DB)
2. 发送消息通知库存扣减
如何保证 "写 DB" 和 "发消息" 的原子性?RocketMQ 事务消息
Producer Broker Consumer
|--- Half Message -------->| |
|<-- Half OK --------------| |
| | |
| [执行本地事务] | |
| | |
|--- Commit/Rollback ----->| |
| |--- Deliver Message ----->|java
// RocketMQ 事务消息示例
TransactionMQProducer producer = new TransactionMQProducer("group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
orderService.createOrder(order); // 执行本地事务
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 事务回查:检查本地事务是否成功
if (orderService.exists(msg.getTransactionId())) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
});本地消息表
sql
-- 业务操作和消息写入在同一个事务中
BEGIN;
INSERT INTO orders (id, ...) VALUES (...);
INSERT INTO outbox_messages (id, topic, payload, status)
VALUES (uuid, 'order_created', '{}', 'PENDING');
COMMIT;
-- 定时任务扫描 outbox 表,发送消息
-- 发送成功后更新状态为 SENT死信队列(DLQ)
什么是死信
消费失败达到最大重试次数后,消息进入死信队列,等待人工处理。
死信处理流程
正常队列 → 消费失败 → 重试队列(延迟重试) → 再次失败 → 死信队列 → 告警 → 人工处理python
# 带重试次数的消费逻辑
MAX_RETRY = 3
def consume_with_retry(msg):
retry_count = msg.get_property("retry_count", 0)
try:
process(msg)
except Exception as e:
if retry_count < MAX_RETRY:
# 延迟重试,指数退避
delay = 2 ** retry_count # 1s, 2s, 4s
msg.set_property("retry_count", retry_count + 1)
producer.send_delay(msg, delay_seconds=delay)
else:
# 进入死信队列
producer.send("DLQ_TOPIC", msg)
alert(f"消息进入死信队列: {msg.msg_id}")消息顺序性
全局有序
- 单分区/单队列保证
- 吞吐量受限,适用于要求严格有序的场景
分区有序
- 相同业务 Key 的消息路由到同一分区
- 兼顾有序性和吞吐量
python
# Kafka 按订单 ID 分区,保证同一订单的消息有序
producer.send(
topic="order_events",
key=order_id.encode(), # 相同 order_id 路由到同一分区
value=event_payload
)延迟消息
应用场景
- 订单超时未支付自动取消(30 分钟延迟)
- 定时任务触发
- 重试退避
实现方案
| 方案 | 原理 | 精度 |
|---|---|---|
| RocketMQ 延迟级别 | 预定义延迟级别(1s/5s/10s/...) | 固定级别 |
| Kafka + 时间轮 | 自定义延迟队列 | 秒级 |
| Redis ZSET | Score 为到期时间,定时扫描 | 秒级 |
| 数据库轮询 | 定时扫描待执行记录 | 分钟级 |
事件溯源(Event Sourcing)
核心思想
不存储当前状态,而是存储所有状态变更事件,当前状态通过回放事件得到。
Event Store:
OrderCreated { orderId: 1, items: [...], time: T1 }
OrderPaid { orderId: 1, amount: 100, time: T2 }
OrderShipped { orderId: 1, trackingNo: "xxx", time: T3 }
当前状态 = replay(events)优势
- 完整审计日志
- 支持时间旅行(回放到任意时刻)
- 天然适合 CQRS(命令查询职责分离)
相关专题
消息系统选型决策
选择消息系统不是"哪个最好",而是"哪个最适合你的场景"。以下决策树帮助快速定位:
mermaid
flowchart TB
Start["需要消息系统"] --> Q1{"核心需求?"}
Q1 -->|"高吞吐 + 大数据/流处理"| Kafka["✅ Kafka<br/>百万级 TPS<br/>日志/埋点/流计算"]
Q1 -->|"事务消息 + 延迟消息"| RocketMQ["✅ RocketMQ<br/>金融/电商<br/>分布式事务"]
Q1 -->|"灵活路由 + 优先级队列"| RabbitMQ["✅ RabbitMQ<br/>复杂路由<br/>中小规模"]
Q1 -->|"多租户 + 存算分离"| Pulsar["✅ Pulsar<br/>云原生<br/>超大规模"]
Kafka --> K1{"需要 Exactly Once?"}
K1 -->|"是"| K2["Kafka Streams/Flink<br/>+ 幂等 Producer"]
K1 -->|"否"| K3["标准 Consumer Group"]
RocketMQ --> R1{"需要强一致?"}
R1 -->|"是"| R2["DLedger 模式<br/>同步双写"]
R1 -->|"否"| R3["异步复制<br/>(性能优先)"]四大消息系统核心对比
| 维度 | Kafka | RocketMQ | RabbitMQ | Pulsar |
|---|---|---|---|---|
| 吞吐量 | 百万/s | 十万/s | 万/s | 百万/s |
| 延迟 | ms 级 | ms 级 | μs 级 | ms 级 |
| 事务消息 | ✅ (Streams) | ✅ (原生半消息) | ✅ (AMQP 事务) | ✅ |
| 延迟消息 | ❌ (需自实现) | ✅ (18 级) | ✅ (TTL+DLX) | ✅ (任意延迟) |
| 消息回溯 | ✅ (按 offset/时间) | ✅ (按时间) | ❌ | ✅ |
| 优先级队列 | ❌ | ❌ | ✅ | ❌ |
| 协议 | 自定义 | 自定义 | AMQP 0-9-1 | 自定义 |
| 存储 | 分区文件 | CommitLog | 内存+磁盘 | BookKeeper |
| 运维复杂度 | 中 | 中高 | 低 | 高 |
| 社区生态 | 极强 (大数据) | 阿里系 | 广泛 | 新兴 |
背压(Backpressure)机制
当消费者处理速度跟不上生产者发送速度时,如果没有背压机制,消息会无限堆积直到系统崩溃。不同消息系统的背压策略差异很大:
背压机制对比
mermaid
flowchart TB
subgraph Kafka["Kafka 背压"]
KP["Producer"] -->|"发送"| KB["Broker<br/>(磁盘存储)"]
KB -->|"Consumer 按自身速度 Pull"| KC["Consumer"]
KNote["Pull 模型天然背压:<br/>Consumer 不拉就不消费<br/>消息在 Broker 磁盘堆积"]
end
subgraph RocketMQ["RocketMQ 背压"]
RP["Producer"] -->|"发送"| RB["Broker"]
RB -->|"Pull (长轮询)"| RC["Consumer"]
RNote["流控机制:<br/>1. Broker 内存>水位 → 拒绝写入<br/>2. Consumer 本地缓冲>阈值 → 暂停拉取"]
end
subgraph RabbitMQ["RabbitMQ 背压"]
QP["Producer"] -->|"Push"| QB["Broker"]
QB -->|"Push"| QC["Consumer"]
QNote["QoS prefetch:<br/>Consumer 设置 prefetch_count<br/>未 ACK 数达上限 → 停止推送<br/>队列满 → 阻塞 Producer Connection"]
end各系统背压配置
python
# Kafka Consumer: 通过 max.poll.records 控制每次拉取量
# 如果处理太慢,自动减少拉取频率
consumer_config = {
"max.poll.records": 500, # 每次 poll 最多 500 条
"max.poll.interval.ms": 300000, # 5 分钟内必须 poll,否则被踢出组
"fetch.max.bytes": 52428800, # 单次 fetch 最大 50MB
}
# 手动暂停/恢复(精细控制)
consumer.pause(partitions) # 暂停消费某些分区
# ... 处理积压 ...
consumer.resume(partitions) # 恢复消费java
// RabbitMQ: prefetch 控制
channel.basicQos(100); // 最多 100 条未确认消息
// 超过 100 条未 ACK → RabbitMQ 停止向该 Consumer 推送
// → 消息在队列中堆积
// → 队列达到 x-max-length → 新消息被拒绝或进入死信java
// RocketMQ: Consumer 端流控
consumer.setPullThresholdForQueue(1000); // 本地缓冲队列最大 1000 条
consumer.setPullThresholdSizeForQueue(100); // 本地缓冲最大 100MB
// 超过阈值 → 暂停拉取 → 等待消费完再继续消息积压排查与处理
消息积压是生产环境最常见的消息系统问题。排查流程:
mermaid
flowchart TB
Alert["🔔 消息积压告警<br/>lag > 阈值"] --> Check{"确认积压位置"}
Check -->|"生产端积压"| ProducerIssue["Producer 发送失败<br/>→ 检查 Broker 是否可用<br/>→ 检查网络/磁盘"]
Check -->|"消费端积压"| ConsumerCheck{"消费者状态?"}
ConsumerCheck -->|"消费者正常但慢"| Slow["消费逻辑慢"]
ConsumerCheck -->|"消费者挂了"| Down["消费者宕机"]
ConsumerCheck -->|"Rebalance 频繁"| Rebalance["频繁重平衡"]
Slow --> SlowFix["1. 增加消费者实例<br/>2. 增加消费线程<br/>3. 优化消费逻辑<br/>4. 批量消费"]
Down --> DownFix["1. 重启消费者<br/>2. 检查 OOM/异常<br/>3. 检查依赖服务"]
Rebalance --> RbFix["1. 增大 session.timeout<br/>2. 增大 max.poll.interval<br/>3. 减少单次处理时间"]
SlowFix --> Emergency{"积压是否紧急?"}
Emergency -->|"是"| EmergencyFix["紧急扩容方案"]
Emergency -->|"否"| Normal["等待自然消化"]
EmergencyFix --> E1["方案 1: 临时扩容消费者<br/>(需要分区数 >= 消费者数)"]
EmergencyFix --> E2["方案 2: 转发到新 Topic<br/>(更多分区 + 更多消费者)"]
EmergencyFix --> E3["方案 3: 跳过非关键消息<br/>(记录后丢弃)"]积压监控 PromQL
promql
# Kafka Consumer Lag(通过 kafka_exporter)
sum by (consumergroup, topic) (kafka_consumergroup_lag)
# 告警:lag 持续增长
kafka_consumergroup_lag > 10000
and
increase(kafka_consumergroup_lag[5m]) > 0
# RocketMQ 积压(通过 rocketmq_exporter)
rocketmq_consumer_diff > 5000紧急扩容方案
python
# 方案 2: 转发到临时 Topic(快速消化积压)
# 原 Topic: order_events (4 分区, 4 消费者 → 处理不过来)
# 临时 Topic: order_events_tmp (32 分区, 32 消费者)
# 步骤 1: 启动转发程序(只转发,不处理业务)
def forward_messages():
for msg in consumer.poll("order_events"):
# 按原 key 发到临时 Topic(保持顺序)
producer.send("order_events_tmp", key=msg.key, value=msg.value)
consumer.commit()
# 步骤 2: 启动 32 个消费者消费临时 Topic
# 步骤 3: 积压消化完后,恢复原消费者,停止转发Go 版本消息消费幂等实现
go
package consumer
import (
"context"
"fmt"
"time"
"github.com/IBM/sarama"
"github.com/redis/go-redis/v9"
)
type IdempotentConsumer struct {
rdb *redis.Client
}
// 幂等检查: 用 msg key + partition + offset 判定是否已处理过
func (c *IdempotentConsumer) IsDuplicate(msg *sarama.ConsumerMessage) (bool, error) {
key := fmt.Sprintf("msg:dedup:%s:%d:%d",
msg.Topic, msg.Partition, msg.Offset)
ok, err := c.rdb.SetNX(context.Background(), key, "1", 24*time.Hour).Result()
if err != nil {
return false, err
}
return !ok, nil
}
func (c *IdempotentConsumer) ConsumeWithIdempotency(
msg *sarama.ConsumerMessage,
handler func([]byte) error,
) error {
dup, err := c.IsDuplicate(msg)
if err != nil {
log.Warn("redis dedup failed, falling through", "err", err)
}
if dup {
return nil
}
if err := handler(msg.Value); err != nil {
c.rdb.Del(context.Background(),
fmt.Sprintf("msg:dedup:%s:%d:%d", msg.Topic, msg.Partition, msg.Offset))
return err
}
return nil
}go
// 死信队列 (DLQ) 处理
type DLQHandler struct {
maxRetries int
dlqTopic string
}
func (h *DLQHandler) handle(msg *sarama.ConsumerMessage, retryCount int) error {
err := processMessage(msg.Value)
if err == nil {
return nil
}
if retryCount < h.maxRetries {
return err // 不 commit → 重新投递
}
// 超限 → 死信队列
dlqMsg := &sarama.ProducerMessage{
Topic: h.dlqTopic,
Value: sarama.ByteEncoder(msg.Value),
Headers: []sarama.RecordHeader{
{Key: []byte("original-topic"), Value: []byte(msg.Topic)},
{Key: []byte("error-at"), Value: []byte(time.Now().Format(time.RFC3339))},
},
}
producer.SendMessage(dlqMsg)
return nil // commit 原 offset
}mermaid
flowchart TB
A["Kafka 消息到达"] --> B{"幂等检查<br/>Redis SETNX"}
B -->|"key 不存在"| C["执行业务逻辑"]
B -->|"key 已存在"| SKIP["跳过"]
C --> D{"成功?"}
D -->|"是"| E["Commit ✅"]
D -->|"否"| F{"重试 < 3?"}
F -->|"是"| G["不 Commit → 重投"]
F -->|"否"| H["DLQ Topic<br/>Commit Offset"]
G --> C
SKIP --> E| 组件 | 作用 | 推荐配置 |
|---|---|---|
| Redis 幂等 Key | 防重复消费 | TTL 24h |
| Kafka Retry | 临时失败重试 | max 3 次 |
| DLQ | 永久失败归集 | 保留 7 天 |
登录后即可发表评论 👇