跳转至

消息队列与分布式中间件

image.png

消息顺序性与分区策略

📬 1. 如何保证消息的顺序性?

顺序性的根源:分区与并行

首先要明白,消息队列的“顺序”是有作用域的。Kafka 中,单个分区内消息天然有序,因为消息是顺序追加到日志文件中的。但在多个分区之间,全局没有顺序保证。RocketMQ 也是类似的,一个队列(相当于分区)内有序。

Topic: "order-events"
├── Partition 0: [msg1] → [msg2] → [msg3]   ← 这里严格有序
├── Partition 1: [msg4] → [msg5] → [msg6]
└── Partition 2: [msg7] → [msg8] → [msg9]

如果顺序处理 msg1 和 msg5,全局无序。

保证顺序的策略:

① 局部有序(推荐)

将需要顺序的消息发送到同一个分区(或队列)。在 Kafka 中,通过指定相同的 Key 来实现,具有相同 Key 的消息会路由到同一分区。在 RocketMQ 中,可以发送到指定的 MessageQueue。

Kafka 示例:使用订单 ID 作为 Key

// 生产者:相同 orderId 的消息进入同一分区
ProducerRecord<String, String> record = new ProducerRecord<>("order-events", orderId, message);
kafkaProducer.send(record);
// 消费者:正常消费即可,分区内有序
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
    processInOrder(record);
}

RocketMQ 示例:使用 MessageQueueSelector

// 生产者:选择同一个订单 ID 对应的队列
SendResult result = rocketMQProducer.send(msg, new MessageQueueSelector() {
    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        String orderId = (String) arg;
        int index = Math.abs(orderId.hashCode()) % mqs.size();
        return mqs.get(index);
    }
}, orderId);
// 消费者:使用 MessageListenerOrderly 顺序消费
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
    for (MessageExt msg : msgs) {
        process(msg);
    }
    return ConsumeOrderlyStatus.SUCCESS;
});

② 全局有序(代价高)

如果你必须让所有消息全局有序,那只能让 Topic 只有一个分区(或队列)。这样并行度就是 1,吞吐量极低。只有在对顺序要求极其严苛且数据量很小的情况下才考虑,比如数据库 Binlog 的同步。

③ 业务层面补偿

有时可以不依赖 MQ 的顺序性,而是在消费端利用 版本号 或 状态机 来处理乱序消息。例如订单状态从“已创建”到“已支付”,如果你收到“已支付”但没找到“已创建”,可以把它暂存起来,等“已创建”到达后再处理,或者查询数据库校验状态。这种方式提高了系统的柔韧性。

收束: 顺序性的保证,本质上是把需要顺序的事情通过分区键收敛到一条线上。不要追求全局有序,那是性能和扩展性的天敌。用业务 ID 做 Key,让同一个实体的消息在单个分区内排队,是分布式系统中性价比最高的顺序方案。


🔁 2. 消息消费失败怎么处理?死信队列和重试机制怎么用?

消费失败不可怕,可怕的是失败后束手无策。MQ 提供了两大工具:重试和死信队列。

重试机制

当消费者处理消息失败时,可以让 MQ 自动重试。重试通常是指数退避的,比如 10s、30s、60s、120s... 直到达到最大重试次数。

死信队列

如果重试次数耗尽仍然失败,消息会被转移到 死信队列。死信队列是一个特殊的 Topic 或 Queue,用于存放那些无法被正常消费的消息。运维人员可以监听死信队列,进行人工排查和补偿。

RocketMQ 的处理 RocketMQ 的重试和死信是内建的。消费失败时,返回 ConsumeConcurrentlyStatus.RECONSUME_LATER,消息会进入重试队列,之后根据延迟等级重新投递。达到最大重试次数后,消息进入 DLQ(死信队列),Topic 为 %DLQ%ConsumerGroup

// RocketMQ 消费端:处理失败返回 RECONSUME_LATER
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
    for (MessageExt msg : msgs) {
        try {
            process(msg);
        } catch (Exception e) {
            // 返回 RECONSUME_LATER,消息将自动重试
            return ConsumeConcurrentlyStatus.RECONSUME_LATER;
        }
    }
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

你也可以在消息上设置 delayTimeLevel 来控制重试间隔,例如 msg.setDelayTimeLevel(3) 表示延迟 10 秒。

Kafka 的处理 Kafka 没有原生的自动重试机制,但可以通过 Spring-Kafka 的 SeekToCurrentErrorHandlerDeadLetterPublishingRecoverer 来实现。

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
        ConsumerFactory<String, String> consumerFactory, KafkaTemplate<String, String> kafkaTemplate) {
    ConcurrentKafkaListenerContainerFactory<String, String> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);

    // 死信恢复者:重试 3 次后,将消息发送到 DLQ Topic
    DeadLetterPublishingRecoverer recoverer =
        new DeadLetterPublishingRecoverer(kafkaTemplate,
            (record, ex) -> new TopicPartition(record.topic() + ".DLQ", record.partition()));

    // 错误处理:指数退避重试,最多 3 次
    ExponentialBackOff backOff = new ExponentialBackOff(1000L, 2.0);
    backOff.setMaxElapsedTime(9000L); // 最多重试 9 秒
    SeekToCurrentErrorHandler errorHandler =
        new SeekToCurrentErrorHandler(recoverer, backOff);
    errorHandler.setCommitRecovered(true);

    factory.setErrorHandler(errorHandler);
    return factory;
}

消息消费失败的完整处理流程:

  1. 消费者处理消息,若抛出异常,Spring-Kafka 会捕捉到。

  2. SeekToCurrentErrorHandler 执行重试,每次重试前等待一定时间(指数退避)。

  3. 若重试耗尽,DeadLetterPublishingRecoverer 将消息发送到 DLQ Topic(如 order-events.DLQ)。

  4. 运维人员监控 DLQ,排查问题后可以手工重发或丢弃。

消费异常的业务兜底

除了重试和死信,还可以在消费端实现业务级别的补偿。例如,对于已支付订单的积分发放失败,你可以在本地记录一条补偿任务,定时任务扫描补偿;或者利用对账机制,在日终核对订单状态和积分流水,发现不一致再补发。这比单纯依赖 MQ 重试更灵活。

收束: 重试是给瞬时故障一次悔改的机会,死信是给不可恢复的错误一个隔离的病房。用好它们,系统的抗风险能力会提升一个量级,但切莫忘记业务兜底——那才是最终的安全垫。


🚨 3. 线上消息积压了,怎么排查和优化?

消息积压意味着消费者的处理速度跟不上生产者的发送速度,消息在 MQ 中越堆越多,最终可能导致系统不可用。排查和优化需要有条不紊地进行。

🔍 排查步骤

第一步:查看消费积压量

# Kafka 查看消费者组的积压情况
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group your-group --describe

# 关注 LAG 列,表示每个分区还有多少消息未消费
# 如果 LAG 持续增长,说明积压严重

第二步:确认是生产突增还是消费变慢

  • 生产端:检查上游系统的 QPS 是否突然暴涨(例如营销活动开始)。如果是,可以通过限流或增加分区数来缓解,但分区数调整需要评估对顺序性的影响。

  • 消费端:检查消费者线程的 CPU、内存、GC 情况。是否有线程卡死、频繁 Full GC、或者外部调用(如 DB、Redis)超时导致处理变慢。

第三步:分析消费逻辑

查看消费者处理单条消息的耗时。可以增加日志记录处理时间,或者使用 APM 工具追踪。如果是因为消费逻辑中调用了一个慢 SQL 或没有设置超时的 HTTP 请求,那就是根因。

🛠️ 解决方案

① 临时应急:扩容消费者

如果消费能力不足,可以增加消费者实例数量,但前提是分区数足够(Kafka 的分区数决定了消费者的最大并行度)。如果是 RocketMQ,则直接加机器即可。

② 临时应急:跳过非关键消息

如果积压已经严重影响到核心业务,可以临时编写消费者程序,快速消费掉积压的消息,只处理核心消息(比如支付成功的),其余消息直接丢弃或转存到日志系统,待高峰期过后再慢慢补偿。

// 应急消费者:只处理核心消息,其余跳过
if (record.value().contains("\"type\":\"PAYMENT_SUCCESS\"")) {
    process(record);
} else {
    // 直接提交偏移量,不处理
    consumer.commitSync();
}

③ 长期优化:消费逻辑提速

  • 批量处理:将逐条更新数据库改为批量 update,减少 IO 次数。

  • 异步化:如果消费逻辑中有通知、记录日志等非关键操作,可以放入本地队列异步处理,让 poll 循环快速返回。

  • 合并写入:例如对同一个 Key 的数据做本地聚合,减少数据库写入频率。

④ 长期优化:增加资源与分区

  • 如果 CPU 或内存确实不足,增加消费者实例的硬件配置。

  • 如果分区数不足以支撑更多消费者并行处理,谨慎评估后增加分区数(注意:Kafka 分区数只能增加不能减少,且会影响顺序性)。

⑤ 限流与熔断

在生产者端引入限流,比如令牌桶算法,控制每秒发送到 MQ 的消息数量,防止突发流量冲垮消费者。同时,消费者端对下游依赖做熔断保护,当下游不可用时,暂停消费或快速失败,避免资源耗尽。

# 生产者限流示例:令牌桶
from time import time, sleep
class TokenBucket:
    def __init__(self, rate, capacity):
        self.rate = rate
        self.capacity = capacity
        self.tokens = capacity
        self.last = time()
    def consume(self, tokens=1):
        now = time()
        self.tokens = min(self.capacity, self.tokens + (now - self.last) * self.rate)
        self.last = now
        if self.tokens >= tokens:
            self.tokens -= tokens
            return True
        return False

bucket = TokenBucket(100, 200)  # 每秒100条,突发200条
while True:
    if bucket.consume():
        producer.send(msg)
    else:
        sleep(0.01)

收束: 线上积压的排查就像抢险——先看水位(LAG),再查上游是否暴雨(生产量),下游是否河道淤塞(消费慢)。临时措施是增泵(扩容)、引流(跳过非关键),长期治理是拓宽河道(优化逻辑、增加分区)和加固堤坝(限流熔断)。只有把这整套动作练熟,才能让消息队列成为可靠的泄洪渠,而不是决堤口。


6.4 消费者组与重平衡

问题 4:Kafka 消费者组重平衡是什么?如何避免频繁重平衡?

难度:⭐⭐⭐(Rebalance 原理、避免策略、性能调优)

1️⃣ Common Answer

Rebalance 就是消费者组的分区重新分配。比如有消费者加入或退出的时候就会触发。避免的话可以设置好 session timeout,不要让它频繁触发。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. Rebalance 的触发条件
- 消费者加入或离开(正常或超时)
- 消费者 session 超时(session.timeout.ms)
- 消费者主动离开(leaveGroup)
- 分区数变化(topic 扩容)
  1. Rebalance 的过程(影响为什么大)
Stop the World:重平衡期间所有消费者停止消费
1. Coordinator 发送 GroupHeartbeat
2. 选出一个 Leader Consumer
3. Leader 计算分区分配方案
4. 同步给所有消费者
5. 各消费者开始消费
  1. 频繁重平衡的原因

查看内嵌表格

  1. 关键参数调优
session.timeout.ms=45000    # 超时时间
heartbeat.interval.ms=15000  # 心跳间隔
max.poll.interval.ms=300000  # 最大处理间隔
max.poll.records=500         # 每次 poll 数量
  1. 静态成员机制(Kafka 2.3+)
group.instance.id:固定消费者 ID
重启不会触发 Rebalance
适合频繁重启的场景(如容器化)
  1. 实战经验 有一次容器滚动更新,每次重启都触发 Rebalance,导致消费中断。后来用了静态成员机制,配合合理的超时设置,基本避免了这个问题。

3️⃣ Key Differences

查看内嵌表格


6.5 消息追踪与可观测性

问题 5:如何构建消息系统的可观测性?消息追踪怎么做?

难度:⭐⭐⭐⭐(全链路追踪、监控体系、问题诊断)

1️⃣ Common Answer

可观测性的话,可以打日志,然后用 ELK 查。也可以加一些监控指标,用 Grafana 展示。消息追踪可以用一些开源工具,或者自己实现。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. 可观测性的三支柱
Metrics(指标):吞吐量、延迟、lag、错误率
Logs(日志):结构化日志 + 关键字段
Traces(追踪):全链路 TraceID 透传
  1. 消息追踪的核心设计
生产端:生成 TraceID → 注入消息头
Broker:记录消息元数据(可选)
消费端:提取 TraceID → 关联业务日志
链路:集成 SkyWalking/Jaeger 等 APM
  1. 关键监控指标

查看内嵌表格

  1. 问题诊断的完整链路
发现 lag 上涨
→ 看 Grafana 指标(哪个消费者、哪个分区)
→ 查 APM Trace(哪一步慢)
→ 搜日志(TraceID 关联)
→ 定位根因(代码/网络/下游)
  1. 生产环境的实践
- 统一日志格式:JSON + TraceID + 业务关键字段
- Prometheus + Grafana:实时监控
- SkyWalking:分布式追踪
- 自定义 Dashboard:消息健康度大盘
  1. 高阶能力

  2. 消息血缘:知道每条消息从哪来到哪去

  3. 延迟分析:P95/P99 延迟分位数

  4. 容量规划:根据历史数据预测峰值

3️⃣ Key Differences

查看内嵌表格


6.7 消息队列选型对比

问题 6:Kafka、RocketMQ、RabbitMQ 怎么选?各自适用什么场景?

难度:⭐⭐(消息队列选型、吞吐量、延迟、可靠性、生态对比)

1️⃣ Common Answer

Kafka 吞吐量最高,适合大数据场景;RabbitMQ 延迟最低,适合实时性要求高的场景;RocketMQ 是阿里开源的,功能比较全面,适合电商场景。我们公司用 Kafka 做日志收集,用 RocketMQ 做订单处理。

2️⃣ Impressive Answer

我会从这几个角度思考:选型核心维度、三款 MQ 的详细对比、场景化选型建议、实战决策经验

  1. 选型的核心维度

消息队列选型主要考虑 5 个维度:

  • 吞吐量:单位时间能处理的消息数量

  • 延迟:消息从发送到接收的时间间隔

  • 可靠性:消息不丢失、不重复的保障

  • 生态成熟度:社区活跃度、文档完善度、第三方工具支持

  • 运维成本:部署、监控、故障排查的复杂度

  • 三款 MQ 的详细对比

查看内嵌表格

  1. 场景化选型建议

  2. 日志采集、用户行为分析 → Kafka

  3. 原因:吞吐量极高,天然适配流式处理(Flink/Spark)
  4. 典型场景:埋点日志、监控指标采集、实时数仓

  5. 订单交易、支付场景 → RocketMQ

  6. 原因:支持事务消息,保障最终一致性;可靠性高,消息不丢失
  7. 典型场景:订单创建、支付回调、库存扣减

  8. 复杂路由、实时通信 → RabbitMQ

  9. 原因:Exchange-Queue 模型灵活,延迟极低
  10. 典型场景:即时通讯、复杂规则路由、微服务解耦

  11. 实战决策经验

去年我们做电商订单系统时,选型过程是这样的:

  1. 需求分析:订单创建后需要同步到 5 个下游系统(库存、物流、积分、风控、数据仓库),要求消息不丢失、支持事务消息

  2. 技术选型:排除了 Kafka(事务消息支持弱)和 RabbitMQ(吞吐量不够),最终选择 RocketMQ

  3. 落地效果:日均 500 万订单,消息延迟控制在 100ms 以内,通过事务消息保障了订单与库存的一致性

关键结论:没有最好的 MQ,只有最适合的。一定要结合业务场景、团队能力、现有技术栈综合决策。

3️⃣ Key Differences

查看内嵌表格


6.8 消息可靠性保证

问题 8:如何保证消息不丢失?从生产端、Broker、消费端三个环节分析

难度:⭐⭐⭐(消息可靠性、ACK机制、持久化、事务消息)

1️⃣ Common Answer

开启持久化、手动 ACK 就行了。生产端设置重试,Broker 开启刷盘,消费端确认后再提交 offset。

2️⃣ Impressive Answer

我会从这几个角度思考:全链路风险分析生产端保障机制Broker 可靠性配置消费端幂等设计Kafka/RocketMQ 实战配置

  1. 全链路风险分析
Producer → Broker → Consumer
   ↓         ↓          ↓
 网络抖动   磁盘故障    消费失败
 重试失败   副本丢失    ACK 超时

