消息队列积压排查实战:从消费堵塞到背压治理
2026/7/3大约 15 分钟
消息队列积压排查实战:从消费堵塞到背压治理
生产环境消息积压是分布式系统中最常见的高频故障之一。本文从监控发现、根因定位、紧急处置到长期治理,完整还原一次消息队列积压的全链路排查过程,涵盖 RocketMQ 与 Kafka 两种主流 MQ 的排查差异。
一、问题背景
某电商平台大促期间,订单服务通过 RocketMQ 异步处理订单消息。大促开始后 10 分钟,监控系统连续告警:
[RocketMQ 消费延迟告警] Topic=OrderTopic, ConsumerGroup=OrderConsumerGroup
积压量: 1,250,000 条 (阈值: 100,000)
消费 TPS: 200 (正常: 2000)
生产 TPS: 5,000同时业务侧反馈用户下单后长时间未收到确认通知,订单状态长时间停留在"处理中"。
二、积压发现手段
2.1 核心监控指标
消息队列积压的发现依赖于完善的监控体系,以下是必须覆盖的监控维度:
┌─────────────────────────────────────────────────────────────────┐
│ 消息队列监控指标体系 │
├─────────────────┬───────────────────────────────────────────────┤
│ 指标类别 │ 具体指标 │
├─────────────────┼───────────────────────────────────────────────┤
│ 生产端 │ 发送 TPS、发送延迟、发送失败率、RT │
├─────────────────┼───────────────────────────────────────────────┤
│ Broker 端 │ 积压量、磁盘使用率、PageCache 命中率、网络带宽 │
├─────────────────┼───────────────────────────────────────────────┤
│ 消费端 │ 消费 TPS、消费延迟、消费失败率、消费者活跃数 │
├─────────────────┼───────────────────────────────────────────────┤
│ 业务层 │ 消息处理 RT、下游服务调用 RT、重试队列深度 │
└─────────────────┴───────────────────────────────────────────────┘2.2 告警阈值设计
# 告警规则配置示例(Prometheus + AlertManager)
groups:
- name: mq_alert
rules:
# 积压量告警 - 基于绝对值
- alert: MQBacklogAbsolute
expr: rocketmq_consumer_lag_msgs > 100000
for: 3m
labels:
severity: critical
annotations:
summary: "MQ积压超过10万条,当前: {{ $value }}"
# 积压量告警 - 基于消费时间
- alert: MQBacklagTime
expr: rocketmq_consumer_lag_time_seconds > 300
for: 2m
labels:
severity: critical
annotations:
summary: "MQ消费延迟超过5分钟"
# 消费 TPS 骤降
- alert: MQConsumeTPSDrop
expr: rate(rocketmq_consumer_tps[5m]) < 100
for: 5m
labels:
severity: warning
# 消费者实例数减少
- alert: MQConsumerInstanceDown
expr: rocketmq_consumer_online_count < 3
for: 1m
labels:
severity: critical2.3 积压量计算方式
RocketMQ 积压量计算:
积压量 = Broker 端最大 Offset - 消费者消费 Offset
# 查看积压
mqadmin consumerProgress -g OrderConsumerGroup -t OrderTopic
# 输出示例
Topic Broker MaxOffset ConsumerOffset Diff
OrderTopic broker-a 5000000 3750000 1250000
OrderTopic broker-b 5000000 4500000 500000Kafka 积压量计算:
# Kafka 消费延迟查看
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer-group
# 输出示例
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID
OrderTopic 0 3750000 5000000 1250000 consumer-1
OrderTopic 1 4500000 5000000 500000 consumer-2三、消费堵塞根因排查
3.1 根因分类决策树
消息积压根因定位
│
┌──────────┼──────────┐
▼ ▼ ▼
消费者少 消费者慢 生产者快
│ │ │
┌────────┤ ┌────┴────┐ │
▼ ▼ ▼ ▼ ▼
实例挂了 扩容 下游慢 GC停顿 评估是否
被驱逐 不足 /限流 /FullGC 正常流量
│
┌─────┼─────┐
▼ ▼ ▼
DB慢 RPC慢 外部
查询 调用 接口慢3.2 第一步:检查消费者实例状态
# RocketMQ 查看消费者状态
mqadmin consumerStatus -g OrderConsumerGroup -t OrderTopic
# 输出示例
# ConsumerClientInfo
consumerAddr clientId consumerType
10.0.1.101:52412 10.0.1.101@12345 CONSUME_PASSIVELY
10.0.1.102:52413 10.0.1.102@12346 CONSUME_PASSIVELY
# 预期 5 个实例,实际只有 2 个在线发现: 预期 5 个消费者实例,实际只有 2 个在线,3 个实例因 K8s Pod 调度问题未被拉起。
3.3 第二步:检查消费 TPS 与处理耗时
// 在消费者代码中添加耗时监控
@RocketMQMessageListener(topic = "OrderTopic", consumerGroup = "OrderConsumerGroup")
public class OrderMessageListener implements RocketMQListener<OrderMessage> {
@Override
public void onMessage(OrderMessage message) {
long start = System.currentTimeMillis();
try {
// 业务处理
orderService.processOrder(message);
} finally {
long cost = System.currentTimeMillis() - start;
// 记录处理耗时
Metrics.timer("order.consume.time", "status", "ok")
.record(cost, TimeUnit.MILLISECONDS);
// 超过 500ms 打印慢消费日志
if (cost > 500) {
log.warn("慢消费! msgId={}, cost={}ms", message.getMsgId(), cost);
}
}
}
}通过监控发现,单个消息处理耗时从正常的 50ms 飙升到 800ms:
# 慢消费日志
[WARN] 慢消费! msgId=AC11000000000000, cost=823ms
[WARN] 慢消费! msgId=AC11000000000001, cost=756ms
[WARN] 慢消费! msgId=AC11000000000002, cost=912ms
[WARN] 慢消费! msgId=AC11000000000003, cost=889ms3.4 第三步:定位慢消费根因
// 使用 Arthas 追踪方法耗时
[arthas@12345]$ trace com.example.service.OrderService processOrder '#cost > 200'
# 输出
`---[823ms] com.example.service.OrderService:processOrder()
+---[15ms] com.example.service.OrderService:validateOrder()
+---[680ms] com.example.service.InventoryService:deductInventory() ← 瓶颈!
| `---[670ms] com.example.client.InventoryClient:deduct()
| `---[665ms] java.net.SocketHttpClient:execute() ← HTTP 调用
+---[85ms] com.example.service.CouponService:useCoupon()
`---[43ms] com.example.service.NotificationService:notify()根因定位: 库存服务的 HTTP 调用耗时 680ms,是正常值(30ms)的 22 倍。进一步排查发现库存服务所在节点 CPU 使用率达到 95%,导致响应变慢。
3.5 消费堵塞常见根因汇总
| 根因类别 | 现象 | 排查手段 | 典型表现 |
|---|---|---|---|
| 消费者实例不足 | TPS 低、实例数少 | 检查 Pod/进程状态 | K8s 调度失败、OOM 被杀 |
| 下游服务慢 | 单消息耗时高 | Arthas trace | DB 慢查询、RPC 超时 |
| 下游限流 | 间歇性消费失败 | 检查限流日志 | Sentinel/RateLimiter 拒绝 |
| GC 停顿 | 消费者间歇性无响应 | GC 日志、jstat | Full GC > 3s |
| 消费串行化 | TPS 无法提升 | 检查并发配置 | consumeThreadMin=1 |
| 消息大小异常 | 间歇性慢消费 | 检查消息体大小 | 单条消息 > 1MB |
| 消费失败重试 | 重复消费同一条 | 检查死信队列 | DLQ 堆积 |
3.6 GC 停顿导致的消费堵塞
# 查看 GC 日志
# 启动参数添加: -Xlog:gc*:file=/var/log/gc.log:time,uptime,level,tags
# 典型的 Full GC 停顿
[2026-07-03T10:23:15.123+0800] GC(42) Pause Full (G1 Compaction Pause)
Young regions: 0->0
Old regions: 256->180
Humongous regions: 12->8
Metaspace: 78643K->78643K
2048M->1024M(2048M) 3520.135ms ← 3.5 秒停顿!
# 多次 Full GC 导致消费者长时间无心跳,被 Broker 判定为离线// GC 停顿导致的问题复现
// 消费者心跳线程与业务线程共享 JVM,Full GC 会导致心跳也停顿
// Broker 端超时未收到心跳 → 认为消费者离线 → 触发 Rebalance
// RocketMQ 消费者心跳超时配置
consumer.setMaxReconsumeTimes(5); // 最大重试次数
consumer.setConsumeTimeout(15L); // 消费超时时间(分钟)
consumer.setHeartbeatBrokerInterval(30); // 心跳间隔(秒),默认30s四、紧急扩容方案
4.1 扩容决策矩阵
┌────────────────────────────────────────────────────────────────┐
│ 积压紧急处置决策 │
├──────────────┬─────────────────────────────────────────────────┤
│ 场景 │ 处置方案 │
├──────────────┼─────────────────────────────────────────────────┤
│ 消费者少 │ 紧急拉起更多消费者实例 │
│ 且 Topic │ RocketMQ: 增加消费者实例(≤队列数) │
│ 队列数足够 │ Kafka: 增加消费者实例(≤分区数) │
├──────────────┼─────────────────────────────────────────────────┤
│ 队列数不足 │ 先扩队列/分区,再扩消费者 │
│ │ RocketMQ: 修改 readQueueNums │
│ │ Kafka: 增加 partition 数(注意不可减少) │
├──────────────┼─────────────────────────────────────────────────┤
│ 下游服务慢 │ 1. 扩容下游服务 │
│ │ 2. 降低消费并发,避免压垮下游 │
│ │ 3. 开启背压/限流 │
├──────────────┼─────────────────────────────────────────────────┤
│ 无法快速恢复 │ 1. 降级非核心消费逻辑 │
│ │ 2. 跳过非关键消息(记录后跳过) │
│ │ 3. 启动旁路消费快速清理 │
└──────────────┴─────────────────────────────────────────────────┘4.2 RocketMQ 紧急扩容
# 1. 查看当前 Topic 队列数
mqadmin topicList -t OrderTopic | grep readQueueNums
# readQueueNums=16
# 2. 消费者实例数不能超过队列数
# 当前 2 个消费者,16 个队列,可以扩到 16 个消费者实例
# 3. K8s 紧急扩容消费者
kubectl scale deployment order-consumer --replicas=10
# 4. 如果队列数不足,先扩队列(注意:只能扩不能缩)
mqadmin updateTopic -t OrderTopic -c ClusterA \
-r 32 -w 32 # readQueueNums=32, writeQueueNums=32
# 5. 验证扩容效果
mqadmin consumerProgress -g OrderConsumerGroup -t OrderTopic4.3 Kafka 紧急扩容
# 1. 查看当前分区数
kafka-topics.sh --bootstrap-server kafka:9092 \
--describe --topic OrderTopic
# PartitionCount: 16
# 2. Kafka 增加分区(只能增加不能减少!)
kafka-topics.sh --bootstrap-server kafka:9092 \
--alter --topic OrderTopic --partitions 32
# 3. 扩容消费者
kubectl scale deployment order-consumer --replicas=16
# 4. 验证
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer-group4.4 旁路消费方案
当下游服务暂时无法恢复,但消息不能一直积压时,可以使用旁路消费快速清理:
/**
* 旁路消费者:快速消费,只记录不处理
* 用于紧急清理积压,后续通过补偿任务处理
*/
public class BypassOrderConsumer {
private final FastStorage fastStorage; // Redis/本地文件
@RocketMQMessageListener(
topic = "OrderTopic",
consumerGroup = "OrderConsumerGroup_Bypass",
consumeMode = ConsumeMode.CONCURRENTLY,
consumeThreadMax = 64
)
public void onMessage(OrderMessage message) {
// 只做快速持久化,不做业务处理
fastStorage.lpush("pending:order:processing",
JSON.toJSONString(message));
}
}# 旁路消费后,后续通过补偿任务慢慢处理
# 监控 pending 队列大小
redis-cli LLEN pending:order:processing五、背压机制设计
5.1 背压的核心思想
┌─────────────┐ 5000 TPS ┌─────────────┐ 200 TPS ┌─────────────┐
│ Producer │ ──────────────▶ │ MQ Broker │ ──────────────▶ │ Consumer │
│ 生产者 │ │ 消息队列 │ │ 消费者 │
└─────────────┘ └─────────────┘ └─────────────┘
│
积压 125 万条
│
▼
┌─────────────────┐
│ 背压机制 │
│ Backpressure │
│ │
│ 1. 限流生产者 │
│ 2. 降级消费 │
│ 3. 动态调整并发 │
└─────────────────┘5.2 消费端背压实现
/**
* 自适应背压消费者
* 根据下游服务响应时间动态调整消费速率
*/
public class BackpressureConsumer {
private final Semaphore consumePermits;
private final AtomicInteger pendingCount = new AtomicInteger(0);
private final SlidingWindowMeter rtMeter;
private volatile int maxConcurrency = 32;
public BackpressureConsumer(int initialConcurrency) {
this.maxConcurrency = initialConcurrency;
this.consumePermits = new Semaphore(initialConcurrency);
this.rtMeter = new SlidingWindowMeter(100, 10); // 10个窗口,每个100ms
}
public void onMessage(OrderMessage message) {
// 获取许可,控制并发
if (!consumePermits.tryAcquire()) {
// 获取许可失败,延迟后重试
throw new ConsumeLaterException("背压: 等待消费许可");
}
pendingCount.incrementAndGet();
try {
long start = System.nanoTime();
// 设置下游调用超时
CompletableFuture<OrderResult> future = orderService
.processOrderAsync(message)
.orTimeout(500, TimeUnit.MILLISECONDS);
OrderResult result = future.join();
long costMs = (System.nanoTime() - start) / 1_000_000;
rtMeter.record(costMs);
// 动态调整并发度
adjustConcurrency();
} catch (TimeoutException e) {
// 下游超时,降低消费速率
decreaseConcurrency();
throw new ConsumeLaterException("下游超时,触发背压");
} finally {
pendingCount.decrementAndGet();
consumePermits.release();
}
}
/**
* 自适应并发调整算法
* 类似 TCP 拥塞控制的 AIMD(加性增,乘性减)
*/
private void adjustConcurrency() {
double p99RT = rtMeter.p99();
double avgRT = rtMeter.average();
if (p99RT > 800 || avgRT > 500) {
// P99 > 800ms 或 AVG > 500ms,乘性减少
decreaseConcurrency();
} else if (p99RT < 200 && avgRT < 100) {
// P99 < 200ms 且 AVG < 100ms,加性增加
increaseConcurrency();
}
}
private synchronized void decreaseConcurrency() {
int newConcurrency = Math.max(4, (int)(maxConcurrency * 0.5));
if (newConcurrency < maxConcurrency) {
int delta = maxConcurrency - newConcurrency;
consumePermits.acquireUninterruptibly(delta);
maxConcurrency = newConcurrency;
log.warn("背压降级: 并发度 {} -> {}", maxConcurrency + delta, newConcurrency);
}
}
private synchronized void increaseConcurrency() {
if (maxConcurrency < 64) {
maxConcurrency++;
consumePermits.release();
log.info("背压恢复: 并发度 -> {}", maxConcurrency);
}
}
}5.3 生产端背压实现
/**
* 生产端背压:根据积压量动态控制生产速率
*/
@Component
public class BackpressureProducer {
@Autowired
private DefaultMQProducer producer;
@Autowired
private MQMetricsCollector metricsCollector;
// 积压阈值
private static final long BACKLOG_WARN = 100_000;
private static final long BACKLOG_CRITICAL = 500_000;
public SendResult send(OrderMessage message) {
long backlog = metricsCollector.getBacklog("OrderTopic");
if (backlog > BACKLOG_CRITICAL) {
// 积压严重,拒绝生产(让上游重试或降级)
throw new BackpressureException(
"积压严重: " + backlog + ",触发生产端背压");
}
if (backlog > BACKLOG_WARN) {
// 积压告警,延迟生产
try {
Thread.sleep(Math.min(backlog / 1000, 100)); // 动态延迟
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
return producer.send(message);
}
}5.4 基于 Sentinel 的流控背压
// 使用 Sentinel 实现消费端流控
@SentinelResource(
value = "consumeOrder",
blockHandler = "consumeBlocked",
flowRules = {
@Rule(resource = "consumeOrder", grade = RuleConstant.FLOW_GRADE_QPS,
count = 500) // 限制每秒消费 500 条
}
)
public void onMessage(OrderMessage message) {
orderService.processOrder(message);
}
// 被限流时的处理
public void consumeBlocked(OrderMessage message, BlockException ex) {
log.warn("消费被限流: msgId={}", message.getMsgId());
// 抛出重试异常,消息稍后重新投递
throw new ConsumeLaterException("消费限流,稍后重试");
}六、RocketMQ 与 Kafka 积压排查差异
6.1 架构差异对比
┌──────────────────────────────────────────────────────────────────┐
│ RocketMQ vs Kafka 积压排查 │
├───────────────┬──────────────────────┬──────────────────────────┤
│ 维度 │ RocketMQ │ Kafka │
├───────────────┼──────────────────────┼──────────────────────────┤
│ 消费模型 │ Push + Pull 混合 │ Pull │
│ 积压查看 │ mqadmin consumerProg │ kafka-consumer-groups │
│ │ -ress │ --describe │
├───────────────┼──────────────────────┼──────────────────────────┤
│ 消费者扩容 │ ≤ queueNums │ ≤ partitionNums │
│ 队列扩容 │ 可随时增减 │ 只增不减 │
├───────────────┼──────────────────────┼──────────────────────────┤
│ 消费失败 │ 自动重试 + 死信队列 │ 需手动处理 │
│ 顺序消息 │ MessageQueueSelector │ 分区内有序 │
├───────────────┼──────────────────────┼──────────────────────────┤
│ 延迟消息 │ 原生支持(18个级别) │ 需要外部实现 │
│ 消费位点提交 │ 自动/定时提交 │ 自动/手动提交 │
├───────────────┼──────────────────────┼──────────────────────────┤
│ Rebalance │ 心跳超时触发 │ 协调器触发 │
│ │ 默认 120s 超时 │ session.timeout.ms │
│ │ │ 默认 10s/45s │
└───────────────┴──────────────────────┴──────────────────────────┘6.2 RocketMQ 特有排查
# 1. 查看消费者运行状态
mqadmin consumerStatus -g OrderConsumerGroup
# 2. 查看订阅关系一致性(常见坑:不同消费者订阅同一 Group 但不同 Topic)
mqadmin examineSubscriptionGroup -g OrderConsumerGroup
# 3. 查看死信队列
mqadmin queryMsgByKey -t %DLQ%OrderConsumerGroup -k "*"
# 4. 重置消费位点(慎用!)
# 将消费位点回溯到 1 小时前
mqadmin resetOffsetByTime -g OrderConsumerGroup -t OrderTopic \
-s "2026-07-03#10:00:00:000"
# 5. 暂停消费(不消费但保持心跳,不触发 Rebalance)
mqadmin consumeMessage -g OrderConsumerGroup -t OrderTopic -p false6.3 Kafka 特有排查
# 1. 查看消费者组详情
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer-group \
--members --verbose # 显示每个成员的分配情况
# 2. 查看消费者组状态
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer-group --state
# 3. 重置消费位点(需要先停止消费)
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--reset-offsets --group order-consumer-group \
--topic OrderTopic --to-datetime 2026-07-03T10:00:00.000 \
--execute
# 4. 查看 Broker 端消费者组协调器
kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
--describe --group order-consumer-group --verbose
# 5. Kafka 消费延迟监控(Burrow)
# Burrow 提供更精准的积压监控
burrow client -brokers kafka:9092 -groups order-consumer-group6.4 Kafka 消费者关键参数调优
# 消费者核心参数
max.poll.records=500 # 单次 poll 最大记录数
max.poll.interval.ms=300000 # 两次 poll 最大间隔(超时则 Rebalance)
session.timeout.ms=30000 # 心跳超时时间
heartbeat.interval.ms=10000 # 心跳间隔
fetch.min.bytes=1 # 最小拉取字节数
fetch.max.wait.ms=500 # 最大等待时间
# 关键调优点:
# 1. max.poll.records 不宜过大,否则单批处理时间可能超过 max.poll.interval.ms
# 2. 如果处理慢,降低 max.poll.records 而不是增大 max.poll.interval.ms
# 3. session.timeout.ms 不宜过大,否则消费者真正挂了需要很久才能发现七、长期治理方案
7.1 消费能力评估模型
┌────────────────────────────────────────────────────────────────┐
│ 消费能力容量规划模型 │
├────────────────────────────────────────────────────────────────┤
│ │
│ 所需消费者数 = 峰值生产 TPS / (单消费者 TPS × 安全系数) │
│ │
│ 示例: │
│ 峰值生产 TPS = 5000 │
│ 单消费者 TPS = 400(单消息 50ms,20 并发) │
│ 安全系数 = 0.7(留 30% 余量) │
│ │
│ 所需消费者数 = 5000 / (400 × 0.7) = 17.8 ≈ 18 个 │
│ │
│ 队列/分区数 = 32(≥ 消费者数,留扩容空间) │
│ │
└────────────────────────────────────────────────────────────────┘7.2 消费者健康度评分
/**
* 消费者健康度评分系统
* 综合评估消费能力,提前预警
*/
@Component
public class ConsumerHealthScorer {
private static final double WEIGHT_BACKLOG = 0.3;
private static final double WEIGHT_RT = 0.3;
private static final double WEIGHT_ERROR_RATE = 0.2;
private static final double WEIGHT_TPS_RATIO = 0.2;
public HealthScore calculate(String consumerGroup) {
ConsumerMetrics metrics = metricsCollector.collect(consumerGroup);
// 1. 积压得分(积压越少分越高)
double backlogScore = Math.max(0, 100 - metrics.getBacklog() / 1000.0);
// 2. 处理耗时得分(越快分越高)
double rtScore = Math.max(0, 100 - metrics.getP99RT() / 10.0);
// 3. 错误率得分(越低分越高)
double errorScore = Math.max(0, 100 - metrics.getErrorRate() * 100);
// 4. 消费速率得分(消费 TPS / 生产 TPS)
double ratio = metrics.getConsumeTPS() / Math.max(1, metrics.getProduceTPS());
double tpsScore = Math.min(100, ratio * 100);
// 综合评分
double total = backlogScore * WEIGHT_BACKLOG
+ rtScore * WEIGHT_RT
+ errorScore * WEIGHT_ERROR_RATE
+ tpsScore * WEIGHT_TPS_RATIO;
return new HealthScore(total, backlogScore, rtScore, errorScore, tpsScore);
}
}7.3 消费者自动化运维
# K8s HPA 基于 MQ 积压自动扩缩容
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: order-consumer-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: order-consumer
minReplicas: 3
maxReplicas: 20
metrics:
# 自定义指标:MQ 积压量
- type: External
external:
metric:
name: mq_backlog
selector:
matchLabels:
consumer_group: OrderConsumerGroup
target:
type: AverageValue
averageValue: "50000" # 每个实例处理 5 万条积压
behavior:
scaleUp:
stabilizationWindowSeconds: 30 # 快速扩容
policies:
- type: Percent
value: 100
periodSeconds: 30
scaleDown:
stabilizationWindowSeconds: 300 # 慢速缩容
policies:
- type: Percent
value: 10
periodSeconds: 60八、避坑指南
8.1 RocketMQ 常见坑
坑1: 消费者组订阅不一致
问题: 同一 ConsumerGroup 下不同消费者订阅不同 Topic
后果: 订阅关系被覆盖,部分消息不被消费
解决: 同一 ConsumerGroup 只订阅相同 Topic
坑2: 消费位点自动提交导致丢消息
问题: 消费失败但位点已提交
解决: 使用 ConsumeStatusType.RECONSUME_LATER 返回重试
坑3: 顺序消息消费阻塞
问题: 顺序消息某条消费失败,阻塞整个队列
解决: 设置最大重试次数,超过后发死信队列
坑4: 广播模式下的积压
问题: 广播模式无积压概念(每个消费者独立位点)
解决: 广播模式不适用于高吞吐场景8.2 Kafka 常见坑
坑1: Rebalance 风暴
问题: 消费者频繁加入/退出,导致持续 Rebalance
后果: Rebalance 期间所有消费者停止消费
解决: 合理设置 session.timeout 和 max.poll.interval.ms
坑2: 消费者处理超时
问题: max.poll.records 太大,处理时间超过 max.poll.interval.ms
后果: 消费者被踢出组,触发 Rebalance
解决: 减小 max.poll.records 或异步处理
坑3: 分区数不能减少
问题: Kafka 分区只能增加不能减少
后果: 早期规划不当导致后续无法缩容
解决: 初始规划留足余量,一般按 2 倍峰值估算
坑4: __consumer_offsets 主题膨胀
问题: 消费者组过多导致内部主题膨胀
解决: 定期清理不活跃的消费者组8.3 通用避坑清单
| 序号 | 坑点 | 影响 | 预防措施 |
|---|---|---|---|
| 1 | 消费者单线程处理 | 消费 TPS 低 | 设置合理的 consumeThreadMax |
| 2 | 消费者同步调用下游 | 下游慢拖垮消费 | 异步化 + 超时控制 |
| 3 | 无死信队列监控 | 消息丢失 | DLQ 告警 + 定期巡检 |
| 4 | 消费幂等未做 | 重复消费 | 基于唯一 ID 做幂等 |
| 5 | 积压告警阈值过高 | 发现太晚 | 按消费时间设阈值 |
| 6 | 扩容不看队列数 | 消费者空转 | 扩容前检查队列/分区数 |
| 7 | 重试无退避 | 雪崩效应 | 指数退避重试策略 |
九、面试要点
Q1: 消息队列积压怎么排查?
回答框架:
- 看积压量:通过 MQ 管理工具查看积压量和消费 TPS
- 看消费者:检查消费者实例数是否正常,是否有实例挂掉
- 看消费耗时:通过监控或日志查看单消息处理耗时
- 看下游服务:Arthas trace 定位慢调用根因
- 看 GC 日志:排除 Full GC 导致的停顿
Q2: 消息积压怎么紧急处理?
回答框架:
- 快速扩容:增加消费者实例(注意不能超过队列/分区数)
- 队列扩容:如果队列数不足,先扩队列再扩消费者
- 旁路消费:下游不可用时,先快速消费存储,后续补偿
- 降级策略:跳过非核心逻辑,保核心链路
- 背压限流:限制生产速率,防止积压持续恶化
Q3: 如何设计背压机制?
回答框架:
- 消费端背压:根据下游 RT 动态调整消费并发度(AIMD 算法)
- 生产端背压:根据积压量动态控制生产速率
- 限流兜底:Sentinel/RateLimiter 做硬性限流
- 自适应调整:基于滑动窗口的 P99 RT 反馈调节
Q4: RocketMQ 和 Kafka 积压排查有什么区别?
回答要点:
- 积压查看命令不同
- RocketMQ 有原生重试和死信队列,Kafka 需手动处理
- RocketMQ 队列可增减,Kafka 分区只能增
- Kafka 的 Rebalance 机制更复杂(协同器 + 分区分配策略)
- Kafka 消费者需要关注 max.poll.interval.ms,RocketMQ 不需要
Q5: 如何预防消息积压?
回答要点:
- 容量规划:按 2 倍峰值估算消费者和队列数
- 压测验证:上线前压测验证消费能力
- 监控告警:积压量 + 消费时间双重告警
- 自动扩缩容:K8s HPA 基于