消息队列与分布式中间件¶

消息顺序性与分区策略¶
📬 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 的 SeekToCurrentErrorHandler 和 DeadLetterPublishingRecoverer 来实现。
@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;
}
消息消费失败的完整处理流程:
-
消费者处理消息,若抛出异常,Spring-Kafka 会捕捉到。
-
SeekToCurrentErrorHandler执行重试,每次重试前等待一定时间(指数退避)。 -
若重试耗尽,
DeadLetterPublishingRecoverer将消息发送到 DLQ Topic(如order-events.DLQ)。 -
运维人员监控 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¶
我会从这几个角度思考:
- Rebalance 的触发条件
- Rebalance 的过程(影响为什么大)
Stop the World:重平衡期间所有消费者停止消费
1. Coordinator 发送 GroupHeartbeat
2. 选出一个 Leader Consumer
3. Leader 计算分区分配方案
4. 同步给所有消费者
5. 各消费者开始消费
- 频繁重平衡的原因
- 关键参数调优
session.timeout.ms=45000 # 超时时间
heartbeat.interval.ms=15000 # 心跳间隔
max.poll.interval.ms=300000 # 最大处理间隔
max.poll.records=500 # 每次 poll 数量
- 静态成员机制(Kafka 2.3+)
- 实战经验 有一次容器滚动更新,每次重启都触发 Rebalance,导致消费中断。后来用了静态成员机制,配合合理的超时设置,基本避免了这个问题。
3️⃣ Key Differences¶
6.5 消息追踪与可观测性¶
问题 5:如何构建消息系统的可观测性?消息追踪怎么做?¶
难度:⭐⭐⭐⭐(全链路追踪、监控体系、问题诊断)
1️⃣ Common Answer¶
可观测性的话,可以打日志,然后用 ELK 查。也可以加一些监控指标,用 Grafana 展示。消息追踪可以用一些开源工具,或者自己实现。
2️⃣ Impressive Answer¶
我会从这几个角度思考:
- 可观测性的三支柱
- 消息追踪的核心设计
- 关键监控指标
- 问题诊断的完整链路
- 生产环境的实践
- 统一日志格式:JSON + TraceID + 业务关键字段
- Prometheus + Grafana:实时监控
- SkyWalking:分布式追踪
- 自定义 Dashboard:消息健康度大盘
-
高阶能力
-
消息血缘:知道每条消息从哪来到哪去
-
延迟分析:P95/P99 延迟分位数
-
容量规划:根据历史数据预测峰值
3️⃣ Key Differences¶
6.7 消息队列选型对比¶
问题 6:Kafka、RocketMQ、RabbitMQ 怎么选?各自适用什么场景?¶
难度:⭐⭐(消息队列选型、吞吐量、延迟、可靠性、生态对比)
1️⃣ Common Answer¶
Kafka 吞吐量最高,适合大数据场景;RabbitMQ 延迟最低,适合实时性要求高的场景;RocketMQ 是阿里开源的,功能比较全面,适合电商场景。我们公司用 Kafka 做日志收集,用 RocketMQ 做订单处理。
2️⃣ Impressive Answer¶
我会从这几个角度思考:选型核心维度、三款 MQ 的详细对比、场景化选型建议、实战决策经验。
- 选型的核心维度
消息队列选型主要考虑 5 个维度:
-
吞吐量:单位时间能处理的消息数量
-
延迟:消息从发送到接收的时间间隔
-
可靠性:消息不丢失、不重复的保障
-
生态成熟度:社区活跃度、文档完善度、第三方工具支持
-
运维成本:部署、监控、故障排查的复杂度
-
三款 MQ 的详细对比
-
场景化选型建议
-
日志采集、用户行为分析 → Kafka
- 原因:吞吐量极高,天然适配流式处理(Flink/Spark)
-
典型场景:埋点日志、监控指标采集、实时数仓
-
订单交易、支付场景 → RocketMQ
- 原因:支持事务消息,保障最终一致性;可靠性高,消息不丢失
-
典型场景:订单创建、支付回调、库存扣减
-
复杂路由、实时通信 → RabbitMQ
- 原因:Exchange-Queue 模型灵活,延迟极低
-
典型场景:即时通讯、复杂规则路由、微服务解耦
-
实战决策经验
去年我们做电商订单系统时,选型过程是这样的:
-
需求分析:订单创建后需要同步到 5 个下游系统(库存、物流、积分、风控、数据仓库),要求消息不丢失、支持事务消息
-
技术选型:排除了 Kafka(事务消息支持弱)和 RabbitMQ(吞吐量不够),最终选择 RocketMQ
-
落地效果:日均 500 万订单,消息延迟控制在 100ms 以内,通过事务消息保障了订单与库存的一致性
关键结论:没有最好的 MQ,只有最适合的。一定要结合业务场景、团队能力、现有技术栈综合决策。
3️⃣ Key Differences¶
6.8 消息可靠性保证¶
问题 8:如何保证消息不丢失?从生产端、Broker、消费端三个环节分析¶
难度:⭐⭐⭐(消息可靠性、ACK机制、持久化、事务消息)
1️⃣ Common Answer¶
开启持久化、手动 ACK 就行了。生产端设置重试,Broker 开启刷盘,消费端确认后再提交 offset。
2️⃣ Impressive Answer¶
我会从这几个角度思考:全链路风险分析、生产端保障机制、Broker 可靠性配置、消费端幂等设计、Kafka/RocketMQ 实战配置。
- 全链路风险分析
每个环节都可能丢消息,需要层层防护。
-
生产端保障
-
同步发送:
future.get()确保消息到达 Broker -
重试机制:配置
retries=3,指数退避避免雪崩 -
事务消息: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;
}
});
- Broker 保障
-
消费端保障
-
手动 ACK:消费成功后再提交 offset
-
幂等消费:业务层去重(后续章节详述)
-
消费位移管理:
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; // 确认消费
};
- 实战配置组合拳
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¶
我会从这几个角度思考:幂等性必要性、五种方案对比、适用场景分析、实战案例设计、幂等与重试关系。
-
为什么需要幂等
-
At Least Once 语义:消息可能重复投递
-
网络重试:生产端重试、消费端重试
-
Rebalance 重复消费:消费者组重平衡时重复投递
-
业务异常:ACK 超时导致重复消费
-
五种幂等方案对比
- 方案详解
方案一:数据库唯一键
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;
}
- 实战:订单支付场景幂等设计
组合方案:唯一键 + 状态机
@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);
}
}
-
幂等与重试的关系
-
幂等是重试的前提:只有幂等才能安全重试
-
重试不等于幂等:重试是策略,幂等是保障
-
最佳实践:幂等设计 + 指数退避重试
// 重试 + 幂等
@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 替代方案、实战落地案例。
- 延迟消息的业务场景
延迟消息在业务中非常常见:
-
订单超时关闭:下单后 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¶
我会从这几个角度思考:事务消息的完整流程、半消息的存储原理、事务回查机制、与本地消息表方案的对比、实战应用场景。
- 事务消息的完整流程
Producer RocketMQ Broker Consumer
| | |
|--- 1. 发送半消息 ----------->| |
| |-- 写入 HALF_TOPIC |
|<-- 返回发送成功 -------------| |
| | |
|--- 2. 执行本地事务 | |
| (如:创建订单) | |
| | |
|--- 3a. 事务成功 → Commit --->| |
| |-- 投递到真实 Topic ----->|
| | |-- 4. 正常消费
|--- 3b. 事务失败 → Rollback ->| |
| |-- 删除半消息 |
| | |
| (Producer 崩溃/超时) | |
| |-- 4. 事务回查 --------->|
|<-- 查询本地事务状态 ---------| |
|--- 返回 Commit/Rollback ---->| |
- 半消息的存储原理
半消息并不会直接写入用户指定的 Topic,而是先写入一个系统内部 Topic:RMQ_SYS_TRANS_HALF_TOPIC。这个 Topic 对消费者不可见,确保了在本地事务执行完成前,消息不会被消费。当事务提交后,RocketMQ 会将消息从 RMQ_SYS_TRANS_HALF_TOPIC 移动到真实的 Topic 中,消费者才能消费。
- 事务回查机制
RocketMQ 通过事务回查机制解决生产者崩溃或超时的问题:
-
回查条件:当生产者发送半消息后,超过一定时间(默认 60 秒)未收到 Commit/Rollback 响应时触发
-
回查次数:默认最多回查 15 次
-
回查间隔:每次间隔 60 秒
-
回查逻辑:Broker 调用生产者实现的
TransactionListener.checkLocalTransaction()方法,查询本地事务状态 -
状态返回:
COMMIT_MESSAGE:提交事务ROLLBACK_MESSAGE:回滚事务-
UNKNOWN:继续回查 -
与本地消息表方案的对比
- 实战:订单创建 + 积分发放的事务消息方案
// 订单服务生产者
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 高可用全景、两者高可用方案对比、脑裂问题及解决方案、实战故障恢复经验。
- 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 查询路由信息
-
两者高可用方案对比
- 脑裂问题及解决方案
脑裂场景:网络分区导致集群出现多个"主节点",导致数据不一致。
Kafka 解决方案:
-
ZooKeeper 模式:通过临时节点的 Watch 机制,同一时刻只有一个 Controller
-
KRaft 模式:Raft 协议的多数派投票机制,确保只有获得多数票的节点才能成为 Leader
RocketMQ 解决方案:
-
Dledger 模式:Raft 协议的多数派投票机制
-
配置建议:集群节点数建议为奇数(3、5、7),避免偶数节点导致投票僵局
通用防护策略:
-
部署奇数个节点
-
配置合理的超时时间
-
监控网络分区和节点状态
-
使用 fencing token(围栏令牌)机制
-
实战:集群故障恢复经验
场景 1:Kafka Leader 挂掉
-
Controller 检测到 Leader 挂掉
-
从 ISR 中选择 LEO 最高的副本作为新 Leader
-
更新 ZooKeeper 元数据
-
通知所有 Broker 和客户端
-
恢复时间:通常在 10-30 秒内
场景 2:RocketMQ Master 挂掉(Dledger 模式)
-
其他 Broker 检测到 Master 挂掉
-
触发 Raft 选举,选出新 Master
-
客户端自动切换到新 Master
-
恢复时间:通常在 5-15 秒内
优化建议:
-
监控 ISR/OSR 变化,及时处理副本滞后
-
配置合理的
replica.lag.time.max.ms和min.insync.replicas -
定期演练故障恢复流程
-
使用 Prometheus + Grafana 监控集群健康度
3️⃣ Key Differences¶
6.13 消息队列在微服务中的应用¶
问题 13:消息队列在微服务架构中承担什么角色?有哪些典型应用场景?¶
难度:⭐⭐(异步解耦、削峰填谷、事件驱动、CQRS、Saga)
1️⃣ Common Answer¶
消息队列在微服务中主要是解耦、异步、削峰。比如订单服务发消息,库存服务消费消息,这样就解耦了。还有流量削峰,防止服务被打挂。
2️⃣ Impressive Answer¶
我会从这几个角度思考:消息队列的四大核心角色、每个角色的业务场景与实现、事件驱动架构设计、实战编排案例。
- 消息队列在微服务中的四大核心角色
- 每个角色的业务场景与代码级实现
角色一:异步解耦 - 订单创建流程
传统同步调用的问题:
// 同步调用,链路过长,任何一环失败都会影响订单创建
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
- 事件驱动架构(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¶
我会从这几个角度思考:
- Kafka 写入路径:顺序追加写 Producer 发送消息 → 写入 Page Cache → 定期刷盘到磁盘(顺序追加写)
顺序写比随机写快 100 倍的原因:
-
磁盘磁头不需要频繁移动寻道(寻道时间 5-10ms,是主要瓶颈)
-
顺序写可以达到磁盘的物理极限速度(机械盘 100-200MB/s,SSD 更高)
-
批量写入:Kafka 默认 batch.size=16KB,积攒一批后一次性写入
-
零拷贝原理:传统 IO vs sendfile
传统 4 次拷贝流程:
sendfile 2 次拷贝流程:
关键改进:
-
减少 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¶
- 四种策略的分配逻辑和算法
假设场景: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 策略
CooperativeSticky 策略
- 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 个少)
- CooperativeSticky(增量式重平衡)
传统重平衡(Eager):
增量式重平衡(Cooperative):
优势:
-
减少消费停顿时间
-
降低重平衡对吞吐的影响
-
配置方式代码示例和实战选择建议
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¶
- 三种语义对比
- 幂等生产者原理
核心机制:PID(Producer ID)+ Sequence Number 去重
工作流程:
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
- 消费端配合: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 的消息都回滚
- 性能影响和实战取舍
性能开销:
-
事务开销约 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 我会从这几个角度思考:
-
CommitLog:统一顺序写入
-
所有 Topic 的消息都顺序写入同一个 CommitLog 文件
-
单文件固定 1GB,写满后滚动创建新文件
-
顺序写入,磁盘 IO 性能最优,写入吞吐高
-
ConsumeQueue:逻辑消费队列
-
每个 Topic 的每个 Queue 对应一个 ConsumeQueue
-
每条记录固定 20 字节:8 字节 CommitLog offset + 4 字节消息大小 + 8 字节 Tag hashcode
-
极小文件,可以全量加载到内存,快速定位消息
-
IndexFile:基于 Key 的消息检索
-
HashMap 结构,支持按消息 Key(如订单 ID)快速查找
-
索引文件包含 Key hashcode、CommitLog offset、时间戳等
-
用于业务查询、消息追踪场景
-
写入流程
- 读取流程
- 与 Kafka 存储模型对比
Trade-off 分析:
-
RocketMQ 牺牲读取性能换取极致写入性能和消息检索能力
-
Kafka 牺牲检索能力换取简单架构和良好读取性能
-
RocketMQ 适合订单、支付等需要按 Key 查询的场景
-
Kafka 适合日志采集、流计算等纯消费场景
-
实战经验
-
文件清理策略:默认 72 小时,磁盘空间达到 85% 时强制删除,可配置
deleteWhen和fileReservedTime -
磁盘水位告警:监控
CommitLog和ConsumeQueue目录使用率,超过 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 我会从这几个角度思考:
- 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"); // 支持或运算
- 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'"));
- 性能对比
性能测试数据:
-
Tag 过滤:单 Broker 可支持 10 万+ TPS
-
SQL92 过滤:单 Broker 约 2-3 万 TPS(受表达式引擎性能限制)
-
使用场景
Tag 过滤适用场景:
-
消息分类简单(如订单状态:待支付、已支付、已发货)
-
高吞吐场景(如日志采集、实时计算)
-
消息属性固定,无需动态过滤
SQL92 过滤适用场景:
-
复杂业务规则(如
price > 100 and region = 'HZ' and status in (1,2)) -
低吞吐、高精度过滤
-
需要动态调整过滤条件(不重启 Consumer)
-
配置方式
Broker 端配置(默认开启):
Java 代码示例:
// Tag 过滤
consumer.subscribe("TopicOrder", "TagPaid || TagShipped");
// SQL92 过滤
consumer.subscribe("TopicOrder",
MessageSelector.bySql("TAGS in ('TagPaid', 'TagShipped') and price > 100"));
- 实战:过滤导致的消费倾斜问题
问题现象:
-
某个 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 我会从这几个角度思考:
- 批量发送:batch.size + linger.ms 的配合逻辑
Kafka 配置:
配合逻辑:
-
linger.ms控制等待时间,batch.size控制批次大小 -
两者任一达到阈值即发送(先到先发)
-
高吞吐场景:
linger.ms=50-100ms,batch.size=1MB -
低延迟场景:
linger.ms=0-10ms,batch.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); // 批量发送
- 压缩算法对比
Kafka 配置:
RocketMQ 配置:
// 消息体超过阈值时压缩(默认 4KB,建议 1KB-10KB)
producer.setCompressMsgBodyOverHowmuch(1024);
// 压缩算法(默认 ZIP,建议 LZ4)
producer.setCompressLevel(5); // ZIP 压缩级别(0-9)
实战经验:
-
日志采集场景:使用 Zstd,压缩率高,节省带宽
-
实时计算场景:使用 LZ4,低延迟,高吞吐
-
订单交易场景:使用 Snappy,平衡性能与可靠性
-
acks 参数:吞吐与可靠性的 trade-off
Kafka 配置:
RocketMQ 配置:
// 同步发送(等待 Broker 确认)
SendResult result = producer.send(message); // 默认同步
// 异步发送(不等待确认,回调处理)
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
// 发送成功回调
}
@Override
public void onException(Throwable e) {
// 发送失败回调
}
});
- 异步发送 + 回调:提升吞吐的最佳实践
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);
}
});
- 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);
- 实战:从 5000 TPS 调优到 50000 TPS 的参数组合
初始配置(5000 TPS):
优化配置(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 我会从这几个角度思考:
- 消费线程模型对比
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 倍
-
异步落库:消费 → 内存队列 → 批量写入数据库
架构设计:
代码示例:
// 内存队列
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 模型对比
- 消费者并发度调优
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:增加消费者数量
Step 2:启用批量消费
Step 3:异步落库
Step 4:优化 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 我会从这几个角度思考:
- Saga 模式的两种实现
编排式(Orchestration):
-
由一个中心协调器(Saga Coordinator)控制整个流程
-
协调器调用各个服务的正向操作和补偿操作
-
优点:流程清晰、易于监控、补偿逻辑集中
-
缺点:协调器单点、耦合度高
协同式(Choreography):
-
没有中心协调器,各服务通过事件驱动协作
-
每个服务消费上游事件,执行本地事务,发布下游事件
-
优点:去中心化、松耦合、扩展性好
-
缺点:流程分散、难以追踪、补偿逻辑分散
-
基于消息队列的协同式 Saga
架构设计:
流程说明:
-
订单服务创建订单 → 发布
OrderCreated事件 -
库存服务消费
OrderCreated→ 扣减库存 → 发布InventoryReserved事件 -
支付服务消费
InventoryReserved→ 扣款 → 发布PaymentCompleted事件 -
物流服务消费
PaymentCompleted→ 创建物流单 → 发布ShipmentCreated事件
补偿流程(任一步骤失败):
-
支付失败 → 发布
PaymentFailed事件 -
库存服务消费
PaymentFailed→ 回滚库存 → 发布InventoryReleased事件 -
订单服务消费
InventoryReleased→ 取消订单 → 发布OrderCancelled事件 -
补偿事务设计
正向操作 + 逆向补偿:
// 正向操作:扣减库存
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); // 释放预留
}
- 与本地消息表的对比
本地消息表示例:
// 本地消息表实现
@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);
}
}
}
- 实战:订单 → 库存 → 支付的 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 我会从这几个角度思考:
- 缓存一致性的经典问题
方案一:先更新 DB,再删缓存
问题:
-
步骤 1 成功,步骤 2 失败 → 缓存是旧数据
-
并发场景:线程 A 更新 DB(未删缓存),线程 B 读取缓存(旧数据)→ 不一致
方案二:先删缓存,再更新 DB
问题:
-
步骤 1 成功,步骤 2 失败 → 缓存为空,DB 未更新 → 下次读取缓存空,加载 DB 新数据(可接受)
-
并发场景:线程 A 删缓存,线程 B 读缓存(空)→ 加载 DB 旧数据,线程 A 更新 DB → 缓存是旧数据
方案三:Cache Aside(读时写缓存)
问题:
-
仍然存在并发不一致问题
-
删除缓存失败 → 缓存永久不一致
-
Cache Aside 模式的局限性
问题 1:删除缓存失败
- 解决方案:重试 + 本地消息表 + MQ
问题 2:并发读写不一致
- 解决方案:分布式锁(性能差)、延迟双删
问题 3:主从延迟导致不一致
-
场景:主库更新,从库未同步,缓存删除 → 读请求打到从库(旧数据)→ 写入缓存(旧数据)
-
解决方案:Canal 订阅从库 binlog
-
Canal + MQ 方案
架构设计:
流程说明:
-
Canal Server 监听 MySQL binlog(增量日志)
-
解析 binlog,提取数据变更(INSERT/UPDATE/DELETE)
-
发送变更消息到 MQ
-
消费者消费消息,删除或更新缓存
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);
}
}
- 延迟双删的消息化替代
传统延迟双删:
问题:
-
步骤 3 休眠阻塞线程,性能差
-
步骤 4 删除缓存可能失败
MQ 延迟双删:
代码示例:
@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 延迟等级配置:
- 方案对比
- 实战:电商商品缓存一致性的完整方案
需求:
-
商品信息缓存到 Redis
-
商品更新后,缓存同步更新
-
支持高并发(10 万 QPS)
-
保证最终一致性
方案选择:Canal + MQ + 延迟双删
架构设计:
代码实现:
商品服务:
@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());
}
}
优化点:
-
批量删除:Canal 消费者批量删除缓存(减少 Redis IO)
-
消息去重:使用 Redis Set 记录已删除的缓存 key,避免重复删除
-
监控告警:监控 Canal 消费延迟、MQ 消费延迟、缓存命中率
-
降级策略:Canal 或 MQ 故障时,降级为定时任务刷新缓存
监控指标:
-
Canal binlog 延迟(秒)
-
MQ 消费延迟(秒)
-
缓存命中率(%)
-
缓存更新成功率(%)
3️⃣ Key Differences
6.24 消息队列多租户隔离与流量治理¶
问题 24:消息队列的多租户隔离与流量治理怎么做?Topic 命名规范、限流、权限控制?¶
难度:⭐⭐⭐(命名空间隔离、ACL 权限控制、生产/消费限流、配额管理)
1️⃣ Common Answer¶
嗯...多租户隔离的话,主要就是靠命名空间和 Topic 命名规范来区分不同团队。权限控制的话,Kafka 有 ACL,RocketMQ 也有权限机制。限流的话,生产端可以控制发送速率,消费端可以控制拉取频率。还有就是配额管理,不能让某个团队占用太多资源。其实核心就是规范好命名,配好权限,再加点限流措施就行了。
2️⃣ Impressive Answer¶
我会从这几个角度思考:
- 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
- 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
- 生产/消费限流
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);
- 灰度发布场景
消息路由 + 消费者分组:
// 生产端:根据版本路由
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)
));
- 实战:多团队共用集群的治理方案
治理架构:
┌─────────────────────────────────────────────┐
│ 共享 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¶
我会从这几个角度思考:
- 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
- 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
- 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. 停止旧集群消费端
- 数据一致性校验
消息对账方案:
// 对账任务
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);
}
});
}
}
- 迁移过程中的风险控制
风险控制措施:
// 限速控制
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");
}
}
- 实战:从旧集群迁移到新集群的完整 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¶
我会从这几个角度思考:
- 消息丢失的三个环节
三个环节都可能丢失:
-
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);
}
}
}
}
- 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
- 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);
});
}
});
- 全链路消息对账方案
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;
- 实战:一次线上消息丢失的完整排查 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¶
我会从这几个角度思考:
- 消费倾斜的表现
监控各分区 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
- 常见原因
原因一:热点 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
- 排查方法
检查 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);
}
}
}
}
}
- 解决方案一:自定义分区策略,打散热点 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:增加分区数
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 个分区)
- 解决方案三:本地队列二次分发
消费者内部再做负载均衡:
// 主消费者
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 个线程)
- 实战:大促期间热点商品导致消费倾斜的处理经验
实战案例:
【背景】
双 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 小时)