每个环节都可能丢消息,需要层层防护。

  1. 生产端保障

  2. 同步发送future.get() 确保消息到达 Broker

  3. 重试机制:配置 retries=3,指数退避避免雪崩

  4. 事务消息:RocketMQ 事务消息,本地事务 + 消息发送原子性

// Kafka 生产端配置
Properties props = new Properties();
props.put("acks", "all");  // 等待所有副本确认
props.put("retries", 3);
props.put("enable.idempotence", true);  // 开启幂等

// RocketMQ 事务消息
TransactionMQProducer producer = new TransactionMQProducer("group");
producer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 执行本地事务
        return success ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE;
    }
});
  1. Broker 保障

查看内嵌表格

  1. 消费端保障

  2. 手动 ACK:消费成功后再提交 offset

  3. 幂等消费:业务层去重(后续章节详述)

  4. 消费位移管理enable.auto.commit=false

// Kafka 消费端手动提交
props.put("enable.auto.commit", "false");
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        process(record);  // 业务处理
    }
    consumer.commitSync();  // 手动提交
}

// RocketMQ 手动 ACK
MessageListenerConcurrently listener = (msgs, context) -> {
    for (MessageExt msg : msgs) {
        processMessage(msg);
    }
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;  // 确认消费
};
  1. 实战配置组合拳

Kafka 高可靠性配置

# 生产端
acks=all
retries=3
enable.idempotence=true

# Broker 端
min.insync.replicas=2
num.replica.fetchers=2
log.flush.interval.messages=1

# 消费端
enable.auto.commit=false
auto.offset.reset=earliest

RocketMQ 高可靠性配置

# Broker 端
flushDiskType=SYNC_FLUSH
brokerRole=SYNC_MASTER
defaultTopicQueueNums=8

# 消费端
consumeThreadMin=20
consumeThreadMax=64

3️⃣ Key Differences

查看内嵌表格


6.9 消息幂等性设计

问题 9:消费者如何保证幂等?有哪些常见方案?

难度:⭐⭐⭐(幂等设计、去重表、状态机、Token机制)

1️⃣ Common Answer

用数据库唯一键或者 Redis 去重。

2️⃣ Impressive Answer

我会从这几个角度思考:幂等性必要性五种方案对比适用场景分析实战案例设计幂等与重试关系

  1. 为什么需要幂等

  2. At Least Once 语义:消息可能重复投递

  3. 网络重试:生产端重试、消费端重试

  4. Rebalance 重复消费:消费者组重平衡时重复投递

  5. 业务异常:ACK 超时导致重复消费

  6. 五种幂等方案对比

查看内嵌表格

  1. 方案详解

方案一:数据库唯一键

CREATE TABLE `order` (
  `order_id` VARCHAR(32) PRIMARY KEY,
  `status` TINYINT,
  `amount` DECIMAL(10,2),
  UNIQUE KEY `uk_order_no` (`order_no`)
);

-- 幂等插入
INSERT INTO `order` VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE status = status;

方案二:去重表

CREATE TABLE `dedup_log` (
  `biz_id` VARCHAR(64) PRIMARY KEY,
  `create_time` DATETIME
);

-- 事务保证
BEGIN;
INSERT INTO `dedup_log` VALUES (?, NOW());
-- 业务操作
UPDATE `account` SET balance = balance - ? WHERE user_id = ?;
COMMIT;

方案三:Redis setnx

public boolean process(String messageId) {
    String key = "dedup:" + messageId;
    // setnx + 过期时间
    boolean success = redis.setnx(key, "1", 300);  // 5分钟过期
    if (!success) {
        return false;  // 重复消息
    }
    try {
        doBusiness();
        return true;
    } finally {
        // 不删除 key,等待自动过期
    }
}

方案四:状态机

public enum OrderStatus {
    CREATED, PAID, SHIPPED, COMPLETED, CANCELLED;

    public boolean canTransitionTo(OrderStatus target) {
        switch (this) {
            case CREATED: return target == PAID || target == CANCELLED;
            case PAID: return target == SHIPPED;
            case SHIPPED: return target == COMPLETED;
            default: return false;
        }
    }
}

// 幂等更新
public void payOrder(String orderId) {
    Order order = orderDao.findById(orderId);
    if (!order.getStatus().canTransitionTo(OrderStatus.PAID)) {
        throw new IllegalStateException("订单状态不允许支付");
    }
    order.setStatus(OrderStatus.PAID);
    orderDao.update(order);
}

方案五:Token 机制

// 1. 生产端预分配 Token
String token = UUID.randomUUID().toString();
redis.setex("token:" + token, 300, "1");
message.setToken(token);

// 2. 消费端验证 Token
public boolean consume(Message msg) {
    String token = msg.getToken();
    String key = "token:" + token;
    if (redis.del(key) == 0) {
        return false;  // Token 已被消费
    }
    doBusiness(msg);
    return true;
}
  1. 实战:订单支付场景幂等设计

组合方案:唯一键 + 状态机

@Service
public class PaymentService {

    @Transactional
    public void payOrder(PaymentRequest request) {
        // 1. 唯一键去重(数据库层面)
        Order order = orderDao.findByOrderNo(request.getOrderNo());
        if (order == null) {
            throw new OrderNotFoundException();
        }

        // 2. 状态机检查(业务层面)
        if (!order.getStatus().canTransitionTo(OrderStatus.PAID)) {
            log.warn("订单状态不允许支付: {}", order.getStatus());
            return;  // 幂等返回
        }

        // 3. 执行支付
        paymentService.charge(request.getUserId(), request.getAmount());

        // 4. 更新状态
        order.setStatus(OrderStatus.PAID);
        order.setPayTime(LocalDateTime.now());
        orderDao.update(order);

        // 5. 发送后续消息(幂等发送)
        sendOrderPaidMessage(order);
    }
}
  1. 幂等与重试的关系

  2. 幂等是重试的前提:只有幂等才能安全重试

  3. 重试不等于幂等:重试是策略,幂等是保障

  4. 最佳实践:幂等设计 + 指数退避重试

// 重试 + 幂等
@Retryable(value = {Exception.class}, maxAttempts = 3,
    backoff = @Backoff(delay = 1000, multiplier = 2))
public void processMessage(Message msg) {
    // 幂等消费逻辑
    if (isProcessed(msg.getId())) {
        return;
    }
    doBusiness(msg);
    markProcessed(msg.getId());
}

3️⃣ Key Differences

查看内嵌表格


6.10 延迟消息与定时消息

问题 10:延迟消息的实现原理?RocketMQ 和 Kafka 分别怎么做?

难度:⭐⭐(延迟消息、时间轮、延迟级别、业务场景)

1️⃣ Common Answer

RocketMQ 支持延迟消息,发送消息时设置 delayLevel 就行,有 18 个延迟级别。Kafka 没有延迟消息功能,需要自己实现或者用外部调度。

2️⃣ Impressive Answer

我会从这几个角度思考:延迟消息的业务场景、RocketMQ 实现原理、Kafka 替代方案、实战落地案例

  1. 延迟消息的业务场景

延迟消息在业务中非常常见:

  • 订单超时关闭:下单后 30 分钟未支付自动取消

  • 延迟通知:活动开始前 10 分钟提醒用户

  • 定时任务:每天凌晨 2 点生成报表

  • 异步回调重试:调用失败后延迟 5 分钟重试

  • RocketMQ 延迟级别机制原理

RocketMQ 的延迟消息通过 18 个预定义延迟级别 实现:

查看内嵌表格

核心实现流程

Producer 发送消息(设置 delayLevel=14)
    → 消息写入 SCHEDULE_TOPIC_XXXX(而非目标 Topic)
    → Broker 定时任务每秒扫描延迟消息
    → 判断是否到达投递时间
        → 否:继续等待
        → 是:修改原始 Topic,清除延迟属性
    → 投递到真实 Topic,消费者正常消费

关键点

  • Producer 发送消息时设置 message.setDelayTimeLevel(14)(30 分钟)

  • 消息不会直接写入目标 Topic,而是写入 SCHEDULE_TOPIC_XXXX

  • Broker 有定时任务,每秒扫描延迟消息,判断是否到达投递时间

  • 到达时间后,将消息的 Topic 改回原始 Topic,清除延迟属性,投递到真实队列

  • Kafka 无内置延迟消息的替代方案

Kafka 原生不支持延迟消息,常用 3 种替代方案:

方案一:时间轮(Netty HashedWheelTimer)

// 生产者本地延迟发送
HashedWheelTimer timer = new HashedWheelTimer();
timer.newTimeout(timeout -> {
    producer.send(new ProducerRecord<>("order-topic", orderId));
}, 30, TimeUnit.MINUTES);
  • 优点:简单,无需 Kafka 支持

  • 缺点:生产者重启会丢失任务,不适合高可靠场景

方案二:外部调度系统(XXL-Job、Quartz)

// 定时任务扫描数据库
@Scheduled(fixedDelay = 60000)
public void scanTimeoutOrders() {
    List<Order> orders = orderMapper.findTimeoutOrders();
    orders.forEach(order -> {
        kafkaTemplate.send("order-cancel-topic", order.getId());
    });
}
  • 优点:可靠性高,任务持久化

  • 缺点:延迟精度依赖调度频率,数据库压力大

方案三:死信队列 + 延迟重试

// 消费失败后发送到延迟重试 Topic
@KafkaListener(topics = "order-topic")
public void handleOrder(Order order) {
    try {
        processOrder(order);
    } catch (Exception e) {
        // 发送到延迟重试 Topic,5 分钟后重试
        kafkaTemplate.send("order-retry-5m", order);
    }
}
  • 优点:利用 Kafka 原生能力

  • 缺点:需要多个 Topic 配合,复杂度高

  • 实战:订单 30 分钟未支付自动关闭

我们系统的完整方案:

// 1. Producer 发送延迟消息
Message message = new Message("order-topic",
    order.getId().getBytes());
message.setDelayTimeLevel(14); // 30 分钟
producer.send(message);

// 2. 消费者处理超时订单
@RocketMQMessageListener(topic = "order-topic",
    consumerGroup = "order-timeout-group")
public class OrderTimeoutConsumer implements
    RocketMQListener<Order> {

    @Override
    public void onMessage(Order order) {
        // 检查订单状态
        Order latest = orderService.getById(order.getId());
        if (latest.getStatus() == OrderStatus.UNPAID) {
            // 取消订单,释放库存
            orderService.cancelOrder(latest.getId());
        }
    }
}

关键优化

  • 幂等性保障:消费者处理前检查订单状态,避免重复取消

  • 监控告警:延迟消息堆积超过阈值时触发告警

  • 降级方案:RocketMQ 故障时降级到数据库定时任务扫描

3️⃣ Key Differences

查看内嵌表格


6.11 事务消息与最终一致性

问题 11:RocketMQ 事务消息的原理?如何实现分布式事务的最终一致性?

难度:⭐⭐⭐(事务消息、半消息、事务回查、最终一致性)

1️⃣ Common Answer

RocketMQ 有事务消息,先发半消息再提交。半消息就是暂时不能被消费者消费的消息,等本地事务执行成功后再提交,这样就能保证分布式事务的最终一致性。

2️⃣ Impressive Answer

我会从这几个角度思考:事务消息的完整流程、半消息的存储原理、事务回查机制、与本地消息表方案的对比、实战应用场景

  1. 事务消息的完整流程
Producer                    RocketMQ Broker              Consumer
   |                              |                         |
   |--- 1. 发送半消息 ----------->|                         |
   |                              |-- 写入 HALF_TOPIC       |
   |<-- 返回发送成功 -------------|                         |
   |                              |                         |
   |--- 2. 执行本地事务           |                         |
   |   (如:创建订单)             |                         |
   |                              |                         |
   |--- 3a. 事务成功 → Commit --->|                         |
   |                              |-- 投递到真实 Topic ----->|
   |                              |                         |-- 4. 正常消费
   |--- 3b. 事务失败 → Rollback ->|                         |
   |                              |-- 删除半消息            |
   |                              |                         |
   |   (Producer 崩溃/超时)       |                         |
   |                              |-- 4. 事务回查 --------->|
   |<-- 查询本地事务状态 ---------|                         |
   |--- 返回 Commit/Rollback ---->|                         |
  1. 半消息的存储原理

半消息并不会直接写入用户指定的 Topic,而是先写入一个系统内部 Topic:RMQ_SYS_TRANS_HALF_TOPIC。这个 Topic 对消费者不可见,确保了在本地事务执行完成前,消息不会被消费。当事务提交后,RocketMQ 会将消息从 RMQ_SYS_TRANS_HALF_TOPIC 移动到真实的 Topic 中,消费者才能消费。

  1. 事务回查机制

RocketMQ 通过事务回查机制解决生产者崩溃或超时的问题:

  • 回查条件:当生产者发送半消息后,超过一定时间(默认 60 秒)未收到 Commit/Rollback 响应时触发

  • 回查次数:默认最多回查 15 次

  • 回查间隔:每次间隔 60 秒

  • 回查逻辑:Broker 调用生产者实现的 TransactionListener.checkLocalTransaction() 方法,查询本地事务状态

  • 状态返回

  • COMMIT_MESSAGE:提交事务
  • ROLLBACK_MESSAGE:回滚事务
  • UNKNOWN:继续回查

  • 与本地消息表方案的对比

查看内嵌表格

  1. 实战:订单创建 + 积分发放的事务消息方案
// 订单服务生产者
public class OrderTransactionListener implements TransactionListener {

    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            // 执行本地事务:创建订单
            Long orderId = createOrder(msg);
            // 将订单ID存入消息属性,供回查使用
            msg.putUserProperty("orderId", orderId.toString());
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            log.error("创建订单失败", e);
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String orderId = msg.getUserProperty("orderId");
        // 查询订单状态
        Order order = orderService.getById(orderId);
        if (order != null && order.getStatus() == OrderStatus.PAID) {
            return LocalTransactionState.COMMIT_MESSAGE;
        } else if (order != null && order.getStatus() == OrderStatus.FAILED) {
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
        return LocalTransactionState.UNKNOWN;
    }
}

// 积分服务消费者
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "points-consumer")
public class PointsConsumer implements RocketMQListener<OrderMessage> {

    @Override
    public void onMessage(OrderMessage message) {
        // 发放积分
        pointsService.addPoints(message.getUserId(), message.getPoints());
    }
}

通过事务消息,我们实现了订单创建与积分发放的分布式事务一致性,即使订单服务崩溃,RocketMQ 也会通过事务回查机制确保最终一致性。

3️⃣ Key Differences

查看内嵌表格


6.12 消息队列的高可用架构

问题 12:Kafka/RocketMQ 如何保证高可用?主从同步、ISR 机制、选举策略是怎样的?

难度:⭐⭐⭐⭐(高可用架构、副本机制、ISR、Controller选举、Dledger、脑裂)

1️⃣ Common Answer

Kafka 和 RocketMQ 都有主从备份,主节点挂了从节点会顶上,保证消息不丢失。Kafka 用 ZooKeeper 做协调,RocketMQ 用 NameServer。

2️⃣ Impressive Answer

我会从这几个角度思考:Kafka 高可用全景、RocketMQ 高可用全景、两者高可用方案对比、脑裂问题及解决方案、实战故障恢复经验

  1. Kafka 高可用全景

1.1 副本机制(Leader/Follower、ISR/OSR)

Kafka 的每个 Topic 分区有多个副本,分为 Leader 和 Follower:

  • Leader:负责处理所有读写请求

  • Follower:异步复制 Leader 数据,不处理客户端请求

  • ISR(In-Sync Replicas):与 Leader 保持同步的副本集合,只有 ISR 中的副本才有资格被选为新 Leader

  • OSR(Out-of-Sync Replicas):与 Leader 同步滞后的副本

# 副本配置示例
replication.factor=3          # 总副本数
min.insync.replicas=2         # 最小同步副本数
replica.lag.time.max.ms=10000 # 副本最大滞后时间

1.2 ISR 收缩与扩展条件

ISR 的动态调整基于 replica.lag.time.max.ms 参数:

  • 收缩条件:Follower 在 replica.lag.time.max.ms 时间内未与 Leader 同步完成,则从 ISR 中移除

  • 扩展条件:OSR 中的副本追上 Leader 的 LEO(Log End Offset),则重新加入 ISR

1.3 min.insync.replicas 的作用

min.insync.replicas 确保消息的持久性和可用性:

  • 当 ISR 中的副本数 < min.insync.replicas 时,生产者会收到 NotEnoughReplicasException

  • 配合 acks=all 使用,确保消息至少写入 min.insync.replicas 个副本才认为发送成功

  • 防止在 ISR 收缩时降低数据可靠性

1.4 Controller 选举

