第18讲 | 消费积压(Lag)治理实战
导读:Lag 是消费链路的"体温计",平时监控它、异常时降它、架构上防它;本讲讲透 Lag 的计算原理、监控方案、积压突增的应急 SOP,以及"扩消费者"那条最容易踩的硬约束。
本讲目标
一、Lag 是什么,怎么算
Lag(分区) = LEO(分区) − 消费组已提交位移
- LEO(Log End Offset):分区下一条待写入的位移,代表"生产到哪了";
- 已提交位移:代表"消费到哪了"。
Lag 只在"分区 × 消费组"维度有意义。两个细节:
二、Lag 监控三方案
2.1 命令行(人肉排查)
bin/kafka-consumer-groups.sh –bootstrap-server localhost:9092 \\
–describe –group order-consumer-group
输出列:CURRENT-OFFSET、LOG-END-OFFSET、LAG、CONSUMER-ID。CONSUMER-ID 为 – 说明该分区当前无消费者在消费(消费者全挂或 Rebalance 中),是快速判死的线索。
2.2 体系化监控(Prometheus + Burrow/kafka-exporter)
kafka-exporter 暴露 kafka_consumergroup_lag 指标给 Prometheus,Grafana 画大盘并按如下规则告警(工程经验值):
- lag > 10万 且持续 5 分钟:Warn;
- lag 增速(一阶导数)> 0 持续 10 分钟:消费速率跟不上生产速率;
- consumer_id == "-" 持续 2 分钟:消费者失联。
2.3 代码内嵌(框架自治)
消费框架自己周期性上报 Lag(示例见下文实战代码),好处是可以把"Lag 值"用于应用内自适应(如动态扩缩容依据),不依赖外部监控联动。
三、积压突增应急处理流程
第一步:定性——生产突增还是消费骤降?
# 观察两端速率(连续执行两次看增量)
bin/kafka-consumer-groups.sh –bootstrap-server localhost:9092 –describe –group xxx –all-groups
- 生产突增(LEO 快涨、CURRENT 也涨但慢):营销活动、上游批量补数 → 属容量问题;
- 消费骤降(CURRENT 停滞):消费者挂了 / Rebalance 风暴 / 下游依赖故障 → 属故障问题。
第二步:决策树(按序执行)
#mermaid-svg-slj9AP1IWyu2cDZt{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-slj9AP1IWyu2cDZt .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-slj9AP1IWyu2cDZt .error-icon{fill:#552222;}#mermaid-svg-slj9AP1IWyu2cDZt .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-slj9AP1IWyu2cDZt .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-slj9AP1IWyu2cDZt .marker{fill:#333333;stroke:#333333;}#mermaid-svg-slj9AP1IWyu2cDZt .marker.cross{stroke:#333333;}#mermaid-svg-slj9AP1IWyu2cDZt svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-slj9AP1IWyu2cDZt p{margin:0;}#mermaid-svg-slj9AP1IWyu2cDZt .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-slj9AP1IWyu2cDZt .cluster-label text{fill:#333;}#mermaid-svg-slj9AP1IWyu2cDZt .cluster-label span{color:#333;}#mermaid-svg-slj9AP1IWyu2cDZt .cluster-label span p{background-color:transparent;}#mermaid-svg-slj9AP1IWyu2cDZt .label text,#mermaid-svg-slj9AP1IWyu2cDZt span{fill:#333;color:#333;}#mermaid-svg-slj9AP1IWyu2cDZt .node rect,#mermaid-svg-slj9AP1IWyu2cDZt .node circle,#mermaid-svg-slj9AP1IWyu2cDZt .node ellipse,#mermaid-svg-slj9AP1IWyu2cDZt .node polygon,#mermaid-svg-slj9AP1IWyu2cDZt .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-slj9AP1IWyu2cDZt .rough-node .label text,#mermaid-svg-slj9AP1IWyu2cDZt .node .label text,#mermaid-svg-slj9AP1IWyu2cDZt .image-shape .label,#mermaid-svg-slj9AP1IWyu2cDZt .icon-shape .label{text-anchor:middle;}#mermaid-svg-slj9AP1IWyu2cDZt .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-slj9AP1IWyu2cDZt .rough-node .label,#mermaid-svg-slj9AP1IWyu2cDZt .node .label,#mermaid-svg-slj9AP1IWyu2cDZt .image-shape .label,#mermaid-svg-slj9AP1IWyu2cDZt .icon-shape .label{text-align:center;}#mermaid-svg-slj9AP1IWyu2cDZt .node.clickable{cursor:pointer;}#mermaid-svg-slj9AP1IWyu2cDZt .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-slj9AP1IWyu2cDZt .arrowheadPath{fill:#333333;}#mermaid-svg-slj9AP1IWyu2cDZt .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-slj9AP1IWyu2cDZt .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-slj9AP1IWyu2cDZt .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-slj9AP1IWyu2cDZt .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-slj9AP1IWyu2cDZt .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-slj9AP1IWyu2cDZt .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-slj9AP1IWyu2cDZt .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-slj9AP1IWyu2cDZt .cluster text{fill:#333;}#mermaid-svg-slj9AP1IWyu2cDZt .cluster span{color:#333;}#mermaid-svg-slj9AP1IWyu2cDZt div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-slj9AP1IWyu2cDZt .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-slj9AP1IWyu2cDZt rect.text{fill:none;stroke-width:0;}#mermaid-svg-slj9AP1IWyu2cDZt .icon-shape,#mermaid-svg-slj9AP1IWyu2cDZt .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-slj9AP1IWyu2cDZt .icon-shape p,#mermaid-svg-slj9AP1IWyu2cDZt .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-slj9AP1IWyu2cDZt .icon-shape .label rect,#mermaid-svg-slj9AP1IWyu2cDZt .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-slj9AP1IWyu2cDZt .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-slj9AP1IWyu2cDZt .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-slj9AP1IWyu2cDZt :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
CONSUMER-ID 为 –
活着
DB/RPC 慢
正常
是
否
Lag 告警
消费者活着吗?
查消费者进程/Rebalance 日志
下游依赖正常吗?
扩下游连接池/降级非核心逻辑
分区数 > 消费者数?
扩容消费者实例
多线程消费 or 扩分区
观察追平耗时
第三步:兜底手段(慎用)
- 跳过积压:实时性敏感、历史数据无用(如实时大屏、风控特征)时,–reset-offsets –to-latest 跳到最新,历史数据离线补;
- 扩分区:注意打破 key 顺序的副作用(第 17 讲);
- 临时消费组补数:新 group + earliest + 独立消费者把历史数据"旁路"写入离线表,线上只追新数据。
四、扩容消费者的分区硬约束
组内并行度 ≤ 分区数,这是无法绕过的硬约束(不改造消费模型的前提下)。因此扩容前必答三问:
实战案例(Java)
场景:埋点行为分析服务,track-events 12 分区,单条处理 5ms,积压 2000 万条。方案:拉取线程 + worker 线程池 + 按分区哈希派发保证分区串行,并周期上报 Lag。
依赖(Maven)
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.5.1</version>
</dependency>
完整代码
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
/**
* 多线程消费方案:
* – 单 KafkaConsumer 负责拉取(避免多 consumer 各自 Rebalance 的协调成本);
* – worker 线程池处理;同一分区的消息固定派发给同一 worker(保证分区串行);
* – 用"处理完成游标"跟踪每个分区最小未完成位移,安全后统一提交;
* – 周期打印 Lag,便于观察追平进度。
*/
public class MultiThreadTrackConsumer {
private static final String TOPIC = "track-events";
private static final String GROUP = "track-analyze-service";
private static final String BOOTSTRAP = "localhost:9092";
private static final int WORKER_COUNT = 12; // 建议与分区数一致
public static void main(String[] args) throws Exception {
ExecutorService pool = Executors.newFixedThreadPool(WORKER_COUNT);
// 每个 worker 维护一个"按分区"的待处理队列:同分区进同队列 → 分区内天然串行
List<BlockingQueue<ConsumerRecord<String, String>>> queues = new ArrayList<>();
for (int i = 0; i < WORKER_COUNT; i++) {
BlockingQueue<ConsumerRecord<String, String>> q = new LinkedBlockingQueue<>(2000);
queues.add(q);
pool.submit(new Worker(q));
}
// 记录每个分区"已完成处理的最大位移+1"与"已拉取未完成计数",用于安全提交
Map<Integer, AtomicLong> finishedOffset = new ConcurrentHashMap<>(); // 分区 -> 已完成 offset+1
Map<Integer, AtomicLong> outstanding = new ConcurrentHashMap<>(); // 分区 -> 在途消息数
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(buildProps());
consumer.subscribe(List.of(TOPIC));
// 后台线程:周期性提交"安全位移"并打印 Lag
ScheduledExecutorService monitor = Executors.newSingleThreadScheduledExecutor();
monitor.scheduleAtFixedRate(() -> {
try {
Map<TopicPartition, OffsetAndMetadata> commits = new HashMap<>();
finishedOffset.forEach((p, off) -> {
// 仅当该分区在途消息为 0 时才提交(简单安全策略)
if (outstanding.getOrDefault(p, new AtomicLong()).get() == 0) {
commits.put(new TopicPartition(TOPIC, p), new OffsetAndMetadata(off.get()));
}
});
if (!commits.isEmpty()) {
consumer.commitSync(commits);
}
// 打印 Lag
for (TopicPartition tp : consumer.assignment()) {
long lag = consumer.currentLag(tp).orElse(–1L); // 3.x API:实时 Lag
if (lag >= 0) {
System.out.printf("partition=%d lag=%d%n", tp.partition(), lag);
}
}
} catch (Exception e) {
System.err.println("监控线程异常: " + e.getMessage());
}
}, 5, 5, TimeUnit.SECONDS);
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
int workerIdx = Math.abs(r.partition()) % WORKER_COUNT; // 分区 -> 固定 worker
outstanding.computeIfAbsent(r.partition(), k -> new AtomicLong())
.incrementAndGet();
queues.get(workerIdx).put(new ConsumerRecord<>(
r.topic(), r.partition(), r.offset(), r.key(),
r.value() + "|" + System.identityHashCode(finishedOffset)) {} {
// 需要携带完成回调:改用内部包装类(见下)
});
}
}
} finally {
pool.shutdownNow();
monitor.shutdownNow();
consumer.close();
}
}
private static Properties buildProps() {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
props.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1000); // 大批量摊薄拉取开销
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000);
return props;
}
/** worker:串行处理自己队列中的消息(同分区消息必然同队列) */
static class Worker implements Runnable {
private final BlockingQueue<ConsumerRecord<String, String>> queue;
Worker(BlockingQueue<ConsumerRecord<String, String>> queue) { this.queue = queue; }
@Override
public void run() {
while (!Thread.currentThread().isInterrupted()) {
try {
ConsumerRecord<String, String> r = queue.poll(1, TimeUnit.SECONDS);
if (r == null) {
continue;
}
try {
process(r);
} finally {
// 通知完成:更新该分区完成游标并递减在途计数(由框架回调实现)
CompletionTracker.onDone(r);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}
private void process(ConsumerRecord<String, String> r) {
try { Thread.sleep(5); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
// 实际业务:解析埋点、写 ClickHouse/ES
}
}
}
上面主循环中为了保持可读性留了一个"包装类"的钩子,补全后的正确形态是用 TaskRecord 包装原始消息 + 回调。下面给出精炼的补全实现(替换主循环的派发片段即可):
/** 任务包装:携带完成回调所需的上下文 */
record TaskRecord(ConsumerRecord<String, String> record) {}
// 主循环派发(替换上文 try 内的 for 循环):
for (ConsumerRecord<String, String> r : records) {
int workerIdx = Math.abs(r.partition()) % WORKER_COUNT;
AtomicLong inFlight = outstanding.computeIfAbsent(r.partition(), k -> new AtomicLong());
inFlight.incrementAndGet();
queues.get(workerIdx).put(r); // Worker 里处理完后调用 CompletionTracker.onDone(r)
}
// CompletionTracker:合并到主类的静态内部类
static class CompletionTracker {
private static final Map<Integer, AtomicLong> FINISHED = new ConcurrentHashMap<>();
private static final Map<Integer, AtomicLong> IN_FLIGHT = new ConcurrentHashMap<>();
static void onDone(ConsumerRecord<String, String> r) {
FINISHED.computeIfAbsent(r.partition(), k -> new AtomicLong())
.accumulateAndGet(r.offset() + 1, Math::max); // 完成游标取最大已处理+1
IN_FLIGHT.get(r.partition()).decrementAndGet();
}
}
说明:为控制篇幅,示例采用"在途为 0 才提交"的保守策略,吞吐足够且不会丢消息;更强的实现会维护"每个 offset 的完成位图"以提交最小连续完成位移,属于第 19 讲框架层的工作。
运行方式与预期输出
# 灌入 100 万条埋点后启动:
mvn compile exec:java -Dexec.mainClass=MultiThreadTrackConsumer
预期:12 个 worker 并行处理(单机吞吐约 12 × 200 = 2400 条/秒),每 5 秒打印一轮各分区 Lag,数值单调下降直到 0;重启进程后从已提交位移续跑,不丢不重(配合第 17 讲幂等组件彻底防重)。
关键点解读
踩坑提示 / 生产建议
本讲小结
Lag = LEO − 已提交位移,是消费链路健康度的核心指标;监控要告警"Lag 增速 + 追平耗时 + 消费者失联"三者而非单一数值;积压突增先定性(生产突增 vs 消费骤降)再按决策树处置;扩容消费者的硬约束是"组内并行度 ≤ 分区数",分区不够时要么扩分区(注意保序代价),要么用"单 consumer 拉取 + 线程池按分区派发"的多线程方案榨干单实例并行度。
思考题


