凌晨两点,订单履约监控突然告警:下游消费延迟从秒级涨到 40 分钟,用户投诉「下单后一直不发货」。Kafka 消息积压了——这不是第一次,但每次根因都不同。本文复盘这次 Kafka 消息积压与重复消费的完整排查过程,从 consumer lag 定位到临时扩容止血,再到用幂等消费根治重复,给你一套可复用的排查套路。
一、事故背景:订单履约延迟告警
我们的订单链路是「下单服务 → Kafka → 履约消费者 → 仓储系统」。那天大促预热,下单量涨了 3 倍,履约消费者却没跟上。监控看到的是履约延迟一路飙升,但 Kafka 本身没报错,Broker 磁盘和 CPU 都正常——典型的消费端跟不上生产端,也就是消息积压。
报警来自 Prometheus 的履约延迟看板,阈值设的是 5 分钟,而当时已经逼近 40 分钟。更糟的是,积压不会自己消失——只要生产速率不减,lag 会线性累积,越拖恢复时间越长。这决定了我们的响应策略:先扩消费能力把 lag 拉平,再慢慢查为什么单条处理这么慢,而不是一上来就关生产端。
二、先看懂 Kafka 的”积压”是什么
Kafka 里每条消息在分区上有个递增的 offset。消费者组维护一个 committed offset(已消费到的位置),而生产者不断写入 log-end offset(最新位置)。两者之差就是 consumer lag(消费滞后量),也就是积压量。lag 持续上涨,说明消费速度 < 生产速度。
理解分区并行度很关键:一个消费者组内,同一个分区同一时刻只能被一个消费者实例持有。换句话说,分区数是 6,消费者实例最多也只能并行消费 6 个——实例开再多也白搭,多余的实例会处于空闲状态。这也是后面我们扩容能立刻见效的前提:当时只有 2 个实例,远没吃满 6 个分区的并行度,等于白白浪费了 4 个分区的吞吐。
2.1 用命令行看 lag
# 查看消费者组的积压(Kafka 2.4+ 自带命令)
kafka-consumer-groups.sh \
--bootstrap-server broker1:9092 \
--describe --group order-fulfill-group
# 输出关键列:
# TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# orders 0 1002300 1009800 7500 <-- 这就是积压
# orders 1 998800 1000100 1300
三、定位根因:为什么消费不过来
积压只是现象,根因通常在消费端。我们依次排查:消费者实例有没有挂(rebalance 风暴?)、单条消息处理是不是太慢(同步调用下游?)、是不是消费线程被锁住。用监控看每秒消费条数最直观。
其中 rebalance 风暴尤其隐蔽:消费者心跳超时,或单次 poll 处理时间超过 max.poll.interval.ms,就会被协调者踢出组,触发全组重新分配分区。分配期间所有消费者暂停消费,lag 瞬间飙升;如果慢调用持续,会陷入「踢出—重平衡—再踢出」的死循环。我们用 kafka-consumer-groups.sh 的 STATE 列确认 group 当时正处于频繁 Rebalance,这正是 lag 抖动型上涨的典型信号。
3.1 看消费 TPS 是否上不去
# 临时加日志或用 Micrometer 打点,观察单实例消费速率
@KafkaListener(topics = "orders", groupId = "order-fulfill-group")
public void listen(OrderEvent e) {
long t0 = System.nanoTime();
fulfillService.process(e); // 业务处理
long cost = (System.nanoTime() - t0) / 1_000_000;
metrics.timer("consume.cost").record(cost, TimeUnit.MILLISECONDS);
}
我们发现单条处理平均 80ms,且里面同步调用了一次外部 HTTP(风控校验)。生产 3 倍流量下,单实例消费 TPS 卡在 ~120,而分区有 6 个、消费者只有 2 个实例——明显是消费者实例数 < 分区并行度 + 单条处理有慢调用。
3.2 常见根因对照
| 现象 | 最可能根因 | 验证方式 |
|---|---|---|
| lag 匀速上涨 | 消费速度 < 生产速度(慢调用/实例少) | 看消费 TPS 与处理耗时 |
| lag 突然暴涨后平稳 | 某批大消息 / 批量任务灌入 | 查生产端是否有定时任务 |
| lag 反复抖动 | consumer rebalance 风暴 | 看 group 是否频繁重平衡 |
| 部分分区 lag 高 | 消费不均(key 倾斜) | 对比各 partition lag |
四、止血:先让队列动起来
线上第一优先级是止血,不是找根因。我们做了两件事:临时把消费者实例扩容到 6 个(吃满分区并行度),并把风控校验改成异步+降级(超时直接放行)。10 分钟内 lag 归零。
4.1 临时扩容与跳过积压
# 紧急扩容消费者副本(K8s)
kubectl scale deployment order-fulfill-consumer --replicas=6
# 极端情况:积压太久且老数据已无用,可重置 offset 跳过
kafka-consumer-groups.sh --bootstrap-server broker1:9092 \
--group order-fulfill-group --reset-offsets \
--topic orders --to-latest --execute # ⚠️ 会丢消息,仅止损用
五、重复消费:比积压更隐蔽的坑
止血后我们发现一个更麻烦的问题:有订单被履约了两次。Kafka 的投递语义默认是「至少一次(at-least-once)」——消息处理成功、但 offset 提交前消费者崩溃,重启后会从旧 offset 重新拉取,于是这条消息被消费了两次。
5.1 为什么会有重复
关键在位移提交时机。如果先处理业务、再异步提交 offset,机器在「处理完但没提交」时挂了,重平衡后新消费者从已提交位置拉,就会重复。我们正是踩了这个坑:业务处理里调了外部 HTTP,偶发超时导致处理慢、期间实例被踢出触发 rebalance。
5.2 幂等三板斧
// 方案 A:业务唯一键(最稳)。用订单号做去重表/唯一索引
@Transactional
public void fulfill(OrderEvent e) {
int n = jdbcTemplate.update(
"INSERT IGNORE INTO fulfill_log(order_no) VALUES(?)", e.getOrderNo());
if (n == 0) return; // 已处理过,直接跳过 = 幂等
fulfillService.doFulfill(e);
}
// 方案 B:Redis SETNX 去重(适合高并发、可短暂容忍)
String key = "fulfill:" + e.getOrderNo();
if (!redis.setnx(key, "1", Duration.ofHours(24))) return;
fulfillService.doFulfill(e);
去重方案要按业务一致性要求选:金融/订单类用数据库唯一键最稳;日志/计数类用 Redis 足够。核心原则——消费逻辑必须可重入。
5.3 三种去重方案对照
| 方案 | 一致性 | 成本 | 适用 |
|---|---|---|---|
| 数据库唯一键 | 强(事务保证) | 低 | 订单/交易等核心 |
| Redis SETNX | 最终一致 | 低 | 高并发可容忍短时重复 |
| 本地内存去重 | 弱(重启失效) | 极低 | 仅调试/非关键 |
六、根治:把可靠性写进设计
止血解决眼前,根治靠设计。我们做了三件事:① 慢调用限流+超时+降级;② 用手动提交(ack 后提交)配合幂等,避免「处理了没提交就崩」;③ 接死信队列(DLQ),处理失败 N 次后转入 DLQ,不再阻塞主业。
6.1 手动提交 + 死信队列
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> factory() {
var f = new ConcurrentKafkaListenerContainerFactory<>();
f.setConsumerFactory(consumerFactory());
f.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); // 处理完手动 ack
return f;
}
@KafkaListener(topics = "orders", groupId = "order-fulfill-group")
public void listen(ConsumerRecord<String,String> rec, Acknowledgment ack) {
try {
fulfillService.process(parse(rec.value()));
ack.acknowledge(); // 成功才提交位移
} catch (RetryExhaustedException ex) {
dlqTemplate.send("orders-dlq", rec.value()); // 转入死信队列
}
}
七、监控与告警:别等用户投诉
最好的排查是让它不发生。我们用 Kafka Exporter 把 lag 暴露给 Prometheus,超过阈值就告警,而不是等履约延迟。
# Prometheus 告警规则
- alert: KafkaConsumerLagHigh
expr: kafka_consumergroup_lag > 5000
for: 5m
labels: { severity: warning }
annotations:
summary: "消费者组 {{ $labels.consumergroup }} lag 超 5000"
八、六步排查清单
| 步骤 | 动作 | 命令/指标 |
|---|---|---|
| 1 看 lag | 确认是否真积压 | kafka-consumer-groups –describe |
| 2 看分布 | 定位单分区还是全局 | 对比各 partition LAG |
| 3 看消费 TPS | 消费是否过慢 | 消费速率 / 处理耗时打点 |
| 4 看 rebalance | 是否频繁重平衡 | group 状态 / 日志 |
| 5 止血 | 扩容 / 降级慢调用 | kubectl scale |
| 6 根治 | 幂等 + 手动提交 + DLQ + 告警 | 代码改造 + 监控 |
九、小结
Kafka 消息积压与重复消费,本质是两个独立又相关的问题:积压是速度不匹配,重复是提交时机 + 无幂等。排查时先止血(扩容+降级),再定位(lag/TPS/rebalance),最后根治(幂等消费、手动提交、死信队列、lag 告警)。把这套流程沉淀成监控和清单,下一次大促就能睡个安稳觉。相关排查思路也适用于其他资源瓶颈:比如 HikariCP 连接泄漏、慢 SQL 拖垮生产库、磁盘 I/O 打满雪崩 的排查套路一脉相承。