ZooKeeper 模式(旧版)

  • Controller 是 Kafka 集群的"大脑",负责分区 Leader 选举、副本管理等

  • 通过 ZooKeeper 的临时节点选举,第一个启动的 Broker 成为 Controller

  • Controller 挂掉后,剩余 Broker 通过 ZooKeeper 重新选举

KRaft 模式(新版)

  • 移除 ZooKeeper,使用 Kafka 内部的 Raft 协议(Kafka Raft Metadata)

  • Controller 从 Quorum Controller 集群中选举,元数据存储在内部 Topic __cluster_metadata

  • 提高了性能和可维护性

1.5 Unclean Leader Election 的取舍

unclean.leader.election.enable 参数控制是否允许非 ISR 副本成为 Leader:

  • 开启(true):提高可用性,但可能丢失数据(OSR 副本数据不完整)

  • 关闭(false):保证数据一致性,但可能长时间不可用(ISR 副本全部挂掉)

  • 推荐配置:生产环境建议关闭,优先保证数据一致性

  • RocketMQ 高可用全景

2.1 主从同步(同步复制 vs 异步复制)

查看内嵌表格

<!-- 同步复制配置 -->
<brokerRole>SYNC_MASTER</brokerRole>

<!-- 异步复制配置 -->
<brokerRole>ASYNC_MASTER</brokerRole>

2.2 Dledger 自动选主(基于 Raft 协议)

RocketMQ 4.5+ 引入 Dledger(基于 Raft 协议)实现自动选主:

  • Raft 协议:通过日志复制和投票机制保证一致性

  • Leader 选举

  • 候选者发起选举,向其他节点请求投票
  • 收到多数节点投票(N/2 + 1)则成为 Leader
  • Leader 定期发送心跳维持统治地位

  • 故障转移:Leader 挂掉后,剩余节点自动重新选举

  • 脑裂防护:通过多数派投票机制避免脑裂

<!-- Dledger 配置 -->
<dledgerGroup>broker-a</dledgerGroup>
<dledgerPeers>n1@127.0.0.1:10911;n2@127.0.0.1:10912;n3@127.0.0.1:10913</dledgerPeers>
<dledgerSelfId>n1</dledgerSelfId>

2.3 NameServer 无状态设计

NameServer 是 RocketMQ 的注册中心,采用无状态设计:

  • 功能:管理 Broker 路由信息、Topic 路由信息

  • 无状态:NameServer 之间不通信,各自独立存储路由信息

  • 高可用:部署多个 NameServer 实例,任意一个挂掉不影响集群

  • 客户端:随机选择一个 NameServer 查询路由信息

  • 两者高可用方案对比

查看内嵌表格

  1. 脑裂问题及解决方案

脑裂场景:网络分区导致集群出现多个"主节点",导致数据不一致。

Kafka 解决方案

  • ZooKeeper 模式:通过临时节点的 Watch 机制,同一时刻只有一个 Controller

  • KRaft 模式:Raft 协议的多数派投票机制,确保只有获得多数票的节点才能成为 Leader

RocketMQ 解决方案

  • Dledger 模式:Raft 协议的多数派投票机制

  • 配置建议:集群节点数建议为奇数(3、5、7),避免偶数节点导致投票僵局

通用防护策略

  • 部署奇数个节点

  • 配置合理的超时时间

  • 监控网络分区和节点状态

  • 使用 fencing token(围栏令牌)机制

  • 实战:集群故障恢复经验

场景 1:Kafka Leader 挂掉

  1. Controller 检测到 Leader 挂掉

  2. 从 ISR 中选择 LEO 最高的副本作为新 Leader

  3. 更新 ZooKeeper 元数据

  4. 通知所有 Broker 和客户端

  5. 恢复时间:通常在 10-30 秒内

场景 2:RocketMQ Master 挂掉(Dledger 模式)

  1. 其他 Broker 检测到 Master 挂掉

  2. 触发 Raft 选举,选出新 Master

  3. 客户端自动切换到新 Master

  4. 恢复时间:通常在 5-15 秒内

优化建议

  • 监控 ISR/OSR 变化,及时处理副本滞后

  • 配置合理的 replica.lag.time.max.msmin.insync.replicas

  • 定期演练故障恢复流程

  • 使用 Prometheus + Grafana 监控集群健康度

3️⃣ Key Differences

查看内嵌表格


6.13 消息队列在微服务中的应用

问题 13:消息队列在微服务架构中承担什么角色?有哪些典型应用场景?

难度:⭐⭐(异步解耦、削峰填谷、事件驱动、CQRS、Saga)

1️⃣ Common Answer

消息队列在微服务中主要是解耦、异步、削峰。比如订单服务发消息,库存服务消费消息,这样就解耦了。还有流量削峰,防止服务被打挂。

2️⃣ Impressive Answer

我会从这几个角度思考:消息队列的四大核心角色、每个角色的业务场景与实现、事件驱动架构设计、实战编排案例

  1. 消息队列在微服务中的四大核心角色

查看内嵌表格

  1. 每个角色的业务场景与代码级实现

角色一:异步解耦 - 订单创建流程

传统同步调用的问题:

// 同步调用,链路过长,任何一环失败都会影响订单创建
public Order createOrder(OrderRequest request) {
    Order order = orderRepository.save(request);
    inventoryService.deduct(order);      // 耗时 50ms
    logisticsService.create(order);      // 耗时 100ms
    pointsService.add(order);            // 耗时 30ms
    riskService.check(order);            // 耗时 80ms
    return order;  // 总耗时 260ms+
}

异步解耦方案:

// 订单服务:只负责订单创建,后续流程异步化
@Transactional
public Order createOrder(OrderRequest request) {
    Order order = orderRepository.save(request);

    // 发送订单创建事件
    OrderCreatedEvent event = new OrderCreatedEvent(
        order.getId(), order.getUserId(),
        order.getItems(), order.getTotalAmount()
    );
    rocketMQTemplate.convertAndSend("order-created-topic", event);

    return order;  // 耗时 < 20ms
}

优势

  • 订单创建耗时从 260ms 降到 20ms,提升 10 倍+

  • 下游服务故障不影响订单创建(消息队列缓冲)

  • 易于扩展新的下游服务(新增监听器即可)

角色二:流量削峰 - 秒杀场景

// 秒杀接口:只承接请求,不直接处理
@PostMapping("/seckill/{productId}")
public Result seckill(@PathVariable Long productId) {
    // 1. 快速校验(Redis 缓存库存)
    if (!redisService.hasStock(productId)) {
        return Result.fail("库存不足");
    }

    // 2. 发送秒杀消息到 MQ
    SeckillMessage message = new SeckillMessage(
        productId, getUserId(), System.currentTimeMillis()
    );
    rocketMQTemplate.convertAndSend("seckill-topic", message);

    return Result.success("秒杀请求已提交");
}

削峰效果

  • 秒杀接口可以承接 10 万 QPS(只做 Redis 校验 + 发送消息)

  • 消费端限流到 1000 TPS,保护数据库不被打挂

  • 消息队列缓冲请求,削峰填谷

角色三:事件驱动 - 用户注册流程

// 用户服务:发布用户注册事件
@Service
public class UserService {

    public void register(UserRegisterRequest request) {
        User user = userRepository.save(request);

        // 发布用户注册事件
        UserRegisteredEvent event = new UserRegisteredEvent(
            user.getId(), user.getEmail(), user.getPhone()
        );
        rocketMQTemplate.convertAndSend("user-registered-topic", event);
    }
}

// 邮件服务:监听事件,发送欢迎邮件
// 数据服务:监听事件,初始化用户数据
// 营销服务:监听事件,发放新用户优惠券

事件驱动架构优势

  • 松耦合:新增服务只需监听事件,无需修改用户服务

  • 可扩展:同一事件可以被多个消费者处理

  • 最终一致性:通过事件保证各服务数据最终一致

角色四:数据同步 - CQRS 模式

// 命令端(写操作):更新数据库
@Service
public class OrderCommandService {

    @Transactional
    public void updateOrderStatus(Long orderId, OrderStatus status) {
        orderMapper.updateStatus(orderId, status);

        // 发送状态变更事件
        OrderStatusChangedEvent event = new OrderStatusChangedEvent(
            orderId, status, System.currentTimeMillis()
        );
        rocketMQTemplate.convertAndSend("order-status-changed", event);
    }
}

// 查询端(读操作):消费事件,更新 Redis 缓存和 Elasticsearch
  1. 事件驱动架构(EDA)的设计要点

核心原则

  • 事件即事实:事件不可变,代表已经发生的事实

  • 异步处理:事件发布后不等待消费者处理

  • 最终一致性:通过事件保证各服务数据最终一致

  • 幂等性保障:消费者必须能够处理重复事件

设计模式

  • 事件溯源:通过事件流重建状态(适用于金融、审计场景)

  • CQRS:命令与查询分离,读写使用不同的数据模型

  • Saga 模式:通过事件编排长事务(订单 → 库存 → 支付)

  • 实战:用消息队列实现订单 → 库存 → 物流的异步编排

用户 → 订单服务:创建订单
         ↓ 发送 OrderCreatedEvent
      消息队列
      库存服务:扣减库存
         ↓ 发送 InventoryDeductedEvent
      消息队列
      物流服务:创建运单
         ↓ 发送 WaybillCreatedEvent
      消息队列
      通知服务:发送发货通知 → 用户

关键优化

  • 事务消息:订单创建和事件发送使用 RocketMQ 事务消息,保证一致性

  • 幂等性:每个消费者处理前检查状态,避免重复处理

  • 重试机制:消费失败时自动重试,超过阈值后进入死信队列

  • 监控告警:消息堆积、消费延迟时触发告警

3️⃣ Key Differences

查看内嵌表格


6.15 Kafka 存储引擎深入

问题 15:Kafka 的零拷贝、顺序写、Page Cache 是如何协作实现高吞吐的?

难度:⭐⭐⭐⭐(零拷贝 sendfile、mmap、顺序 I/O、Page Cache 刷盘策略)

Common Answer

嗯...Kafka 的高吞吐主要是因为用了零拷贝和顺序写。零拷贝就是减少数据拷贝次数,不用从内核态拷到用户态。顺序写就是一直往后写,不用随机寻址。Page Cache 就是操作系统层面的缓存,可以加速读写。这些技术加在一起就让 Kafka 很快了。

Impressive Answer

我会从这几个角度思考:

  1. Kafka 写入路径:顺序追加写 Producer 发送消息 → 写入 Page Cache → 定期刷盘到磁盘(顺序追加写)

顺序写比随机写快 100 倍的原因:

  • 磁盘磁头不需要频繁移动寻道(寻道时间 5-10ms,是主要瓶颈)

  • 顺序写可以达到磁盘的物理极限速度(机械盘 100-200MB/s,SSD 更高)

  • 批量写入:Kafka 默认 batch.size=16KB,积攒一批后一次性写入

  • 零拷贝原理:传统 IO vs sendfile

传统 4 次拷贝流程:

磁盘 → 内核缓冲区(DMA 拷贝) → 用户缓冲区(CPU 拷贝) → 内核 Socket 缓冲区(CPU 拷贝) → 网卡(DMA 拷贝)

sendfile 2 次拷贝流程:

磁盘 → 内核缓冲区(DMA 拷贝) → 网卡(DMA 拷贝)

关键改进:

  • 减少 2 次 CPU 拷贝(内核 ↔ 用户空间)

  • 减少两次上下文切换

  • DMA(Direct Memory Access)直接在内核空间完成数据传输,不经过 CPU

  • Page Cache 的作用

OS 级别缓存:

  • Kafka 不自己实现缓存,依赖 OS 的 Page Cache

  • 写入时先写 Page Cache,异步刷盘(默认 log.flush.interval.messages=∞,不强制刷盘)

  • 读取时优先从 Page Cache 读取(热点数据命中率高)

预读机制:

  • OS 会预读后续数据到 Page Cache(顺序读场景下效果显著)

刷盘策略:

  • 同步刷盘:每次写入都刷盘(性能差,可靠性高)

  • 异步刷盘:由 OS 决定刷盘时机(性能好,可能丢数据)

  • Kafka 默认异步刷盘,通过 log.flush.interval.messages / log.flush.interval.ms 控制

  • mmap 内存映射

Index 文件使用 mmap 加速查找:

// 稀疏索引,通过 mmap 映射到内存
FileChannel indexChannel = new RandomAccessFile(indexFile, "r").getChannel();
MappedByteBuffer mappedBuffer = indexChannel.map(FileChannel.MapMode.READ_ONLY, 0, indexFile.length());

优势:

  • 索引文件直接映射到虚拟内存,避免 read() 系统调用

  • OS 自动按需加载页面

  • 适合小文件(Kafka Index 文件通常几 MB)

  • 与 RocketMQ 存储模型对比

查看内嵌表格

Trade-off:

  • Kafka:分区独立存储,扩展性好,但多 Topic 时磁盘碎片

  • RocketMQ:统一存储,磁盘利用率高,但单点 CommitLog 压力大

  • 实战配置调优

# 写入性能优化
batch.size=16384  # 批量发送大小
linger.ms=10      # 等待 10ms 积攒更多消息
compression.type=snappy  # 压缩减少网络 IO

# 刷盘策略(生产环境默认不强制刷盘)
log.flush.interval.messages=9223372036854775807  # 不强制刷盘
log.flush.interval.ms=9223372036854775807       # 不强制刷盘

# Segment 大小
log.segment.bytes=1073741824  # 1GB,减少文件数量

# Page Cache 优化(OS 级别,Kafka 无法直接配置)
# 建议:增加内存、减少 swap、使用 ext4/xfs 文件系统

Key Differences

查看内嵌表格


6.16 Kafka 分区分配策略

问题 16:Kafka 消费者的分区分配策略有哪些?Range、RoundRobin、Sticky、CooperativeSticky 有什么区别?

难度:⭐⭐⭐(分区分配策略、增量式重平衡、Cooperative Rebalance)

Common Answer

嗯...Kafka 有几种分区分配策略,Range 和 RoundRobin 比较常用。Range 是按范围分配,RoundRobin 是轮询分配。Sticky 是粘性分配,尽量保持之前的分配。CooperativeSticky 是改进版,重平衡时不会停止所有消费者。

Impressive Answer

  1. 四种策略的分配逻辑和算法

假设场景:3 个消费者(C0, C1, C2)消费 7 个分区(P0-P6)

Range 策略(默认)

按 Topic 分配,每个 Topic 独立计算:
7 个分区 / 3 个消费者 = 2 余 1
C0: P0, P1, P2 (3个)  # 前 2 个 + 余数 1
C1: P3, P4 (2个)
C2: P5, P6 (2个)

RoundRobin 策略

所有 Topic 的分区混合后轮询分配:
P0→C0, P1→C1, P2→C2, P3→C0, P4→C1, P5→C2, P6→C0
C0: P0, P3, P6 (3个)
C1: P1, P4 (2个)
C2: P2, P5 (2个)

Sticky 策略

首次分配类似 RoundRobin,但重平衡时尽量保持原有分配:
原则:最小化分区移动

CooperativeSticky 策略

增量式重平衡,不再 Stop the World:
- 只撤销需要重新分配的分区
- 分两阶段:ON_JOIN → ON_SYNC
  1. Range 的缺陷(多 Topic 时数据倾斜)

场景:2 个消费者,2 个 Topic(T1 有 3 个分区,T2 有 2 个分区)

T1 分配:
C0: T1-P0, T1-P1 (2个)
C1: T1-P2 (1个)

T2 分配:
C0: T2-P0 (1个)
C1: T2-P1 (1个)

总计:
C0: 3 个分区
C1: 2 个分区

数据倾斜导致:

  • C0 负载高,C1 负载低

  • 消费进度不一致

  • RoundRobin 的改进和局限

改进:

  • 解决了 Range 的数据倾斜问题

  • 多 Topic 下分配更均匀

局限:

  • 重平衡时所有分区都会重新分配

  • 不考虑之前的分配结果,可能频繁移动分区

  • Sticky 的核心思想(最小化分区移动)

原则:

  • 尽量保持消费者原有的分区分配

  • 只在必要时移动分区

示例:

初始分配:
C0: P0, P3, P6
C1: P1, P4
C2: P2, P5

C2 宕机后重平衡(Sticky):
C0: P0, P3, P6, P2  # P2 从 C2 迁移过来
C1: P1, P4, P5      # P5 从 C2 迁移过来

分区移动:2 个(比 RoundRobin 的 7 个少)
  1. CooperativeSticky(增量式重平衡)

传统重平衡(Eager):

