Kafka 消息积压与重复消费排查实录

凌晨两点,订单履约监控突然告警:下游消费延迟从秒级涨到 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 打满雪崩 的排查套路一脉相承。

上一篇 Dev Container 实战:团队开发环境一次配好
下一篇 Text2SQL 落地实战:自然语言转 SQL 与安全执行