消息积压、扩缩容与可观测性:判断恢复速度与下游承受力
一个消费组落后十万条消息,可能十秒就能追平,也可能持续增长到保留期耗尽。差别取决于新消息进入速度、业务实际完成速度、分区分配和下游容量。数量提供当前位置,变化速率与最老消息年龄才能帮助判断还剩多少恢复时间。
消息停在哪个阶段
几种常见积压指标
| 指标 | 观察对象 | 需要继续核对的内容 |
|---|---|---|
| RabbitMQ messages_ready | 等待正常交付的消息 | 是否有在线消费者、路由是否落错队列 |
| RabbitMQ messages_unacknowledged | 已交付但未确认的消息 | prefetch、处理时长、客户端停顿和 ACK |
| Kafka 组提交位点 Lag | 末尾位置与 committed offset 的差 | 是否只是尚未提交,还是业务也尚未完成 |
| 客户端 fetch/position Lag | 当前拉取进度与末尾位置之差 | 进程内队列是否囤积、是否提前读取大量记录 |
| 最老未完成消息年龄 | 从约定起点到当前的时间 | 起点是创建、Broker 接收、交付还是业务开始 |
| 业务未完成数量 | 尚未到达约定业务终态的对象 | 重试、死信、外部结果未知是否被统计进去 |
RabbitMQ ready 与 unacked 是常用队列观察值,但内部死信转交等状态还有自己的保留记录。Kafka Lag 也有多种实现口径。监控面板应明确公式和 API 来源,不能只显示一个没有定义的 lag 字段。RabbitMQ 监控指标、Kafka 监控指标
拉取位置可以先于业务提交前进
Kafka Topic partition 0:记录位置 0 … 11;末尾位置 12
poll 已返回 12 条
├── consumer position = 12
├── committed offset = 0
└── 按组提交位点计算的 Lag = 12
业务完成后 commitSync(nextOffsets)
├── committed offset = 12
└── Lag = 0poll 返回记录时,消费者的 position 会向前移动。组提交位点保存的是恢复时使用的位置,应在业务完成后推进。提前提交可以让面板更好看,却会使进程退出后的消费者跳过尚未完成的记录。KafkaConsumer 位置与提交
Offset 差值也不能无条件当作消息数量。事务控制记录、压缩和其他位置空洞会使某些位置没有可交付业务记录;read_committed 可见范围还受稳定事务位置影响。将 leader 日志末尾、复制水位、事务可见末尾和消费者提交位置混用,会得到难以解释的差值。
下面的实验使用没有事务与压缩的单分区 Topic,明确把 Kafka Admin 返回的末尾位置减去组 committed offset,所以 12 个位置正好对应 12 条普通记录。Admin 位点查询
成功率与年龄应一起观察
只统计消费处理成功次数,会漏掉消息一直未被分配、在延迟队列等待或已经进入死信的情况。对订单通知等业务,应定义从订单提交到通知完成的时限,并分别记录排队、处理、重试和外部等待。
如果消息生产端跨机器,occurredAt 的时钟偏差也会进入年龄计算。Broker 接收时刻适合估算队列等待,业务发生时刻更接近用户感受到的延迟;保留两者并注明来源,比从一个混合时间戳计算 p99 更可靠。
用真实位点和消费组检查扩容空间
准备环境
下载完整工程,解压进入 message-event-lab。使用 Linux shell、Docker Engine 与 Compose v2+,普通宿主用户具有 Docker 权限。Maven 3.9.12 / Java 25 构建,Java 字节码目标为 17;Kafka Broker 与 client 4.3.1,PostgreSQL 18.6。
compose.kafka.yaml 启动一个 controller 与三个 Broker;叠加 compose.cdc.yaml 后可以启动数据库供吞吐实验使用。这里不用启动 Connect。服务仅位于 me12-cdc Docker 网络,数据库应用使用 app 角色,Maven 显式映射宿主 UID/GID 与可写缓存。
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
dc() { docker compose -f compose.kafka.yaml -f compose.cdc.yaml "$@"; }
dc up -d controller kafka1 kafka2 kafka3 postgres
ready=0
for attempt in $(seq 1 60); do
if dc exec -T kafka1 /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka1:9092 --list >/dev/null 2>&1 &&
dc exec -T postgres pg_isready -h 127.0.0.1 -U labadmin -d lab >/dev/null; then
ready=1
break
fi
sleep 2
done
test "$ready" -eq 1 || { dc logs --tail=60; exit 1; }
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/maven-cache \
-v "$PWD:/workspace" -v "$PWD/.m2:/maven-cache" -w /workspace \
maven:3.9.12-eclipse-temurin-25 \
mvn -B -ntp '-Dmaven.repo.local=/maven-cache' '-Duser.home=/tmp' \
-pl kafka -am clean package dependency:copy-dependencies
run_message_java() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-cdc_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:25.0.4_7-jdk \
java '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'kafka/target/classes:kafka/target/dependency/*' "$@"
}
run_message_java example.BacklogLab laglag: end=12 position=12 committed=0 lag=12
lag: committed=12 lag=0程序先提交明确的起点 0,再拉取 12 条记录。在 position=12 时通过独立 Admin 请求取得 committed=0,之后提交本次 poll 的 nextOffsets,再次查询 Lag=0。该模式只演示位点,不把读取记录称为业务处理完成。
三个分区与四个普通消费者
run_message_java example.KafkaGroupLab 4关键结果如下,成员排列次序可能不同:
key: order-42 records=6 partitions=[0]
group: members=1 assigned=[3] union=3 disjoint=true
group: members=4 assigned=[0, 1, 1, 1] union=3 disjoint=true
group: members=1 assigned=[3] union=3 disjoint=trueKafkaGroupLab 为每个线程创建并独占一个真实 KafkaConsumer,使用 subscribe 加入同组,通过实际 assignment/rebalance 回调检查分配稳定。它断言三分区没有被同组普通消费者同时拥有,增加到四成员后必有一个成员没有分区;退出额外成员后,剩余成员重新获得三个分区。
实验固定 classic 协议与 RangeAssignor,以便观察。Kafka 新消费组协议有不同协调流程;Share Consumer 也不是这里的普通分区独占模型。使用哪种 API 和协议,应先在运行配置中确认。消费组重平衡协议
发送到同一 key=order-42 的六条消息实际落在一个分区。即使其他分区空闲,也无法简单把这个 key 的严格串行处理拆到四个普通消费者上。固定 key 对应的分区数扩张还可能改变后续路由,历史与新事件的顺序需要迁移设计;完整 key 与分区机制见 Kafka 运行链。
分区内部异步处理还受提交顺序约束
同一个消费者可以把已获得的记录交给工作线程,但提交位点必须只越过连续完成的前缀。例如 offset 10 已完成、11 仍在等待、12 已完成,此时不能把下一位置提交到 13,否则重启将跳过 11。
下一个实验采用更简单的批次方式:每次 poll 最多取 6 条,由两个工作线程执行数据库事务;等本批所有 Future 成功,才提交这批 nextOffsets。这样便于确认语义,吞吐代价是最快的记录也会等本批最慢者完成。只有主线程访问 KafkaConsumer;工作线程仅使用各自 JDBC 连接。
完整位点与容量实验入口
源码位于 kafka/src/main/java/example/BacklogLab.java。公共 KafkaLabSupport 负责创建三副本 Topic、配置序列化和查询超时,完整实现在 ZIP 中,并在 Kafka 运行链 给出。
package example;
import java.sql.DriverManager;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.kafka.clients.admin.OffsetSpec;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import static example.KafkaLabSupport.*;
public final class BacklogLab {
static long lag(String group,TopicPartition partition) throws Exception {
try(var a=admin()) {
long end=a.listOffsets(Map.of(partition,OffsetSpec.latest())).all().get(15,TimeUnit.SECONDS).get(partition).offset();
var committed=a.listConsumerGroupOffsets(group).partitionsToOffsetAndMetadata().get(15,TimeUnit.SECONDS).get(partition);
check(committed!=null,"group has no committed starting point");return end-committed.offset();
}
}
static void positions() throws Exception {
String topic="me12.backlog.positions",group="me12-backlog-positions";
resetTopic(topic,1);resetGroup(group);var partition=new TopicPartition(topic,0);
try(var p=producer()){for(int i=0;i<12;i++)p.send(new ProducerRecord<>(topic,0,"k","event-"+i)).get(15,TimeUnit.SECONDS);}
try(var c=consumer(group,"read_committed")) {
c.assign(List.of(partition));c.seekToBeginning(c.assignment());c.commitSync(Map.of(partition,new OffsetAndMetadata(0)));
var batch=read(c,12);
long before=lag(group,partition);check(c.position(partition)==12&&before==12,"read changed durable committed offset");
System.out.println("lag: end=12 position=12 committed=0 lag="+before);
c.commitSync(batch.nextOffsets());check(lag(group,partition)==0,"lag not reduced after commit");
System.out.println("lag: committed=12 lag=0");
}
}
static java.sql.Connection database() throws Exception {
return DriverManager.getConnection("jdbc:postgresql://postgres:5432/lab","app","app-lab-only");
}
static void capacity() throws Exception {
String topic="me12.backlog.capacity",group="me12-backlog-capacity";
resetTopic(topic,1);resetGroup(group);var partition=new TopicPartition(topic,0);
try(var db=database();var s=db.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS backlog_effect(event_id text PRIMARY KEY)");s.execute("TRUNCATE backlog_effect");
}
try(var p=producer()){for(int i=0;i<40;i++)p.send(new ProducerRecord<>(topic,0,"k","initial-"+i)).get(15,TimeUnit.SECONDS);}
var stop=new AtomicBoolean();var arrivals=new AtomicInteger();var active=new AtomicInteger();var peak=new AtomicInteger();
var producerFailure=new java.util.concurrent.atomic.AtomicReference<Throwable>();
var sender=new Thread(()->{
try(var p=producer()) {
while(!stop.get()) {
int id=arrivals.get();p.send(new ProducerRecord<>(topic,0,"k","live-"+id)).get(15,TimeUnit.SECONDS);
arrivals.incrementAndGet();Thread.sleep(100);
}
}catch(Throwable error){producerFailure.set(error);}
},"live-producer");
var properties=consumerProperties(group,"read_committed");properties.put("max.poll.records","6");
var workers=Executors.newFixedThreadPool(2);var slots=new Semaphore(2);int processed=0;boolean caughtUp=false;
long started=System.nanoTime(),deadline=started+TimeUnit.SECONDS.toNanos(30);
try(var c=new KafkaConsumer<String,String>(properties)) {
c.assign(List.of(partition));c.seekToBeginning(c.assignment());c.commitSync(Map.of(partition,new OffsetAndMetadata(0)));sender.start();
while(System.nanoTime()<deadline) {
var records=c.poll(Duration.ofMillis(200));var pending=new ArrayList<Future<?>>();
for(var record:records) {
slots.acquire();
pending.add(workers.submit(()->{
int running=active.incrementAndGet();peak.accumulateAndGet(running,Math::max);
try(var db=database();var s=db.createStatement()) {
db.setAutoCommit(false);s.execute("SELECT pg_sleep(0.05)");
try(var insert=db.prepareStatement("INSERT INTO backlog_effect VALUES (?)")){insert.setString(1,record.value());insert.executeUpdate();}
db.commit();
}catch(Exception error){throw new RuntimeException(error);}
finally{active.decrementAndGet();slots.release();}
}));
}
for(var future:pending)future.get(10,TimeUnit.SECONDS);
processed+=records.count();if(!records.isEmpty())c.commitSync(records.nextOffsets());
check(producerFailure.get()==null,"live producer failed: "+producerFailure.get());
if(processed>=40&&lag(group,partition)==0) {
if(!caughtUp){caughtUp=true;stop.set(true);sender.join(5000);check(!sender.isAlive(),"producer did not stop");}
if(processed==40+arrivals.get())break;
}
}
}finally{stop.set(true);sender.join(5000);workers.shutdown();check(workers.awaitTermination(10,TimeUnit.SECONDS),"workers still active");}
long elapsed=TimeUnit.NANOSECONDS.toMillis(System.nanoTime()-started);
check(caughtUp&&arrivals.get()>0&&processed==40+arrivals.get(),"backlog not drained while producer was active");
check(peak.get()==2&&active.get()==0,"downstream concurrency not bounded to two");
try(var db=database();var s=db.createStatement();var rows=s.executeQuery("SELECT count(*) FROM backlog_effect")) {
check(rows.next()&&rows.getInt(1)==processed,"committed business rows differ from processed records");
}
check(lag(group,partition)==0,"final lag not zero");
System.out.println("capacity: initial=40 liveArrivals="+arrivals.get()+" committedRows="+processed+" peakDbConcurrency="+peak.get()+" elapsedMs="+elapsed+" finalLag=0");
}
public static void main(String[] args) throws Exception {
check(args.length==1,"mode: lag|capacity");
switch(args[0]){case "lag"->positions();case "capacity"->capacity();default->throw new IllegalArgumentException("unknown mode");}
}
}带着实时流量恢复积压
先限制下游并发,再测净排空能力
run_message_java example.BacklogLab capacity程序先向 Kafka 写入 40 条消息作为积压,再启动实时生产线程,每次发布完成后等待 100 毫秒。消费端同时读取消息,每条在 PostgreSQL 中等待 50 毫秒后写入独立业务结果并提交;这是可控的下游延迟,不代表真实业务的数据库性能。
一次运行结果:
capacity: initial=40 liveArrivals=67 committedRows=107 peakDbConcurrency=2 elapsedMs=7682 finalLag=0liveArrivals、committedRows 和 elapsedMs 随机器及运行状态变化。程序断言 committedRows=40+liveArrivals、实时生产至少写入一条、数据库并发峰值为 2、最终 Lag=0。它在生产线程仍活动时第一次观测到积压归零,然后停止生产并处理可能正在发送的最后记录。
两张工作许可由 Semaphore 控制,工作线程数也为 2;Kafka 主线程取得许可后才能提交新数据库任务,因此待执行任务不会无限灌入线程池。真实业务应使用有界连接池,并同时限制 SQL/HTTP 调用时间。数据库连接每条新建是这个实验的简化,会显著影响耗时,不能从该数字推导服务器极限吞吐。
pg_sleep 使用服务器端等待,并在真实数据库事务中插入记录;完成数量来自数据库查询。PostgreSQL 时间与等待函数
恢复时间看的是消费与生产的差
在输入相对稳定、消息成本相近时:
λ:每秒新增消息数
μ:每秒成功完成并可确认的消息数
B:当前待完成消息数
净排空速度 = μ - λ
预计排空时间 ≈ B / (μ - λ),前提 μ > λ
例:B = 120000,λ = 800 / 秒,μ = 1000 / 秒
净排空 200 / 秒,预计需要 600 秒μ=λ 只能保持现有积压,μ<λ 时继续增长。计算 μ 要用成功完成业务的速度,失败后立即 requeue 的高交付速率不算恢复能力。若重试占用大量处理次数,应将重试放大纳入实际容量。
并非所有分区都一样快。三个分区分别落后 100、100、100000 条,而最后一个分区只有一个热点 key 时,总体平均吞吐会掩盖最长恢复时间。应按分区估算,整个消费组的恢复完成时间取决于最慢的剩余分区。
扩容、扩分区与降低输入分别解决什么
增加消费者实例可以使用尚未被充分利用的分区,并降低单实例负载;当每个分区已经有活跃普通消费者时,继续增加实例只会产生空闲成员。增加单实例工作并发可以提高每分区业务处理吞吐,但要解决连续提交、内存积累、顺序要求和下游限额。
增加 Topic 分区改变并行空间,也可能改变 key 的后续分布。需要严格按聚合排序时,宜提前规划分区、执行带切换点的迁移,或借助版本检查协调历史与新路由。不能把分区数当作一个无副作用的在线线程数。
下游已接近连接、锁、CPU 或第三方配额上限时,进一步扩消费者可能使 μ 下降。此时优先减少无效重试、优化慢调用、批量处理,必要时限制生产输入或降低非关键消费优先级。
RabbitMQ 消费者还应结合 prefetch 调整在途数量。prefetch 过大时,ready 看似减少,客户端却囤积大量未确认工作,实例故障后又一起返回队列。RabbitMQ consumer prefetch
监控、定位和安全恢复
四层观测保持同一业务关联
业务层:订单通知是否完成、最老未完成年龄、失败终态
应用层:处理耗时、等待线程池/连接池、重试原因、幂等命中
客户端层:发送等待、消费 poll、ACK/commit、rebalance
Broker 层:队列/分区、复制、磁盘、网络、内存与保留空间日志至少保留 event ID、业务对象 ID、消费组、Topic/queue、分区/offset 或 delivery attempt,以及最终操作结果。Trace 可以链接生产与消费的上下文;异步等待很长时,消费 span 通常应使用符合消息语义的关联方式,并避免把所有等待都算作某一次 HTTP 请求耗时。OpenTelemetry 消息语义约定
监控不要把 event ID、订单号或完整异常正文做成高基数 metric label;这些内容更适合日志与跟踪系统。指标保留可聚合的 Topic、消费组、错误类别和部署版本,异常样本再通过 ID 回查。
从形状判断瓶颈位置
| 指标组合 | 主要怀疑方向 | 下一步 |
|---|---|---|
| ready 高、unacked 低、消费者少 | 消费者未启动、订阅/权限或分配异常 | 查成员、连接、消费资源名与启动错误 |
| ready 下降、unacked 持续升高 | 在途工作过多或业务处理卡住 | 查线程栈、连接池等待、prefetch 与调用超时 |
| Kafka position 前进、committed 不动 | 批次未完成、提交失败或长事务 | 找最早未完成 offset,检查 commit 错误 |
| 只有一个分区 Lag 和年龄高 | 热点 key 或该分区消息特别昂贵 | 采样 key 分布、处理成本和实例分配 |
| 实例增加后数据库 p99 更差 | 下游争用已成为限制 | 降低并发,核对锁等待与服务配额 |
| Lag 突然归零、业务结果未增加 | 组被重置、提前提交或统计口径变化 | 查操作历史、group ID、业务表和消费者版本 |
| 消息进入 DLQ 后主队列恢复 | 坏消息已被隔离 | 仍需修复死信和业务结果,不能只关闭积压告警 |
Lag 告警应结合持续时长、增长趋势与业务时限。短时部署重平衡造成的尖峰,与最老消息接近保留期限的持续积压,处理优先级不同。预测排空时间若已经超过业务截止或消息保留期,就要启动有授权的限流、补偿或人工处置。
先保留可恢复数据
发现消息即将过期时,先确认 Topic 保留设置、剩余磁盘、下游吞吐和业务可接受延迟,再决定是否调整保留期。增加保留可能加大磁盘压力;直接把 offset 重置到 latest 会跳过历史记录,需要明确的业务授权和可审计记录。
缩容时让实例停止接收新工作,等待可完成的在途事务,按正确位点提交,再离开消费组。无法完成的工作应保留重投条件。消费位点重置、队列 purge 和批量死信重放都不是普通只读排查动作。
实验模式会删除并重建自己命名的 me12.backlog.* Topic、重置对应实验组,并清空自己的 backlog_effect 表。不要对真实业务同名资源执行。结束使用 dc stop 保留数据;确认 me12-cdc 项目的消息与数据库均可丢弃后才执行 dc down --volumes。