1. 停止所有消费者(Stop the World)
2. 撤销所有分区
3. 重新计算分配
4. 恢复消费

增量式重平衡(Cooperative):

阶段 1(ON_JOIN):
- 新消费者加入,只撤销需要重新分配的分区
- 其他消费者继续消费

阶段 2(ON_SYNC):
- 完成剩余分区的迁移

优势:

  • 减少消费停顿时间

  • 降低重平衡对吞吐的影响

  • 配置方式代码示例和实战选择建议

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");

// 分区分配策略配置
props.put("partition.assignment.strategy",
    "org.apache.kafka.clients.consumer.StickyAssignor");

// 或者使用 CooperativeSticky(Kafka 2.4+)
props.put("partition.assignment.strategy",
    "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

实战选择建议:

  • 单 Topic、分区数能整除消费者数:Range(默认,简单)

  • 多 Topic、需要均匀分配:RoundRobin

  • 重平衡频繁、需要最小化分区移动:Sticky

  • Kafka 2.4+、高吞吐场景:CooperativeSticky(推荐)

Key Differences

查看内嵌表格


6.17 Kafka Exactly-Once 语义

问题 17:Kafka 的 Exactly-Once 语义是如何实现的?幂等生产者和事务的底层原理?

难度:⭐⭐⭐⭐(PID + Sequence Number、事务协调器、两阶段提交、isolation.level)

Common Answer

嗯...Kafka 支持 Exactly-Once 语义,通过幂等生产者和事务机制实现。幂等生产者就是保证消息不重复,事务就是保证一组消息要么都成功要么都失败。配置 enable.idempotence=true 就可以了。

Impressive Answer

  1. 三种语义对比

查看内嵌表格

  1. 幂等生产者原理

核心机制:PID(Producer ID)+ Sequence Number 去重

// 配置幂等生产者
props.put("enable.idempotence", true);  // 自动配置 acks=all, retries=Integer.MAX_VALUE

工作流程:

1. Producer 初始化时向 Broker 申请 PID
2. 每个 TopicPartition 维护一个 Sequence Number(从 0 开始递增)
3. 发送消息时携带 <PID, TopicPartition, Sequence Number>
4. Broker 端维护 PID 表,记录每个 TopicPartition 的最新 Sequence Number
5. 收到消息时检查:
   - 如果 Sequence Number == 预期值(当前最大值 + 1):接受
   - 如果 Sequence Number <= 当前最大值:重复消息,丢弃
   - 如果 Sequence Number > 当前最大值 + 1:数据丢失,抛出异常

Broker 端检测重复的伪代码:

public class PidManager {
    private Map<TopicPartition, Integer> sequenceNumbers = new ConcurrentHashMap<>();

    public boolean acceptMessage(ProducerIdAndEpoch producerId,
                                  TopicPartition partition,
                                  int sequenceNumber) {
        int currentMax = sequenceNumbers.getOrDefault(partition, -1);
        if (sequenceNumber == currentMax + 1) {
            sequenceNumbers.put(partition, sequenceNumber);
            return true;  // 接受
        } else if (sequenceNumber <= currentMax) {
            return false;  // 重复,丢弃
        } else {
            throw new OutOfOrderSequenceException();  // 数据丢失
        }
    }
}

限制:

  • 只能保证单次会话内的幂等(Producer 重启后 PID 会变化)

  • 只能保证单个 Partition 内的顺序和去重

  • 事务机制完整流程

核心组件:

  • TransactionCoordinator:事务协调器(类似 Group Coordinator)

  • _transactionstate Topic:存储事务状态(内部 Topic)

完整流程:

1. initTransactions
   Producer 向 TransactionCoordinator 注册,获取 PID

2. beginTransaction
   标记事务开始

3. send (多个消息)
   发送消息到多个 Partition
   消息带有事务标记(TRANSACTIONAL)

4. commitTransaction / abortTransaction
   向 TransactionCoordinator 提交/回滚事务

TransactionCoordinator 的角色:

  • 管理 Producer 的事务状态

  • 协调多个 Partition 的提交

  • 处理 _transactionstate Topic 的日志

  • 两阶段提交:PREPARE → COMMIT/ABORT

阶段 1:PREPARE

1. Producer 向 TransactionCoordinator 发送 COMMIT 请求
2. TransactionCoordinator 在 __transaction_state Topic 写入 PREPARE 状态
3. TransactionCoordinator 向所有涉及 Partition 的 Leader 发送 PREPARE 请求
4. Partition Leader 在本地日志写入事务标记(TRANSACTIONAL)

阶段 2:COMMIT / ABORT

COMMIT:
1. TransactionCoordinator 向 Partition Leader 发送 COMMIT 请求
2. Partition Leader 标记事务为 COMMITTED
3. TransactionCoordinator 在 __transaction_state Topic 写入 COMMIT_COMPLETE

ABORT:
1. TransactionCoordinator 向 Partition Leader 发送 ABORT 请求
2. Partition Leader 标记事务为 ABORTED
3. TransactionCoordinator 在 __transaction_state Topic 写入 ABORT_COMPLETE
  1. 消费端配合:isolation.level=read_committed
// 消费者配置
props.put("isolation.level", "read_committed");  // 只读已提交的消息
// 默认是 read_uncommitted,会读到未提交的消息

工作原理:

  • Consumer 过滤掉 ABORTED 的事务消息

  • 只消费 COMMITTED 的事务消息

  • 非事务消息正常消费

  • 端到端 Exactly-Once 的完整链路(consume-transform-produce 模式)

场景:从 Topic A 消费,处理后写入 Topic B

1. Consumer 消费 Topic A(isolation.level=read_committed)
2. 处理消息(业务逻辑)
3. Producer 写入 Topic B(事务模式)
4. 提交消费 Offset 和生产消息的事务(原子操作)

关键:
- 消费 Offset 也作为事务的一部分写入 __transaction_state Topic
- 事务提交后,Offset 和 Topic B 的消息同时生效
- 事务回滚后,Offset 和 Topic B 的消息都回滚
  1. 性能影响和实战取舍

性能开销:

  • 事务开销约 10-20%(相比普通生产)

  • 原因:

  • 需要与 TransactionCoordinator 交互
  • 两阶段提交增加网络往返
  • _transactionstate Topic 的写入开销

实战取舍:

  • At Least Once + 幂等消费:大多数场景够用,性能好

  • Exactly Once:金融、订单等强一致性场景,性能可接受

  • Java 代码示例

// Producer 配置
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("transactional.id", "my-transactional-id");  // 必须配置
props.put("enable.idempotence", true);  // 自动启用

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

// 初始化事务
producer.initTransactions();

try {
    // 开始事务
    producer.beginTransaction();

    // 发送多个消息
    producer.send(new ProducerRecord<>("topic1", "key1", "value1"));
    producer.send(new ProducerRecord<>("topic2", "key2", "value2"));

    // 提交事务
    producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
    // 严重错误,关闭 Producer
    producer.close();
} catch (KafkaException e) {
    // 可恢复错误,回滚事务
    producer.abortTransaction();
}

// Consumer 配置
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "localhost:9092");
consumerProps.put("group.id", "test-group");
consumerProps.put("isolation.level", "read_committed");  // 只读已提交消息

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);

Key Differences

查看内嵌表格

消息队列深入问题


6.18 RocketMQ 存储模型深入

问题 18:RocketMQ 的 CommitLog + ConsumeQueue + IndexFile 是如何协作的?与 Kafka 存储模型有什么区别?

难度:⭐⭐⭐⭐(CommitLog 统一存储、ConsumeQueue 逻辑队列、IndexFile 消息检索)

1️⃣ Common Answer 嗯...RocketMQ 的存储有三个部分,CommitLog 存消息,ConsumeQueue 存队列信息,IndexFile 做索引。Kafka 是每个分区独立存文件。RocketMQ 这样设计主要是为了性能吧,写的时候都写到一个文件里,读的时候通过 ConsumeQueue 找到 CommitLog 的位置。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. CommitLog:统一顺序写入

  2. 所有 Topic 的消息都顺序写入同一个 CommitLog 文件

  3. 单文件固定 1GB,写满后滚动创建新文件

  4. 顺序写入,磁盘 IO 性能最优,写入吞吐高

  5. ConsumeQueue:逻辑消费队列

  6. 每个 Topic 的每个 Queue 对应一个 ConsumeQueue

  7. 每条记录固定 20 字节:8 字节 CommitLog offset + 4 字节消息大小 + 8 字节 Tag hashcode

  8. 极小文件,可以全量加载到内存,快速定位消息

  9. IndexFile:基于 Key 的消息检索

  10. HashMap 结构,支持按消息 Key(如订单 ID)快速查找

  11. 索引文件包含 Key hashcode、CommitLog offset、时间戳等

  12. 用于业务查询、消息追踪场景

  13. 写入流程

Producer → Broker 写入 CommitLog(顺序写)
    异步构建 ConsumeQueue(更新 offset)
    异步构建 IndexFile(如果有 Key)
  1. 读取流程
Consumer → 读取 ConsumeQueue(内存)
    获取 CommitLog offset + size
    从 CommitLog 随机读取消息
  1. 与 Kafka 存储模型对比

查看内嵌表格

Trade-off 分析

  • RocketMQ 牺牲读取性能换取极致写入性能和消息检索能力

  • Kafka 牺牲检索能力换取简单架构和良好读取性能

  • RocketMQ 适合订单、支付等需要按 Key 查询的场景

  • Kafka 适合日志采集、流计算等纯消费场景

  • 实战经验

  • 文件清理策略:默认 72 小时,磁盘空间达到 85% 时强制删除,可配置 deleteWhenfileReservedTime

  • 磁盘水位告警:监控 CommitLogConsumeQueue 目录使用率,超过 80% 发告警

  • 性能优化flushDiskType 设置为 ASYNC_FLUSH 提升写入性能(牺牲可靠性)

3️⃣ Key Differences

查看内嵌表格


6.19 RocketMQ 消息过滤机制

问题 19:RocketMQ 的 Tag 过滤和 SQL92 过滤的原理与性能差异?

难度:⭐⭐⭐(Broker 端过滤 vs Consumer 端过滤、Tag hashcode、SQL92 表达式引擎)

1️⃣ Common Answer 嗯...RocketMQ 可以用 Tag 过滤,比如发送消息时指定 Tag,消费时订阅特定 Tag。还有 SQL92 过滤,可以写更复杂的条件,比如 a > 10 and b = 'xxx'。Tag 过滤应该更快一点,SQL92 更灵活。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. Tag 过滤原理

Broker 端初筛

  • ConsumeQueue 每条记录存储 8 字节 Tag hashcode

  • Consumer 拉取消息时,Broker 用 Tag hashcode 快速比对

  • Hashcode 冲突概率低(使用 MurmurHash),误判率极低

Consumer 端精确匹配

  • 消息实际拉取到 Consumer 后,再比对完整 Tag 字符串

  • 双重过滤确保准确性

代码示例

// Producer 发送消息
Message message = new Message("TopicTest", "TagA", "OrderID123", body);

// Consumer 订阅 Tag
consumer.subscribe("TopicTest", "TagA || TagB"); // 支持或运算
  1. SQL92 过滤原理

Broker 端解析与过滤

  • Consumer 订阅时发送 SQL92 表达式,如 TAGS is not null and TAGS in ('TagA', 'TagB')

  • Broker 解析表达式为过滤语法树

  • 每条消息拉取时,Broker 执行表达式引擎判断是否匹配

  • 支持的属性:TAGS、自定义属性(如 orderStatus > 0

代码示例

// Producer 发送带属性的消息
Message message = new Message("TopicTest", "TagA", body);
message.putUserProperty("orderStatus", "1");
message.putUserProperty("region", "HZ");

// Consumer 订阅 SQL92
consumer.subscribe("TopicTest",
    MessageSelector.bySql("TAGS is not null and orderStatus > 0 and region = 'HZ'"));
  1. 性能对比

查看内嵌表格

性能测试数据

  • Tag 过滤:单 Broker 可支持 10 万+ TPS

  • SQL92 过滤:单 Broker 约 2-3 万 TPS(受表达式引擎性能限制)

  • 使用场景

Tag 过滤适用场景

  • 消息分类简单(如订单状态:待支付、已支付、已发货)

  • 高吞吐场景(如日志采集、实时计算)

  • 消息属性固定,无需动态过滤

SQL92 过滤适用场景

  • 复杂业务规则(如 price > 100 and region = 'HZ' and status in (1,2)

  • 低吞吐、高精度过滤

  • 需要动态调整过滤条件(不重启 Consumer)

  • 配置方式

Broker 端配置(默认开启):

# 开启 SQL92 过滤支持
enablePropertyFilter=true

# SQL92 过滤最大表达式长度
maxMessageSize=4M

Java 代码示例

// Tag 过滤
consumer.subscribe("TopicOrder", "TagPaid || TagShipped");

// SQL92 过滤
consumer.subscribe("TopicOrder",
    MessageSelector.bySql("TAGS in ('TagPaid', 'TagShipped') and price > 100"));
  1. 实战:过滤导致的消费倾斜问题

问题现象

  • 某个 Tag 的消息量特别大,导致消费该 Tag 的 Consumer 负载过高

  • 其他 Tag 的 Consumer 空闲,整体消费效率低

解决方案

// 方案1:增加队列数和 Consumer 数,按 Tag 分组消费
// TagA 消费者订阅 TopicOrder:Queue0-3,TagB 消费者订阅 TopicOrder:Queue4-7

// 方案2:使用 SQL92 动态调整过滤条件
// 根据消费延迟动态调整过滤表达式,平衡负载

// 方案3:消息分流(推荐)
// 将高频 Tag 的消息发送到独立 Topic
producer.send(new Message("TopicOrder_HighVolume", "TagA", body));

3️⃣ Key Differences

查看内嵌表格


6.20 消息队列生产端性能调优

问题 20:Kafka/RocketMQ 的生产端如何调优?batch、压缩、acks、linger.ms 等参数怎么配?

难度:⭐⭐⭐(批量发送、压缩算法对比、acks 与吞吐的 trade-off、异步发送)

1️⃣ Common Answer 嗯...生产端调优主要是批量发送和压缩。batch.size 设置大一点,linger.ms 设置几百毫秒,这样能攒更多消息一起发。压缩用 gzip 或者 snappy。acks 设置成 1 或者 all,保证可靠性。异步发送也很快。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. 批量发送:batch.size + linger.ms 的配合逻辑

Kafka 配置

# 单个批次最大字节数(默认 16KB,建议 32KB-1MB)
batch.size=32768

# 等待批次攒满的最大时间(默认 0ms,建议 10-100ms)
linger.ms=20

配合逻辑

  • linger.ms 控制等待时间,batch.size 控制批次大小

  • 两者任一达到阈值即发送(先到先发)

  • 高吞吐场景:linger.ms=50-100msbatch.size=1MB

  • 低延迟场景:linger.ms=0-10msbatch.size=32KB

RocketMQ 配置

// 批量发送(默认 1KB,建议 4KB-1MB)
producer.setRetryTimesWhenSendFailed(2); // 失败重试
producer.setSendMsgTimeout(3000); // 发送超时 3s

// 批量发送代码示例
List<Message> messages = new ArrayList<>();
for (int i = 0; i < 100; i++) {
    messages.add(new Message("TopicTest", ("Hello" + i).getBytes()));
}
SendResult result = producer.send(messages); // 批量发送
  1. 压缩算法对比

查看内嵌表格

Kafka 配置

# 压缩算法(默认 none)
compression.type=lz4  # 推荐:LZ4 或 Snappy

RocketMQ 配置

// 消息体超过阈值时压缩(默认 4KB,建议 1KB-10KB)
producer.setCompressMsgBodyOverHowmuch(1024);

// 压缩算法(默认 ZIP,建议 LZ4)
producer.setCompressLevel(5); // ZIP 压缩级别(0-9)

实战经验

  • 日志采集场景:使用 Zstd,压缩率高,节省带宽

  • 实时计算场景:使用 LZ4,低延迟,高吞吐

  • 订单交易场景:使用 Snappy,平衡性能与可靠性

  • acks 参数:吞吐与可靠性的 trade-off

Kafka 配置

# 0:不等待确认(最快,可能丢消息)
# 1:等待 Leader 确认(折中)
# all/-1:等待所有 ISR 确认(最可靠)
acks=1

查看内嵌表格

RocketMQ 配置

// 同步发送(等待 Broker 确认)
SendResult result = producer.send(message); // 默认同步

// 异步发送(不等待确认,回调处理)
producer.send(message, new SendCallback() {
    @Override
    public void onSuccess(SendResult sendResult) {
        // 发送成功回调
    }

    @Override
    public void onException(Throwable e) {
        // 发送失败回调
    }
});
  1. 异步发送 + 回调:提升吞吐的最佳实践

Kafka 代码示例

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("acks", "1");
props.put("retries", 3); // 失败重试
props.put("batch.size", 32768);
props.put("linger.ms", 20);
props.put("buffer.memory", 33554432); // 缓冲区大小 32MB
props.put("compression.type", "lz4");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", StringSerializer.class.getName());

Producer<String, String> producer = new KafkaProducer<>(props);

// 异步发送 + 回调
producer.send(new ProducerRecord<>("TopicTest", "key", "value"),
    new Callback() {
        @Override
        public void onCompletion(RecordMetadata metadata, Exception exception) {
            if (exception != null) {
                // 处理失败
                log.error("Send failed", exception);
            } else {
                // 处理成功
                log.info("Send success, offset: {}", metadata.offset());
            }
        }
    });

RocketMQ 代码示例

DefaultMQProducer producer = new DefaultMQProducer("ProducerGroup");
producer.setNamesrvAddr("localhost:9876");
producer.setRetryTimesWhenSendAsyncFailed(2); // 异步重试
producer.start();

// 异步发送 + 回调
Message message = new Message("TopicTest", "TagA", "Hello".getBytes());
producer.send(message, new SendCallback() {
    @Override
    public void onSuccess(SendResult sendResult) {
        log.info("Send success, msgId: {}", sendResult.getMsgId());
    }

    @Override
    public void onException(Throwable e) {
        log.error("Send failed", e);
    }
});
  1. RocketMQ 生产端调优参数
// 发送超时时间(默认 3s,建议 1-5s)
producer.setSendMsgTimeout(3000);

// 消息体超过阈值时压缩(默认 4KB)
producer.setCompressMsgBodyOverHowmuch(1024);

// 同步发送失败重试次数(默认 2)
producer.setRetryTimesWhenSendFailed(2);

// 异步发送失败重试次数(默认 2)
producer.setRetryTimesWhenSendAsyncFailed(2);

// 最大消息体大小(默认 4MB,建议 1-10MB)
producer.setMaxMessageSize(4 * 1024 * 1024);
  1. 实战:从 5000 TPS 调优到 50000 TPS 的参数组合

初始配置(5000 TPS)

# Kafka
batch.size=16384
linger.ms=0
acks=all
compression.type=none

优化配置(50000 TPS)

# Kafka
batch.size=1048576          # 1MB
linger.ms=50                # 50ms
acks=1                      # 降低可靠性提升吞吐
compression.type=lz4        # 启用压缩
buffer.memory=67108864      # 64MB 缓冲区
max.in.flight.requests.per.connection=10  # 提升并发

优化效果

  • 批量发送提升 5 倍:16KB → 1MB

  • 延迟攒批提升 2 倍:0ms → 50ms

  • 压缩节省 50% 带宽:none → LZ4

  • 降低确认延迟:acks=all → acks=1

  • 总体提升:10 倍吞吐(5000 → 50000 TPS)

RocketMQ 优化配置

// 批量发送(100 条/批次)
List<Message> messages = new ArrayList<>();
for (int i = 0; i < 100; i++) {
    messages.add(new Message("TopicTest", body));
}
producer.send(messages);

// 异步发送 + 回调
producer.send(message, callback);

3️⃣ Key Differences

查看内嵌表格


6.21 消费端性能瓶颈排查与优化

问题 21:消费端处理速度跟不上怎么办?如何定位和优化消费者的性能瓶颈?

难度:⭐⭐⭐(消费线程模型、批量消费、异步落库、pull vs push)

1️⃣ Common Answer 嗯...消费慢的话,可以增加消费者数量,或者增加分区数。批量消费也能提升性能,一次处理多条消息。异步处理也可以,消费完放到队列里,后台线程慢慢处理数据库。Kafka 是 pull 模式,RocketMQ 是 push 模式。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. 消费线程模型对比

Kafka:单线程 poll + 多线程处理

// 主线程:单线程 poll 消息
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

    // 多线程处理(线程池)
    executorService.submit(() -> {
        for (ConsumerRecord<String, String> record : records) {
            processMessage(record);
        }
    });
}

// 注意:多线程处理时需要手动提交 offset
consumer.commitSync();

RocketMQ:线程池消费

// 默认使用线程池消费(ConsumeThreadMin=20, ConsumeThreadMax=20)
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroup");
consumer.setConsumeThreadMin(20);  // 最小线程数
consumer.setConsumeThreadMax(20);  // 最大线程数
consumer.setPullBatchSize(32);     // 每次 pull 的消息数

consumer.registerMessageListener(new MessageListenerConcurrently() {
    @Override
    public ConsumeConcurrentlyStatus consumeMessage(
        List<MessageExt> msgs,
        ConsumeConcurrentlyContext context
    ) {
        // msgs 是批量消息(默认 1 条,可配置 pullBatchSize)
        for (MessageExt msg : msgs) {
            processMessage(msg);
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});

对比

  • Kafka:单线程 poll,避免并发提交 offset 问题,处理逻辑需要多线程

  • RocketMQ:线程池消费,框架自动管理并发和 offset 提交

  • 批量消费:减少网络开销和事务提交次数

Kafka 批量消费配置

# 每次 poll 最大拉取条数(默认 500)
max.poll.records=500

# 批量处理代码
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
List<ConsumerRecord<String, String>> batch = new ArrayList<>(records);

// 批量插入数据库(减少事务提交次数)
jdbcTemplate.batchUpdate("INSERT INTO t_order VALUES (?, ?)", batch);

RocketMQ 批量消费配置

consumer.setPullBatchSize(32);  // 每次 pull 32 条

// 批量处理
consumer.registerMessageListener(new MessageListenerConcurrently() {
    @Override
    public ConsumeConcurrentlyStatus consumeMessage(
        List<MessageExt> msgs,
        ConsumeConcurrentlyContext context
    ) {
        // msgs 可能包含多条消息(根据 pullBatchSize)
        List<Order> orders = parseOrders(msgs);
        orderService.batchInsert(orders);  // 批量插入
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});

优化效果

  • 单条消费:1000 条消息 = 1000 次网络 IO + 1000 次事务提交

  • 批量消费(100 条/批次):1000 条消息 = 10 次网络 IO + 10 次事务提交

  • 性能提升:10-50 倍

  • 异步落库:消费 → 内存队列 → 批量写入数据库

架构设计

MQ Consumer → 内存队列(BlockingQueue) → 后台线程批量写入 DB

代码示例

// 内存队列
BlockingQueue<Order> orderQueue = new LinkedBlockingQueue<>(10000);

// 消费者:快速消费到内存队列
consumer.registerMessageListener(new MessageListenerConcurrently() {
    @Override
    public ConsumeConcurrentlyStatus consumeMessage(
        List<MessageExt> msgs,
        ConsumeConcurrentlyContext context
    ) {
        for (MessageExt msg : msgs) {
            Order order = parseOrder(msg);
            orderQueue.offer(order);  // 非阻塞写入队列
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});

// 后台线程:批量写入数据库
ScheduledExecutorService executor = Executors.newScheduledThreadPool(1);
executor.scheduleAtFixedRate(() -> {
    List<Order> batch = new ArrayList<>(500);
    orderQueue.drainTo(batch, 500);  // 最多取 500 条

    if (!batch.isEmpty()) {
        orderService.batchInsert(batch);  // 批量插入
    }
}, 100, 100, TimeUnit.MILLISECONDS);  // 每 100ms 执行一次

优化效果

  • 消费速度提升 5-10 倍(无需等待 DB 写入)

  • DB 写入批量处理,减少事务提交次数

  • 整体消费延迟从分钟级优化到秒级

  • pull vs push 模型对比

查看内嵌表格

  1. 消费者并发度调优

Kafka 调优

# 分区数 = 消费者数(每个消费者消费多个分区)
# 例如:8 个分区,4 个消费者,每个消费者消费 2 个分区

# 消费者数配置
num.consumer.instances=4

# 每个 Consumer 线程数(可选)
consumer.threads.per.instance=2

RocketMQ 调优

// 消费线程数(默认 20)
consumer.setConsumeThreadMin(20);
consumer.setConsumeThreadMax(50);  // 最大 50 线程

// 每次 pull 消息数(默认 32)
consumer.setPullBatchSize(100);  // 每次 pull 100 条

// 消费最小间隔(默认 0)
consumer.setPullInterval(0);  // 无间隔,最快消费

调优原则

  • 消费线程数 = CPU 核心数 × 2

  • pullBatchSize = 消费线程数 × 5-10(确保线程不空闲)

  • 消费者数 ≤ 分区数(Kafka)

  • 实战:消费延迟从分钟级优化到秒级的完整方案

问题现象

  • 消费延迟:10 分钟

  • 消息堆积:100 万条

  • 消费者:2 个实例,单实例 20 线程

优化方案

Step 1:增加消费者数量

// 从 2 个实例增加到 8 个实例
// 每个实例 20 线程,总并发 160 线程

Step 2:启用批量消费

// Kafka
max.poll.records=500

// RocketMQ
consumer.setPullBatchSize(100);  // 每次 pull 100 条

Step 3:异步落库

// 消费 → 内存队列 → 批量写入 DB
// 见上文代码示例

Step 4:优化 DB 写入

// 批量插入(500 条/批次)
jdbcTemplate.batchUpdate(sql, batch);

// 异步写入(非核心业务)
// 或者使用 MQ 二次分发到专门的 DB 写入服务

优化效果

  • 消费延迟:10 分钟 → 30 秒

  • 消息堆积:100 万条 → 0

  • 吞吐量:2000 TPS → 50000 TPS

监控指标

  • 消费延迟(Consumer Lag)

  • 消息堆积量(Message Backlog)

  • 消费 TPS(Messages per Second)

  • 消费线程 CPU 使用率

  • DB 写入耗时

3️⃣ Key Differences

查看内嵌表格


6.22 Saga 分布式事务编排

问题 22:消息队列如何实现 Saga 分布式事务编排?与 TCC、本地消息表方案有什么区别?

难度:⭐⭐⭐⭐(Saga 编排式 vs 协同式、补偿事务、与 TCC 对比、Seata 集成)

1️⃣ Common Answer 嗯...Saga 是一种分布式事务模式,通过一系列本地事务和补偿操作实现最终一致性。TCC 是 Try-Confirm-Cancel,需要三个接口。本地消息表是先写消息再发 MQ。Saga 可以用消息队列实现,每个服务消费上游消息,执行本地事务,再发下游消息。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. Saga 模式的两种实现

编排式(Orchestration)

  • 由一个中心协调器(Saga Coordinator)控制整个流程

  • 协调器调用各个服务的正向操作和补偿操作

  • 优点:流程清晰、易于监控、补偿逻辑集中

  • 缺点:协调器单点、耦合度高

协同式(Choreography)

  • 没有中心协调器,各服务通过事件驱动协作

  • 每个服务消费上游事件,执行本地事务,发布下游事件

  • 优点:去中心化、松耦合、扩展性好

  • 缺点:流程分散、难以追踪、补偿逻辑分散

  • 基于消息队列的协同式 Saga

架构设计

订单服务 → MQ → 库存服务 → MQ → 支付服务 → MQ → 物流服务
   ↓          ↓          ↓          ↓
补偿订单   补偿库存   补偿支付   补偿物流

流程说明

  1. 订单服务创建订单 → 发布 OrderCreated 事件

  2. 库存服务消费 OrderCreated → 扣减库存 → 发布 InventoryReserved 事件

  3. 支付服务消费 InventoryReserved → 扣款 → 发布 PaymentCompleted 事件

  4. 物流服务消费 PaymentCompleted → 创建物流单 → 发布 ShipmentCreated 事件

补偿流程(任一步骤失败):

  1. 支付失败 → 发布 PaymentFailed 事件

  2. 库存服务消费 PaymentFailed → 回滚库存 → 发布 InventoryReleased 事件

  3. 订单服务消费 InventoryReleased → 取消订单 → 发布 OrderCancelled 事件

  4. 补偿事务设计

正向操作 + 逆向补偿

// 正向操作:扣减库存
public void reserveInventory(String orderId, List<OrderItem> items) {
    // 执行本地事务
    transactionTemplate.execute(status -> {
        for (OrderItem item : items) {
            inventoryRepository.decrease(item.getProductId(), item.getQuantity());
        }
        // 记录操作日志(用于补偿)
        sagaLogRepository.save(new SagaLog(orderId, "INVENTORY_RESERVED", items));
        return null;
    });

    // 发布下游事件
    eventPublisher.publish(new InventoryReservedEvent(orderId, items));
}

// 逆向补偿:回滚库存
public void releaseInventory(String orderId) {
    // 查询操作日志
    SagaLog log = sagaLogRepository.findByOrderIdAndOperation(orderId, "INVENTORY_RESERVED");
    List<OrderItem> items = log.getItems();

    // 执行补偿事务
    transactionTemplate.execute(status -> {
        for (OrderItem item : items) {
            inventoryRepository.increase(item.getProductId(), item.getQuantity());
        }
        return null;
    });

    // 发布补偿事件
    eventPublisher.publish(new InventoryReleasedEvent(orderId));
}

补偿原则

  • 补偿操作必须是幂等的(可能重复执行)

  • 补偿操作必须能回滚正向操作的所有影响

  • 补偿操作本身也可能失败,需要重试或人工介入

  • 与 TCC 的对比

查看内嵌表格

对比示例

// Saga:扣减库存(直接扣减,失败时补偿)
public void reserveInventory(String orderId, int productId, int quantity) {
    inventoryRepository.decrease(productId, quantity); // 直接扣减
}

// TCC:Try 阶段(预留库存)
public boolean tryReserveInventory(String orderId, int productId, int quantity) {
    return inventoryRepository.reserve(productId, quantity); // 预留
}

// TCC:Confirm 阶段(确认扣减)
public boolean confirmReserveInventory(String orderId, int productId, int quantity) {
    return inventoryRepository.decreaseReserved(productId, quantity); // 扣减预留
}

// TCC:Cancel 阶段(释放预留)
public boolean cancelReserveInventory(String orderId, int productId, int quantity) {
    return inventoryRepository.releaseReserved(productId, quantity); // 释放预留
}
  1. 与本地消息表的对比

查看内嵌表格

本地消息表示例

// 本地消息表实现
@Transactional
public void createOrder(Order order) {
    // 1. 写入订单
    orderRepository.save(order);

    // 2. 写入本地消息表
    messageRepository.save(new LocalMessage(
        order.getId(),
        "ORDER_CREATED",
        JSON.toJSONString(order)
    ));
}

// 定时任务:扫描消息表并发送 MQ
@Scheduled(fixedDelay = 5000)
public void scanAndSendMessages() {
    List<LocalMessage> messages = messageRepository.findByStatusAndSendTimeBefore(
        "PENDING",
        new Date()
    );

    for (LocalMessage message : messages) {
        try {
            rocketMQTemplate.send(message.getTopic(), message.getContent());
            message.setStatus("SENT");
            messageRepository.save(message);
        } catch (Exception e) {
            log.error("Send message failed", e);
        }
    }
}
  1. 实战:订单 → 库存 → 支付的 Saga 编排完整 Java 代码示例

订单服务

@Service
public class OrderService {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Autowired
    private OrderRepository orderRepository;

    @Autowired
    private SagaLogRepository sagaLogRepository;

    @Transactional
    public void createOrder(CreateOrderRequest request) {
        // 1. 创建订单
        Order order = new Order(request.getUserId(), request.getItems());
        orderRepository.save(order);

        // 2. 记录 Saga 日志
        sagaLogRepository.save(new SagaLog(
            order.getId(),
            "ORDER_CREATED",
            JSON.toJSONString(order)
        ));

        // 3. 发布订单创建事件
        rocketMQTemplate.convertAndSend(
            "order-created-topic",
            new OrderCreatedEvent(order.getId(), order.getItems())
        );
    }

    @RocketMQMessageListener(
        topic = "inventory-released-topic",
        consumerGroup = "order-consumer-group"
    )
    public class OrderCompensationListener implements RocketMQListener<InventoryReleasedEvent> {

        @Override
        public void onMessage(InventoryReleasedEvent event) {
            // 补偿:取消订单
            Order order = orderRepository.findById(event.getOrderId());
            order.setStatus("CANCELLED");
            orderRepository.save(order);

            // 记录补偿日志
            sagaLogRepository.save(new SagaLog(
                event.getOrderId(),
                "ORDER_CANCELLED",
                JSON.toJSONString(order)
            ));
        }
    }
}

库存服务

@Service
public class InventoryService {

    @Autowired
    private InventoryRepository inventoryRepository;

    @Autowired
    private SagaLogRepository sagaLogRepository;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @RocketMQMessageListener(
        topic = "order-created-topic",
        consumerGroup = "inventory-consumer-group"
    )
    public class OrderCreatedListener implements RocketMQListener<OrderCreatedEvent> {

        @Override
        public void onMessage(OrderCreatedEvent event) {
            try {
                // 1. 扣减库存
                transactionTemplate.execute(status -> {
                    for (OrderItem item : event.getItems()) {
                        inventoryRepository.decrease(
                            item.getProductId(),
                            item.getQuantity()
                        );
                    }
                    return null;
                });

                // 2. 记录 Saga 日志
                sagaLogRepository.save(new SagaLog(
                    event.getOrderId(),
                    "INVENTORY_RESERVED",
                    JSON.toJSONString(event.getItems())
                ));

                // 3. 发布库存预留事件
                rocketMQTemplate.convertAndSend(
                    "inventory-reserved-topic",
                    new InventoryReservedEvent(event.getOrderId(), event.getItems())
                );

            } catch (InsufficientStockException e) {
                // 库存不足,发布失败事件
                rocketMQTemplate.convertAndSend(
                    "inventory-failed-topic",
                    new InventoryFailedEvent(event.getOrderId(), e.getMessage())
                );
            }
        }
    }

    @RocketMQMessageListener(
        topic = "payment-failed-topic",
        consumerGroup = "inventory-consumer-group"
    )
    public class PaymentFailedListener implements RocketMQListener<PaymentFailedEvent> {

        @Override
        public void onMessage(PaymentFailedEvent event) {
            // 补偿:回滚库存
            SagaLog log = sagaLogRepository.findByOrderIdAndOperation(
                event.getOrderId(),
                "INVENTORY_RESERVED"
            );
            List<OrderItem> items = JSON.parseArray(log.getContent(), OrderItem.class);

            transactionTemplate.execute(status -> {
                for (OrderItem item : items) {
                    inventoryRepository.increase(
                        item.getProductId(),
                        item.getQuantity()
                    );
                }
                return null;
            });

            // 发布库存释放事件
            rocketMQTemplate.convertAndSend(
                "inventory-released-topic",
                new InventoryReleasedEvent(event.getOrderId())
            );
        }
    }
}

支付服务

@Service
public class PaymentService {

    @Autowired
    private PaymentRepository paymentRepository;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @RocketMQMessageListener(
        topic = "inventory-reserved-topic",
        consumerGroup = "payment-consumer-group"
    )
    public class InventoryReservedListener implements RocketMQListener<InventoryReservedEvent> {

        @Override
        public void onMessage(InventoryReservedEvent event) {
            try {
                // 1. 扣款
                Payment payment = new Payment(
                    event.getOrderId(),
                    calculateTotalAmount(event.getItems())
                );
                paymentRepository.save(payment);

                // 2. 发布支付完成事件
                rocketMQTemplate.convertAndSend(
                    "payment-completed-topic",
                    new PaymentCompletedEvent(event.getOrderId(), payment.getId())
                );

            } catch (PaymentException e) {
                // 支付失败,发布失败事件
                rocketMQTemplate.convertAndSend(
                    "payment-failed-topic",
                    new PaymentFailedEvent(event.getOrderId(), e.getMessage())
                );
            }
        }
    }
}

3️⃣ Key Differences

查看内嵌表格


6.23 消息队列实现缓存一致性

问题 23:如何用消息队列实现分布式缓存更新的最终一致性?Canal + MQ 方案怎么做?

难度:⭐⭐⭐(缓存一致性、Canal binlog 订阅、延迟双删的消息化替代)

1️⃣ Common Answer 嗯...缓存一致性问题,先更新 DB 再删缓存,或者先删缓存再更新 DB,都有问题。可以用 Canal 监听 binlog,然后发消息到 MQ,消费者收到消息后删除缓存。延迟双删也可以用 MQ 实现,发一个延迟消息做二次删除。

2️⃣ Impressive Answer 我会从这几个角度思考:

  1. 缓存一致性的经典问题

方案一:先更新 DB,再删缓存

1. 更新数据库
2. 删除缓存

问题

  • 步骤 1 成功,步骤 2 失败 → 缓存是旧数据

  • 并发场景:线程 A 更新 DB(未删缓存),线程 B 读取缓存(旧数据)→ 不一致

方案二:先删缓存,再更新 DB

1. 删除缓存
2. 更新数据库

问题

  • 步骤 1 成功,步骤 2 失败 → 缓存为空,DB 未更新 → 下次读取缓存空,加载 DB 新数据(可接受)

  • 并发场景:线程 A 删缓存,线程 B 读缓存(空)→ 加载 DB 旧数据,线程 A 更新 DB → 缓存是旧数据

方案三:Cache Aside(读时写缓存)

读缓存:
1. 读缓存,命中 → 返回
2. 读缓存,未命中 → 读 DB → 写缓存 → 返回

写缓存:
1. 更新 DB
2. 删除缓存

问题

  • 仍然存在并发不一致问题

  • 删除缓存失败 → 缓存永久不一致

  • Cache Aside 模式的局限性

问题 1:删除缓存失败

  • 解决方案:重试 + 本地消息表 + MQ

问题 2:并发读写不一致

  • 解决方案:分布式锁(性能差)、延迟双删

问题 3:主从延迟导致不一致

  • 场景:主库更新,从库未同步,缓存删除 → 读请求打到从库(旧数据)→ 写入缓存(旧数据)

  • 解决方案:Canal 订阅从库 binlog

  • Canal + MQ 方案

架构设计

MySQL 主库 → binlog → Canal Server → MQ(RocketMQ/Kafka) → 消费者 → 删除/更新缓存

流程说明

  1. Canal Server 监听 MySQL binlog(增量日志)

  2. 解析 binlog,提取数据变更(INSERT/UPDATE/DELETE)

  3. 发送变更消息到 MQ

  4. 消费者消费消息,删除或更新缓存

Canal 配置示例

# canal.properties
canal.serverMode = rocketMQ
canal.mq.servers = 127.0.0.1:9876
canal.mq.topic = canal-binlog-topic

# instance.properties
canal.instance.master.address = 127.0.0.1:3306
canal.instance.dbUsername = root
canal.instance.dbPassword = password
canal.instance.connectionCharset = UTF-8
canal.instance.filterRegex = .*\\..*  # 监听所有表

消费者代码示例

@RocketMQMessageListener(
    topic = "canal-binlog-topic",
    consumerGroup = "cache-sync-consumer-group"
)
public class CanalBinlogListener implements RocketMQListener<String> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Override
    public void onMessage(String message) {
        // 解析 Canal 消息
        CanalEntry.Entry entry = JSON.parseObject(message, CanalEntry.Entry.class);

        if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
            CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());

            for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                String tableName = entry.getHeader().getTableName();

                // 根据表名和操作类型处理缓存
                if ("t_order".equals(tableName)) {
                    handleOrderCache(rowData);
                } else if ("t_product".equals(tableName)) {
                    handleProductCache(rowData);
                }
            }
        }
    }

    private void handleOrderCache(CanalEntry.RowData rowData) {
        // 提取订单 ID
        String orderId = rowData.getAfterColumnsList().stream()
            .filter(col -> "id".equals(col.getName()))
            .findFirst()
            .get()
            .getValue();

        // 删除缓存
        String cacheKey = "order:" + orderId;
        redisTemplate.delete(cacheKey);

        log.info("Delete cache: {}", cacheKey);
    }
}
  1. 延迟双删的消息化替代

传统延迟双删

1. 删除缓存
2. 更新数据库
3. 休眠 500ms
4. 再次删除缓存

问题

  • 步骤 3 休眠阻塞线程,性能差

  • 步骤 4 删除缓存可能失败

MQ 延迟双删

1. 删除缓存
2. 更新数据库
3. 发送延迟消息到 MQ(延迟 500ms)
4. 消费者消费延迟消息,再次删除缓存

代码示例

@Service
public class OrderService {

    @Autowired
    private OrderRepository orderRepository;

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Transactional
    public void updateOrder(Order order) {
        // 1. 删除缓存
        String cacheKey = "order:" + order.getId();
        redisTemplate.delete(cacheKey);

        // 2. 更新数据库
        orderRepository.save(order);

        // 3. 发送延迟消息(延迟 500ms)
        rocketMQTemplate.syncSend(
            "cache-delay-topic",
            new CacheDeleteMessage(cacheKey),
            3000,  // 超时时间
            16     // 延迟等级(RocketMQ 延迟等级:1-18,对应 1s-2h)
            // 延迟等级 16 = 500ms(需配置)
        );
    }
}

@RocketMQMessageListener(
    topic = "cache-delay-topic",
    consumerGroup = "cache-delay-consumer-group"
)
public class CacheDelayListener implements RocketMQListener<CacheDeleteMessage> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Override
    public void onMessage(CacheDeleteMessage message) {
        // 再次删除缓存
        redisTemplate.delete(message.getCacheKey());
        log.info("Delay delete cache: {}", message.getCacheKey());
    }
}

RocketMQ 延迟等级配置

# broker.conf
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
  1. 方案对比

查看内嵌表格

  1. 实战:电商商品缓存一致性的完整方案

需求

  • 商品信息缓存到 Redis

  • 商品更新后,缓存同步更新

  • 支持高并发(10 万 QPS)

  • 保证最终一致性

方案选择:Canal + MQ + 延迟双删

架构设计

商品服务 → 更新 MySQL → Canal 监听 binlog → MQ → 消费者删除 Redis 缓存
    发送延迟消息 → MQ → 消费者再次删除 Redis 缓存

代码实现

商品服务

@Service
public class ProductService {

    @Autowired
    private ProductRepository productRepository;

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Transactional
    public void updateProduct(Product product) {
        // 1. 更新数据库
        productRepository.save(product);

        // 2. 发送延迟消息(延迟 500ms)
        rocketMQTemplate.syncSend(
            "product-cache-delay-topic",
            new CacheDeleteMessage("product:" + product.getId()),
            3000,
            16     // 延迟等级 16 = 500ms
        );
    }
}

Canal 消费者

@RocketMQMessageListener(
    topic = "canal-binlog-topic",
    consumerGroup = "product-cache-consumer-group"
)
public class ProductCanalListener implements RocketMQListener<String> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Override
    public void onMessage(String message) {
        CanalEntry.Entry entry = JSON.parseObject(message, CanalEntry.Entry.class);

        if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
            String tableName = entry.getHeader().getTableName();

            if ("t_product".equals(tableName)) {
                CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());

                for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                    // 提取商品 ID
                    String productId = rowData.getAfterColumnsList().stream()
                        .filter(col -> "id".equals(col.getName()))
                        .findFirst()
                        .get()
                        .getValue();

                    // 删除缓存
                    String cacheKey = "product:" + productId;
                    redisTemplate.delete(cacheKey);

                    log.info("Delete product cache: {}", cacheKey);
                }
            }
        }
    }
}

延迟消息消费者

@RocketMQMessageListener(
    topic = "product-cache-delay-topic",
    consumerGroup = "product-cache-delay-consumer-group"
)
public class ProductCacheDelayListener implements RocketMQListener<CacheDeleteMessage> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Override
    public void onMessage(CacheDeleteMessage message) {
        // 再次删除缓存
        redisTemplate.delete(message.getCacheKey());
        log.info("Delay delete product cache: {}", message.getCacheKey());
    }
}

优化点

  1. 批量删除:Canal 消费者批量删除缓存(减少 Redis IO)

  2. 消息去重:使用 Redis Set 记录已删除的缓存 key,避免重复删除

  3. 监控告警:监控 Canal 消费延迟、MQ 消费延迟、缓存命中率

  4. 降级策略:Canal 或 MQ 故障时,降级为定时任务刷新缓存

监控指标

  • Canal binlog 延迟(秒)

  • MQ 消费延迟(秒)

  • 缓存命中率(%)

  • 缓存更新成功率(%)

3️⃣ Key Differences

查看内嵌表格

6.24 消息队列多租户隔离与流量治理

问题 24:消息队列的多租户隔离与流量治理怎么做?Topic 命名规范、限流、权限控制?

难度:⭐⭐⭐(命名空间隔离、ACL 权限控制、生产/消费限流、配额管理)

1️⃣ Common Answer

嗯...多租户隔离的话,主要就是靠命名空间和 Topic 命名规范来区分不同团队。权限控制的话,Kafka 有 ACL,RocketMQ 也有权限机制。限流的话,生产端可以控制发送速率,消费端可以控制拉取频率。还有就是配额管理,不能让某个团队占用太多资源。其实核心就是规范好命名,配好权限,再加点限流措施就行了。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. Topic 命名规范
格式:{环境}-{业务线}-{功能}-{版本}
示例:
- prod-order-create-v1
- prod-payment-callback-v2
- test-user-login-v1
- pre-inventory-sync-v1

规范要点:

  • 环境隔离:prod/test/pre/dev 明确区分

  • 业务线归属:order/payment/user/inventory 等

  • 功能语义:create/callback/login/sync 等

  • 版本管理:v1/v2 支持灰度升级

  • 命名空间隔离

RocketMQ Namespace 方案:

<!-- RocketMQ 5.x 支持 Namespace -->
<property name="namespace">team-order</property>
<property name="namesrvAddr">127.0.0.1:9876</property>

<!-- 不同团队使用不同 Namespace -->
team-order: prod-order-create-v1
team-payment: prod-payment-callback-v1

Kafka 多集群方案:

# 不同业务线使用独立集群
kafka-cluster-order: broker1:9092,broker2:9092
kafka-cluster-payment: broker3:9092,broker4:9092
  1. ACL 权限控制

Kafka SASL + ACL 配置:

# 开启 SASL 认证
authorizer.class.name=kafka.security.authorizer.AclAuthorizer
allow.everyone.if.no.acl.found=false
super.users=User:admin

# 创建用户
kafka-configs.sh --zookeeper localhost:2181 --alter --add-config 'SCRAM-SHA-256=[password=secret123]' --entity-type users --entity-name team-order

# 配置 ACL(只允许写特定 Topic)
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
  --add --allow-principal User:team-order \
  --operation Write --topic prod-order-*

# 配置 ACL(只允许读特定 Topic)
kafka-acls.sh --authorizer-properties zookeeper.connect=localhost:2181 \
  --add --allow-principal User:team-order-consumer \
  --operation Read --group team-order-group --topic prod-order-*

RocketMQ ACL 配置:

# broker.conf
aclEnable=true

# plain_acl.yml
accounts:
  - accessKey: team-order
    secretKey: secret123
    whiteRemoteAddress: 192.168.1.*
    admin: false
    defaultTopicPerm: DENY
    defaultGroupPerm: SUB
    topicPerms:
      - topic=prod-order-*;perm=PUB|SUB
  - accessKey: team-payment
    secretKey: secret456
    admin: false
    topicPerms:
      - topic=prod-payment-*;perm=PUB|SUB
  1. 生产/消费限流

Kafka Quota 配置:

# 生产端配额(每秒 10MB)
kafka-configs.sh --alter --add-config 'producer_byte_rate=10485760' \
  --entity-type users --entity-name team-order \
  --bootstrap-server broker1:9092

# 消费端配额(每秒 20MB)
kafka-configs.sh --alter --add-config 'consumer_byte_rate=20971520' \
  --entity-type users --entity-name team-order-consumer \
  --bootstrap-server broker1:9092

# 客户端级别配额
kafka-configs.sh --alter --add-config 'request_percentage=25' \
  --entity-type clients --entity-name client-1 \
  --bootstrap-server broker1:9092

RocketMQ 流控机制:

// 生产端流控
DefaultMQProducer producer = new DefaultMQProducer("producer_group");
producer.setRetryTimesWhenSendFailed(3);
producer.setSendMsgTimeout(3000);
// 通过 maxMessageSize 控制单条消息大小
producer.setMaxMessageSize(4 * 1024 * 1024);

// 消费端流控
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group");
consumer.setPullBatchSize(32); // 每次拉取 32 条
consumer.setConsumeMessageBatchMaxSize(10); // 每次消费 10 条
consumer.setConsumeThreadMin(5);
consumer.setConsumeThreadMax(10);
  1. 灰度发布场景

消息路由 + 消费者分组:

// 生产端:根据版本路由
public class VersionPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                        Object value, byte[] valueBytes, Cluster cluster) {
        String version = ((Message) value).getVersion();
        // v1 发往分区 0-2,v2 发往分区 3-5
        return version.equals("v2") ?
               (int)(Math.random() * 3) + 3 :
               (int)(Math.random() * 3);
    }
}

// 消费端:指定消费分区
consumer.assign(Arrays.asList(
    new TopicPartition("prod-order-create", 3),
    new TopicPartition("prod-order-create", 4),
    new TopicPartition("prod-order-create", 5)
));
  1. 实战:多团队共用集群的治理方案

治理架构:

┌─────────────────────────────────────────────┐
│         共享 MQ 集群                        │
├─────────────────────────────────────────────┤
│  Team Order  │  Team Payment  │  Team User  │
│  Namespace   │  Namespace     │  Namespace  │
│  ┌─────────┐ │  ┌──────────┐  │  ┌────────┐ │
│  │ ACL     │ │  │ ACL      │  │  │ ACL    │ │
│  │ Quota   │ │  │ Quota    │  │  │ Quota  │ │
│  │ Monitor │ │  │ Monitor  │  │  │ Monitor│ │
│  └─────────┘ │  └──────────┘  │  └────────┘ │
└─────────────────────────────────────────────┘

监控指标:

  • 各团队 Topic 数量、消息积压量

  • 各团队带宽使用率、配额命中率

  • 异常访问、越权操作告警

3️⃣ Key Differences

查看内嵌表格


6.25 消息队列数据迁移与集群扩容

问题 25:消息队列如何做到不停机平滑迁移和集群扩容?

难度:⭐⭐⭐⭐(Kafka 分区迁移、RocketMQ Topic 迁移、双写双读、数据一致性校验)

1️⃣ Common Answer

嗯...迁移的话,Kafka 可以用 reassign-partitions 工具,把分区从旧 Broker 迁到新 Broker。RocketMQ 的话,一般用双写双读,新旧集群都写,慢慢切读。扩容的话,就是加新机器,然后重新分配分区。关键是不能停机,要保证数据不丢。迁移的时候最好限速,避免影响性能。最后要做数据校验,确保两边一致。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. Kafka 分区迁移

使用 kafka-reassign-partitions 工具:

# 步骤 1:生成迁移计划(json 文件)
cat > move-partitions.json <<EOF
{
  "version": 1,
  "partitions": [
    {
      "topic": "order-events",
      "partition": 0,
      "replicas": [1, 2, 3],
      "log_dirs": ["any", "any", "any"]
    },
    {
      "topic": "order-events",
      "partition": 1,
      "replicas": [4, 5, 6],
      "log_dirs": ["any", "any", "any"]
    }
  ]
}
EOF

# 步骤 2:执行迁移(限速 50MB/s)
kafka-reassign-partitions.sh --bootstrap-server broker1:9092 \
  --reassignment-json-file move-partitions.json \
  --execute \
  --throttle 52428800

# 步骤 3:验证迁移进度
kafka-reassign-partitions.sh --bootstrap-server broker1:9092 \
  --reassignment-json-file move-partitions.json \
  --verify

# 步骤 4:移除限速
kafka-configs.sh --bootstrap-server broker1:9092 \
  --alter --delete-config leader.replication.throttled.rate \
  --entity-type brokers --entity-name 1,2,3,4,5,6
  1. Kafka 集群扩容

完整扩容流程:

# 步骤 1:新增 Broker(配置 broker.id=7)
server.properties:
  broker.id=7
  listeners=PLAINTEXT://broker7:9092
  log.dirs=/data/kafka-logs

# 步骤 2:启动新 Broker
bin/kafka-server-start.sh config/server.properties &

# 步骤 3:分区重分配(自动生成计划)
kafka-reassign-partitions.sh --bootstrap-server broker1:9092 \
  --topics-to-move-json-file topics-to-move.json \
  --broker-list "0,1,2,3,4,5,6,7" \
  --generate > expand-plan.json

# topics-to-move.json
{
  "version": 1,
  "topics": [
    {"topic": "order-events"},
    {"topic": "payment-events"}
  ]
}

# 步骤 4:执行重分配(限速)
kafka-reassign-partitions.sh --bootstrap-server broker1:9092 \
  --reassignment-json-file expand-plan.json \
  --execute \
  --throttle 104857600

# 步骤 5:验证负载均衡
kafka-topics.sh --bootstrap-server broker1:9092 \
  --describe --topic order-events
  1. RocketMQ Topic 迁移

双写双读方案完整流程:

// 步骤 1:生产端双写
public class DualWriteProducer {
    private DefaultMQProducer oldProducer;
    private DefaultMQProducer newProducer;

    public void send(Message msg) throws Exception {
        // 旧集群写入(同步)
        SendResult oldResult = oldProducer.send(msg);

        // 新集群写入(异步)
        newProducer.send(msg, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                // 记录双写成功
                log.info("Dual write success: msgId={}", sendResult.getMsgId());
            }

            @Override
            public void onException(Throwable e) {
                // 新集群写入失败,降级到旧集群
                log.error("New cluster write failed, fallback to old", e);
            }
        });
    }
}

// 步骤 2:消费端双读(新优先)
public class DualReadConsumer {
    private DefaultMQPushConsumer oldConsumer;
    private DefaultMQPushConsumer newConsumer;
    private Set<String> processedMsgIds = new ConcurrentHashMap<>();

    public void init() {
        // 新集群消费(主)
        newConsumer.subscribe("order-topic", "*");
        newConsumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            for (MessageExt msg : msgs) {
                if (!processedMsgIds.contains(msg.getMsgId())) {
                    processMessage(msg);
                    processedMsgIds.add(msg.getMsgId());
                }
            }
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        });

        // 旧集群消费(兜底)
        oldConsumer.subscribe("order-topic", "*");
        oldConsumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            for (MessageExt msg : msgs) {
                if (!processedMsgIds.contains(msg.getMsgId())) {
                    processMessage(msg);
                    processedMsgIds.add(msg.getMsgId());
                }
            }
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        });
    }
}

// 步骤 3:灰度切流(配置中心控制)
public class TrafficSwitcher {
    private double newClusterRatio; // 0.0 -> 1.0

    public boolean shouldUseNewCluster(String businessKey) {
        // 根据业务 key 哈希 + 灰度比例决定
        int hash = Math.abs(businessKey.hashCode());
        return (hash % 100) < (newClusterRatio * 100);
    }
}

// 步骤 4:完成迁移
// 1. 灰度比例 0% -> 100%
// 2. 观察新集群消费 lag,确认无积压
// 3. 停止旧集群生产端
// 4. 等待旧集群消费完成
// 5. 停止旧集群消费端
  1. 数据一致性校验

消息对账方案:

// 对账任务
public class MessageReconciliationJob {
    public void reconcile() {
        // 步骤 1:获取旧集群消息列表
        List<MessageExt> oldMessages = queryMessagesFromOldCluster(
            "order-topic",
            startTime,
            endTime
        );

        // 步骤 2:获取新集群消息列表
        List<MessageExt> newMessages = queryMessagesFromNewCluster(
            "order-topic",
            startTime,
            endTime
        );

        // 步骤 3:比对消息 ID
        Set<String> oldMsgIds = oldMessages.stream()
            .map(MessageExt::getMsgId)
            .collect(Collectors.toSet());

        Set<String> newMsgIds = newMessages.stream()
            .map(MessageExt::getMsgId)
            .collect(Collectors.toSet());

        // 步骤 4:找出差异
        Set<String> missingInNew = Sets.difference(oldMsgIds, newMsgIds);
        Set<String> missingInOld = Sets.difference(newMsgIds, oldMsgIds);

        // 步骤 5:处理差异消息
        if (!missingInNew.isEmpty()) {
            log.error("Missing messages in new cluster: {}", missingInNew);
            // 补发到新集群
            republishToNewCluster(missingInNew);
        }
    }
}

// offset 校验
public class OffsetValidator {
    public void validateOffsets() {
        // 获取旧集群消费 offset
        Map<TopicPartition, Long> oldOffsets = getConsumerOffsets(
            oldConsumerGroup,
            "order-topic"
        );

        // 获取新集群消费 offset
        Map<TopicPartition, Long> newOffsets = getConsumerOffsets(
            newConsumerGroup,
            "order-topic"
        );

        // 比对进度
        oldOffsets.forEach((tp, oldOffset) -> {
            Long newOffset = newOffsets.get(tp);
            if (newOffset != null && Math.abs(oldOffset - newOffset) > 1000) {
                log.warn("Offset mismatch: topic={}, partition={}, old={}, new={}",
                    tp.topic(), tp.partition(), oldOffset, newOffset);
            }
        });
    }
}
  1. 迁移过程中的风险控制

风险控制措施:

// 限速控制
public class RateLimiter {
    private final RateLimiter rateLimiter = RateLimiter.create(1000); // 1000 条/秒

    public void acquire() {
        rateLimiter.acquire();
    }
}

// 监控告警
public class MigrationMonitor {
    private void monitor() {
        // 监控指标
        metrics.gauge("migration.old_cluster.lag", this::getOldClusterLag);
        metrics.gauge("migration.new_cluster.lag", this::getNewClusterLag);
        metrics.gauge("migration.dual_write.success_rate", this::getDualWriteSuccessRate);
        metrics.gauge("migration.dual_read.duplicate_rate", this::getDuplicateRate);

        // 告警规则
        if (getOldClusterLag() > 100000) {
            alertManager.sendAlert("Old cluster lag too high");
        }

        if (getDualWriteSuccessRate() < 0.99) {
            alertManager.sendAlert("Dual write success rate low");
        }
    }
}

// 回滚方案
public class RollbackPlan {
    public void rollback() {
        // 1. 立即切换生产端回旧集群
        producerConfig.setUseNewCluster(false);

        // 2. 停止新集群消费
        newConsumer.shutdown();

        // 3. 确认旧集群消费正常
        waitForOldClusterCatchUp();

        // 4. 通知相关方
        notifyStakeholders("Rollback completed");
    }
}
  1. 实战:从旧集群迁移到新集群的完整 SOP

迁移 SOP:

【阶段一:准备阶段】(T-7 天)
1. 新集群部署完成,容量规划确认
2. 双写双读代码开发、测试
3. 对账工具开发、测试
4. 监控告警配置完成
5. 回滚方案评审通过

【阶段二:双写阶段】(T-3 天)
1. 生产端开启双写(新集群异步写入)
2. 观察新集群写入成功率(目标 >99.9%)
3. 监控新集群存储增长
4. 确认无性能影响

【阶段三:双读阶段】(T-1 天)
1. 消费端开启双读(新集群优先)
2. 灰度 10% 流量到新集群
3. 观察消费 lag、重复率
4. 逐步提升灰度比例:10% -> 30% -> 50% -> 100%

【阶段四:切流阶段】(T 日)
1. 100% 流量切换到新集群
2. 确认新集群消费正常
3. 停止旧集群生产端
4. 等待旧集群消费完成

【阶段五:清理阶段】(T+7 天)
1. 对账任务持续运行 7 天
2. 确认数据一致性
3. 停止旧集群消费端
4. 下线旧集群

3️⃣ Key Differences

查看内嵌表格


6.26 消息丢失排查

问题 26:生产环境消息丢了,怎么快速定位是哪个环节丢的?

难度:⭐⭐⭐(全链路消息对账、Producer 确认日志、Broker 存储校验、Consumer offset 比对)

1️⃣ Common Answer

嗯...消息丢失的话,先看生产端是不是发送成功了,看日志有没有发送成功的回调。然后看 Broker 端,消息是不是真的存进去了,有没有磁盘故障。最后看消费端,是不是消费了但是没提交 offset,或者 offset 提交时机不对。最好是做个全链路追踪,用 msgId 从生产到消费都记录下来,这样好排查。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. 消息丢失的三个环节
Producer → Broker → Consumer
   ↓         ↓          ↓
 发送确认   持久化     消费提交

三个环节都可能丢失:

  • Producer 端:发送失败、网络抖动、重试不足

  • Broker 端:未落盘、副本未同步、磁盘故障

  • Consumer 端:消费失败、offset 提交过早、Rebalance 导致重复/丢失

  • Producer 端排查

发送确认日志检查:

// Producer 配置(确保可靠发送)
Properties props = new Properties();
props.put("acks", "all"); // 等待所有副本确认
props.put("retries", 3); // 重试 3 次
props.put("max.in.flight.requests.per.connection", 1); // 顺序保证
props.put("enable.idempotence", true); // 幂等性

// 发送回调记录日志
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("order-topic", "key", "value"),
    new Callback() {
        @Override
        public void onCompletion(RecordMetadata metadata, Exception exception) {
            if (exception != null) {
                // 发送失败,记录错误日志
                log.error("Send failed: topic={}, partition={}, offset={}, error={}",
                    metadata.topic(), metadata.partition(),
                    metadata.offset(), exception.getMessage());
                // 触发告警
                alertManager.sendAlert("Producer send failed");
            } else {
                // 发送成功,记录成功日志(包含 msgId)
                log.info("Send success: topic={}, partition={}, offset={}, msgId={}",
                    metadata.topic(), metadata.partition(),
                    metadata.offset(), metadata.toString());
                // 记录到发送成功表
                recordSendSuccess(msgId, metadata.topic(), metadata.offset());
            }
        }
    });

// 排查 SQL
SELECT * FROM send_success_log
WHERE msg_id = 'xxx'
AND send_time BETWEEN '2024-01-01 00:00:00' AND '2024-01-01 01:00:00';

重试日志检查:

// 自定义重试策略
public class CustomRetryPolicy {
    private int maxRetries = 3;
    private long retryInterval = 1000; // 1 秒

    public void sendWithRetry(Message message) {
        int retryCount = 0;
        while (retryCount <= maxRetries) {
            try {
                SendResult result = producer.send(message);
                log.info("Send success: msgId={}, retryCount={}",
                    result.getMsgId(), retryCount);
                return;
            } catch (Exception e) {
                retryCount++;
                log.warn("Send failed, retrying: msgId={}, retryCount={}, error={}",
                    message.getMsgId(), retryCount, e.getMessage());

                if (retryCount > maxRetries) {
                    log.error("Send failed after max retries: msgId={}", message.getMsgId());
                    // 记录到失败表
                    recordSendFailure(message.getMsgId(), e.getMessage());
                    throw e;
                }

                Thread.sleep(retryInterval);
            }
        }
    }
}
  1. Broker 端排查

消息落盘检查:

# 检查 Broker 日志
tail -f /data/kafka/logs/server.log | grep "ERROR"

# 检查磁盘空间
df -h /data/kafka-logs

# 检查消息是否落盘
ls -lh /data/kafka-logs/order-topic-0/
# 应该看到 .log 文件在增长

# 检查消息内容
kafka-run-class.sh kafka.tools.DumpLogSegments \
  --files /data/kafka-logs/order-topic-0/00000000000000000000.log \
  --print-data-log | grep "msgId"

副本同步状态检查:

# 检查 Topic 分区状态
kafka-topics.sh --bootstrap-server broker1:9092 \
  --describe --topic order-topic

# 输出示例:
# Topic: order-topic Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
# 如果 Isr 少于 Replicas,说明有副本同步失败

# 检查副本延迟
kafka-replica-verifier.sh --broker-list broker1:9092,broker2:9092,broker3:9092 \
  --topic order-topic --report-interval-ms 10000

# 检查磁盘故障
dmesg | grep -i error
smartctl -a /dev/sda
  1. Consumer 端排查

offset 提交时机检查:

// 错误示例:先提交 offset 再消费(可能导致丢失)
consumer.subscribe(Arrays.asList("order-topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    consumer.commitSync(); // ❌ 先提交 offset
    for (ConsumerRecord<String, String> record : records) {
        processMessage(record); // 如果这里抛异常,消息就丢了
    }
}

// 正确示例:先消费再提交 offset
consumer.subscribe(Arrays.asList("order-topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        try {
            processMessage(record); // ✅ 先消费
        } catch (Exception e) {
            log.error("Consume failed: offset={}", record.offset(), e);
            // 消费失败,不提交 offset,下次重试
            continue;
        }
    }
    consumer.commitSync(); // ✅ 消费成功后再提交
}

消费日志检查:

// 记录消费日志
public class MessageConsumer {
    public void onMessage(ConsumerRecord<String, String> record) {
        String msgId = record.key();
        long offset = record.offset();

        log.info("Consume start: msgId={}, offset={}", msgId, offset);

        try {
            // 处理消息
            processMessage(record.value());

            // 记录消费成功
            log.info("Consume success: msgId={}, offset={}", msgId, offset);
            recordConsumeSuccess(msgId, offset);

        } catch (Exception e) {
            // 记录消费失败
            log.error("Consume failed: msgId={}, offset={}", msgId, offset, e);
            recordConsumeFailure(msgId, offset, e.getMessage());
        }
    }
}

// 排查 SQL
SELECT * FROM consume_success_log
WHERE msg_id = 'xxx'
ORDER BY consume_time DESC;

Rebalance 导致的重复/丢失:

// 监听 Rebalance 事件
consumer.subscribe(Arrays.asList("order-topic"), new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        // 分区被撤销前,提交当前 offset
        log.info("Partitions revoked: {}", partitions);
        consumer.commitSync();
    }

    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // 分区被分配后,记录日志
        log.info("Partitions assigned: {}", partitions);
        // 可以在这里记录当前 offset,用于排查
        partitions.forEach(tp -> {
            long offset = consumer.position(tp);
            log.info("Current offset: topic={}, partition={}, offset={}",
                tp.topic(), tp.partition(), offset);
        });
    }
});
  1. 全链路消息对账方案

msgId 追踪:

// 生成全局唯一 msgId
public class MsgIdGenerator {
    private static final AtomicLong counter = new AtomicLong(0);

    public static String generate() {
        long timestamp = System.currentTimeMillis();
        long machineId = getMachineId(); // 机器 ID
        long sequence = counter.getAndIncrement();
        return String.format("%d-%d-%d", timestamp, machineId, sequence);
    }
}

// 生产端:设置 msgId
Message message = new Message("order-topic", "key", "value");
message.setMsgId(MsgIdGenerator.generate());

// 消费端:记录 msgId
public void onMessage(ConsumerRecord<String, String> record) {
    String msgId = record.key();
    log.info("Consume: msgId={}", msgId);
    // ... 处理消息
}

生产-消费对账表:

-- 发送记录表
CREATE TABLE send_record (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    msg_id VARCHAR(64) NOT NULL,
    topic VARCHAR(64) NOT NULL,
    partition INT NOT NULL,
    offset BIGINT NOT NULL,
    send_time DATETIME NOT NULL,
    send_status VARCHAR(16) NOT NULL, -- SUCCESS/FAILED
    INDEX idx_msg_id (msg_id),
    INDEX idx_send_time (send_time)
);

-- 消费记录表
CREATE TABLE consume_record (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    msg_id VARCHAR(64) NOT NULL,
    topic VARCHAR(64) NOT NULL,
    partition INT NOT NULL,
    offset BIGINT NOT NULL,
    consume_time DATETIME NOT NULL,
    consume_status VARCHAR(16) NOT NULL, -- SUCCESS/FAILED
    INDEX idx_msg_id (msg_id),
    INDEX idx_consume_time (consume_time)
);

-- 对账 SQL:找出已发送但未消费的消息
SELECT s.msg_id, s.topic, s.partition, s.offset, s.send_time
FROM send_record s
LEFT JOIN consume_record c ON s.msg_id = c.msg_id
WHERE s.send_time >= '2024-01-01 00:00:00'
  AND s.send_time < '2024-01-01 01:00:00'
  AND c.msg_id IS NULL;
  1. 实战:一次线上消息丢失的完整排查 SOP

排查 SOP:

【阶段一:发现问题】(T+0 分钟)
1. 业务方反馈:订单数据缺失
2. 确认时间范围:2024-01-01 10:00:00 - 10:30:00
3. 涉及 Topic:order-topic

【阶段二:生产端排查】(T+5 分钟)
1. 查询发送成功日志
   SELECT * FROM send_record
   WHERE send_time BETWEEN '2024-01-01 10:00:00' AND '2024-01-01 10:30:00'
   AND topic = 'order-topic';

2. 结果:1000 条消息发送成功
3. 结论:生产端无问题

【阶段三:Broker 端排查】(T+10 分钟)
1. 检查 Broker 日志
   grep -i "error" /data/kafka/logs/server.log | grep "2024-01-01 10:0"

2. 检查磁盘空间
   df -h /data/kafka-logs

3. 检查副本同步状态
   kafka-topics.sh --describe --topic order-topic

4. 结果:Broker 正常,副本同步正常
5. 结论:Broker 端无问题

【阶段四:Consumer 端排查】(T+15 分钟)
1. 查询消费成功日志
   SELECT * FROM consume_record
   WHERE consume_time BETWEEN '2024-01-01 10:00:00' AND '2024-01-01 10:30:00'
   AND topic = 'order-topic';

2. 结果:只有 800 条消费成功
3. 缺失 200 条消息

4. 检查消费失败日志
   SELECT * FROM consume_record
   WHERE consume_time BETWEEN '2024-01-01 10:00:00' AND '2024-01-01 10:30:00'
   AND consume_status = 'FAILED';

5. 结果:200 条消费失败,错误原因:数据库连接超时
6. 结论:Consumer 端消费失败,导致消息丢失

【阶段五:定位根因】(T+20 分钟)
1. 检查 Consumer 代码
   发现:offset 提交时机错误(先提交再消费)

2. 检查数据库日志
   发现:10:15:00 数据库连接池耗尽

3. 根因:数据库连接池耗尽 → 消费失败 → offset 已提交 → 消息丢失

【阶段六:修复措施】(T+30 分钟)
1. 修复 Consumer 代码:改为先消费再提交 offset
2. 增加数据库连接池大小
3. 增加消费失败告警
4. 重放丢失的 200 条消息(从发送记录表获取)

【阶段七:验证修复】(T+60 分钟)
1. 观察消费日志,确认无消费失败
2. 对账检查,确认无消息丢失
3. 业务方确认订单数据完整

3️⃣ Key Differences

查看内嵌表格


6.27 消费倾斜与数据热点

问题 27:消息队列的消费倾斜(数据热点)问题怎么解决?某个分区消费特别慢怎么办?

难度:⭐⭐⭐(热点 key、分区不均、自定义分区策略、动态扩分区)

1️⃣ Common Answer

嗯...消费倾斜的话,一般是某个分区数据太多,或者某个消费者处理太慢。先看看是不是热点 key 导致的,比如某个商品特别热,所有消息都到一个分区了。解决的话,可以自定义分区策略,把热点 key 打散到不同分区。或者动态增加分区数,再加几个消费者。还可以在消费者内部用本地队列再做一次负载均衡。

2️⃣ Impressive Answer

我会从这几个角度思考:

  1. 消费倾斜的表现

监控各分区 lag:

// 监控脚本
public class PartitionLagMonitor {
    public void monitorLag() {
        AdminClient adminClient = AdminClient.create(adminConfig);

        // 获取 Topic 分区列表
        DescribeTopicsResult describeTopics = adminClient.describeTopics(
            Collections.singletonList("order-topic")
        );

        TopicDescription topicDescription = describeTopics.allTopics().get().get("order-topic");
        List<TopicPartitionInfo> partitions = topicDescription.partitions();

        // 获取每个分区的 lag
        for (TopicPartitionInfo partition : partitions) {
            int partitionId = partition.partition();
            TopicPartition tp = new TopicPartition("order-topic", partitionId);

            // 获取生产者 offset(最新 offset)
            Map<TopicPartition, Long> endOffsets = adminClient.endOffsets(
                Collections.singletonList(tp)
            );
            long endOffset = endOffsets.get(tp);

            // 获取消费者 offset
            Map<TopicPartition, OffsetAndMetadata> consumerOffsets = adminClient.listConsumerGroupOffsets(
                "order-consumer-group"
            ).partitionsToOffsetAndMetadata().get();
            long consumerOffset = consumerOffsets.get(tp).offset();

            // 计算 lag
            long lag = endOffset - consumerOffset;

            log.info("Partition lag: topic={}, partition={}, lag={}",
                "order-topic", partitionId, lag);

            // 告警:某个分区 lag 远高于其他分区
            if (lag > 100000) {
                alertManager.sendAlert(
                    String.format("High lag detected: partition=%d, lag=%d", partitionId, lag)
                );
            }
        }
    }
}

// 输出示例:
// Partition lag: topic=order-topic, partition=0, lag=100
// Partition lag: topic=order-topic, partition=1, lag=150
// Partition lag: topic=order-topic, partition=2, lag=50000  ← 倾斜!
// Partition lag: topic=order-topic, partition=3, lag=120
  1. 常见原因

原因一:热点 key 导致分区不均

// 示例:订单消息按商品 ID 分区
String productId = order.getProductId();
// 如果某个商品特别热,所有该商品的消息都到一个分区
int partition = Math.abs(productId.hashCode()) % 4;
// 假设商品 "iPhone15" 特别热,所有 iPhone15 订单都到分区 2

原因二:消费者处理能力不均

// 示例:消费者 1 处理逻辑复杂,消费者 2 处理逻辑简单
// Consumer 1
public void processMessage(Message message) {
    // 复杂的数据库操作
    saveToDatabase(message);
    // 复杂的缓存更新
    updateCache(message);
    // 复杂的外部接口调用
    callExternalApi(message);
}

// Consumer 2
public void processMessage(Message message) {
    // 简单的日志记录
    log.info("Message: {}", message);
}

原因三:分区数据量不均

# 检查各分区数据量
kafka-run-class.sh kafka.tools.GetOffsetShell \
  --broker-list broker1:9092 \
  --topic order-topic \
  --time -1

# 输出示例:
# order-topic:0:1000
# order-topic:1:1200
# order-topic:2:50000  ← 数据量远高于其他分区
# order-topic:3:1100
  1. 排查方法

检查 key 分布:

// 分析 key 分布
public class KeyDistributionAnalyzer {
    public void analyzeKeyDistribution() {
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerConfig);
        consumer.subscribe(Collections.singletonList("order-topic"));

        Map<String, Integer> keyCountMap = new HashMap<>();

        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
            for (ConsumerRecord<String, String> record : records) {
                String key = record.key();
                keyCountMap.put(key, keyCountMap.getOrDefault(key, 0) + 1);
            }

            // 定期输出 key 分布
            if (System.currentTimeMillis() % 60000 == 0) {
                log.info("Key distribution: {}", keyCountMap);

                // 找出热点 key
                keyCountMap.entrySet().stream()
                    .sorted(Map.Entry.<String, Integer>comparingByValue().reversed())
                    .limit(10)
                    .forEach(entry ->
                        log.info("Hot key: key={}, count={}", entry.getKey(), entry.getValue())
                    );
            }
        }
    }

    // 输出示例:
    // Hot key: key=iPhone15, count=50000
    // Hot key: key=MacBookPro, count=10000
    // Hot key: key=iPad, count=5000
}

分析消费者处理耗时:

// 监控消费者处理耗时
public class ConsumerPerformanceMonitor {
    public void monitorPerformance() {
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerConfig);
        consumer.subscribe(Collections.singletonList("order-topic"));

        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));

            for (ConsumerRecord<String, String> record : records) {
                long startTime = System.currentTimeMillis();

                try {
                    processMessage(record);
                } finally {
                    long endTime = System.currentTimeMillis();
                    long duration = endTime - startTime;

                    log.info("Process time: partition={}, offset={}, duration={}ms",
                        record.partition(), record.offset(), duration);

                    // 记录到监控系统
                    metrics.histogram("consumer.process.time", duration);
                }
            }
        }
    }
}
  1. 解决方案一:自定义分区策略,打散热点 key

Java 代码示例:

// 自定义分区器
public class HotKeyPartitioner implements Partitioner {
    private Random random = new Random();

    // 热点 key 列表(可配置)
    private Set<String> hotKeys = new HashSet<>(Arrays.asList(
        "iPhone15", "MacBookPro", "iPad"
    ));

    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                        Object value, byte[] valueBytes, Cluster cluster) {
        String keyStr = (String) key;

        // 如果是热点 key,随机分配到不同分区
        if (hotKeys.contains(keyStr)) {
            int partitionCount = cluster.partitionCountForTopic(topic);
            return random.nextInt(partitionCount);
        }

        // 非热点 key,使用默认的 hash 分区
        int partitionCount = cluster.partitionCountForTopic(topic);
        return Math.abs(keyStr.hashCode()) % partitionCount;
    }

    @Override
    public void close() {}

    @Override
    public void configure(Map<String, ?> configs) {
        // 可以从配置中读取热点 key 列表
        String hotKeysConfig = (String) configs.get("hot.keys");
        if (hotKeysConfig != null) {
            hotKeys = new HashSet<>(Arrays.asList(hotKeysConfig.split(",")));
        }
    }
}

// 配置分区器
Properties props = new Properties();
props.put("partitioner.class", "com.example.HotKeyPartitioner");
props.put("hot.keys", "iPhone15,MacBookPro,iPad");

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
  1. 解决方案二:动态扩分区 + 消费者扩容

动态扩分区:

# 步骤 1:增加分区数
kafka-topics.sh --bootstrap-server broker1:9092 \
  --alter --topic order-topic \
  --partitions 8

# 步骤 2:验证分区数
kafka-topics.sh --bootstrap-server broker1:9092 \
  --describe --topic order-topic

# 输出示例:
# Topic: order-topic PartitionCount: 8 ReplicationFactor: 3

消费者扩容:

// 原来有 4 个消费者,现在扩容到 8 个
// Consumer 1
KafkaConsumer<String, String> consumer1 = new KafkaConsumer<>(props);
consumer1.subscribe(Collections.singletonList("order-topic"));

// Consumer 2
KafkaConsumer<String, String> consumer2 = new KafkaConsumer<>(props);
consumer2.subscribe(Collections.singletonList("order-topic"));

// ... Consumer 3-8

// Kafka 会自动重新分配分区
// 原来:4 个消费者,每个消费 1 个分区
// 现在:8 个消费者,每个消费 1 个分区(8 个分区)
  1. 解决方案三:本地队列二次分发

消费者内部再做负载均衡:

// 主消费者
public class MainConsumer {
    private KafkaConsumer<String, String> kafkaConsumer;
    private ExecutorService executorService;
    private BlockingQueue<Message> localQueue;

    public void start() {
        kafkaConsumer = new KafkaConsumer<>(consumerConfig);
        kafkaConsumer.subscribe(Collections.singletonList("order-topic"));

        // 本地队列(内存队列)
        localQueue = new LinkedBlockingQueue<>(10000);

        // 线程池(本地消费者)
        executorService = Executors.newFixedThreadPool(10);

        // 启动本地消费者
        for (int i = 0; i < 10; i++) {
            executorService.submit(new LocalConsumer(localQueue));
        }

        // 主消费者:拉取消息,放入本地队列
        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(100));
            for (ConsumerRecord<String, String> record : records) {
                Message message = new Message(record.key(), record.value());
                localQueue.offer(message);
            }
        }
    }
}

// 本地消费者
public class LocalConsumer implements Runnable {
    private BlockingQueue<Message> queue;

    public LocalConsumer(BlockingQueue<Message> queue) {
        this.queue = queue;
    }

    @Override
    public void run() {
        while (true) {
            try {
                Message message = queue.take();
                processMessage(message);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
    }
}

// 架构图:
// Kafka Topic (4 分区)
//   ↓
// Main Consumer (1 个)
//   ↓
// Local Queue (内存队列)
//   ↓
// Local Consumers (10 个线程)
  1. 实战:大促期间热点商品导致消费倾斜的处理经验

实战案例:

【背景】
双 11 大促,订单消息量激增 10 倍
Topic: order-topic, 4 个分区, 4 个消费者
发现:分区 2 的 lag 远高于其他分区(50000 vs 100)

【排查】
1. 检查 key 分布
   发现:商品 "iPhone15" 的订单占 80%,全部在分区 2

2. 检查消费者处理耗时
   发现:所有消费者处理耗时相近(平均 50ms)

3. 结论:热点商品导致分区倾斜

【解决方案】
方案演进:

阶段一:自定义分区器(快速缓解)
- 实现 HotKeyPartitioner,将热点商品随机分配到 4 个分区
- 效果:lag 从 50000 降到 20000,缓解但不彻底

阶段二:动态扩分区(根本解决)
- 将分区数从 4 扩到 16
- 消费者数从 4 扩到 16
- 效果:lag 降到 1000,基本解决

阶段三:本地队列二次分发(最终优化)
- 每个消费者内部用 10 个线程消费
- 总处理能力:16 消费者 × 10 线程 = 160 并发
- 效果:lag 降到 100,完全解决

【总结】
消费倾斜的解决思路:
1. 先用自定义分区器快速缓解(5 分钟)
2. 再用扩分区 + 扩消费者根本解决(30 分钟)
3. 最后用本地队列二次分发优化(1 小时)

3️⃣ Key Differences

查看内嵌表格

6.28 关联问题速查表

查看内嵌表格