ts系统mq传输延迟三小时 一个真实故障的排查与解决
事情是这样的。
那天是周二凌晨两点,我正在值夜班,监控大屏上突然跳出一大片红色告警。MQ的积压曲线像过山车一样直线飙升,从几千条直接干到几百万条,而且没有任何回落的迹象。更离谱的是,数据延迟竟然已经达到了三小时。
说实话,当时第一反应是懵的。三小时?这意味着生产端生产的数据,到消费端的时候已经过期了。业务方要是这时候来问,我都不知道该怎么解释。
不过慌也没用,既然排查故障这事儿我干过不少,索性按流程一步步来。
一、先搞清楚问题出在哪儿
第一步,不能急着动手,得先看清全局。
我打开监控平台,看到几个关键指标:
- MQ队列堆积量:从正常的几百条飙升到800万+
- 消费端吞吐量:从平时的5000条/秒骤降到接近0
- 生产端:仍在正常发消息,速度没变
- Broker节点状态:全部显示健康,没有宕机
这就很奇怪了。生产端正常,Broker正常,但数据就是流不动,堆积在队列里。
我第一反应是检查消费端。登录到消费集群,发现所有消费节点的CPU使用率都在95%以上,内存也快打满了。
# 查看消费节点资源使用情况
top -p $(pgrep -f consumer)
结果不出所料,消费进程CPU飙升,内存占用很高。我进一步查看进程的堆栈信息:
# 查看Java进程的CPU使用情况
jstack <pid> > consumer_stack.log
# 查看GC情况
jstat -gcutil <pid> 1000 10
jstat的结果显示,Young GC频繁触发,而且每次GC pause时间超过2秒,Full GC每15分钟就触发一次。
问题大概有方向了——消费端因为GC问题导致处理性能急剧下降,消息堆积,然后堆积的消息又进一步加剧了内存压力,形成恶性循环。
但这只是表象,真正的原因是什么呢?
二、深挖根因
2.1 确认消费端的处理逻辑
我登录到一台消费节点,查看最近的日志:
2024-01-16 02:15:23.456 ERROR [consumer-thread-12] c.s.m.MetricsConsumer - Batch processing failed, batch size: 500, error: OutOfMemoryError: Java heap space
2024-01-16 02:15:23.789 WARN [consumer-thread-12] o.a.k.c.c.internals.ConsumerCoordinator - Commit offset failed, will retry...
OutOfMemoryError。看来内存确实是个问题。但为什么突然OOM呢?按理说这个消费集群已经稳定运行了半年,从来没出现过这种问题。
我查看了近三天的变更日志,发现有一个关键信息:
两天前,生产环境的TS系统(时序数据库)升级了协议版本,从Protobuf V2升级到了V3。
这会不会是关联点?
2.2 分析消息大小变化
我随机抓了几条最近消费的消息,对比大小:
# 消费一条消息并统计大小
kafka-console-consumer.sh --bootstrap-server broker1:9092 \
--topic ts-metrics-topic \
--from-beginning \
--max-messages 1 \
--property print.key=true
# 统计消息平均大小
kafka-run-class.sh kafka.tools.GetOffsetShell \
--broker-list broker1:9092 \
--topic ts-metrics-topic \
--time -1
统计结果显示,升级Protobuf V3后,单条消息的平均大小从256字节增加到了1.2KB,增长了将近5倍。
这就解释得通了。
2.3 追踪内存增长路径
为了进一步确认,我使用Java的堆dump工具,在OOM前导出了一次堆内存快照:
# 触发heap dump(如果配置了-XX:+HeapDumpOnOutOfMemoryError会自动生成)
jmap -dump:format=b,file=heap_dump.hprof <pid>
# 然后用MAT工具分析
# Eclipse MAT -> Leak Suspects Report
通过MAT(Memory Analyzer Tool)分析堆快照,我发现了一个关键问题:
消费端在处理每条消息时,都会调用TS系统的SDK进行一次数据校验。而TS SDK在V3版本中,内部维护了一个全局的缓存池,用于缓存schema信息。由于每次消费的消息都携带了完整的schema定义(V3的特性),这个缓存池在持续不断地增长,最终撑爆了堆内存。
核心问题代码段(脱敏后):
// 消费端处理逻辑
public void onMessage(MetricMessage message) {
// 1. 解析消息(V3 Protobuf,包含完整schema)
MetricProto.Metric metric = MetricProto.Metric.parseFrom(message.getBody());
// 2. 调用TS SDK校验(问题就在这里)
TSDatabaseClient client = TSDatabaseClient.getInstance();
// 每次调用都会尝试加载schema缓存
client.validateMetric(metric);
// 3. 写入业务数据库
businessDao.insert(metric);
}
而TS SDK内部的缓存逻辑大致如下:
// TS SDK 内部(简化版)
public class TSDatabaseClient {
// 全局缓存,key为schema ID,value为解析后的schema对象
private static final Cache<String, Schema> SCHEMA_CACHE =
CacheBuilder.newBuilder()
.maximumSize(10000) // 本意是限制10000条
.build();
public boolean validateMetric(Metric metric) {
String schemaKey = metric.getSchemaId();
// 问题:V3版本中,schemaId是动态生成的,几乎每条消息都不同
// 导致缓存命中率极低,缓存持续膨胀
Schema schema = SCHEMA_CACHE.getIfPresent(schemaKey);
if (schema == null) {
schema = parseSchema(metric);
SCHEMA_CACHE.put(schemaKey, schema);
}
return schema.validate(metric);
}
}
问题一目了然:V3版本中,schema ID变成了动态生成的UUID,导致缓存几乎无法命中,每次调用都会新解析一个schema并放入缓存,缓存持续增长直到OOM。
三、紧急止损
确认了根因,但线上还有800万条消息堆积着,不能等慢慢修复,得先止损。
3.1 临时扩容消费端
第一时间,我联系了运维团队,对消费集群进行水平扩容,从原来的10个节点扩展到20个节点,同时调整JVM参数,增加堆内存:
# 扩容后的JVM配置调整
-Xms8g -Xmx8g # 堆内存从4g增加到8g
-XX:MetaspaceSize=512m # 元空间相应调大
-XX:+UseG1GC # 切换到G1 GC,减少STW
-XX:MaxGCPauseMillis=200 # 目标GC暂停时间
-XX:InitiatingHeapOccupancyPercent=35 # 提前触发并发GC
扩容后,消费端的吞吐量从200条/秒提升到了1500条/秒,虽然还是远低于生产速度,但至少停止了堆积的进一步恶化。
3.2 关闭TS SDK的缓存功能
在消费端代码中,我临时添加了一个配置开关,绕过TS SDK的校验逻辑:
// 临时绕过TS SDK校验,直接写入
public void onMessage(MetricMessage message) {
MetricProto.Metric metric = MetricProto.Metric.parseFrom(message.getBody());
// 临时关闭TS校验(通过配置开关)
if (Boolean.parseBoolean(System.getenv("SKIP_TS_VALIDATE"))) {
businessDao.insert(metric);
return;
}
// 正常流程(带缓存校验)
TSDatabaseClient client = TSDatabaseClient.getInstance();
client.validateMetric(metric);
businessDao.insert(metric);
}
通过环境变量SKIP_TS_VALIDATE=true临时关闭了校验,消费端的内存压力骤减,吞吐量恢复到正常的4000条/秒左右。
3.3 消息补偿
在消费端恢复后,我写了一个补偿脚本,将积压的消息重新消费处理:
# 消息补偿脚本(Python)
from kafka import KafkaConsumer
import json
import time
consumer = KafkaConsumer(
'ts-metrics-topic',
bootstrap_servers=['broker1:9092', 'broker2:9092'],
group_id='ts-compensation-group',
auto_offset_reset='earliest',
consumer_timeout_ms=30000 # 30秒无新消息则退出
)
processed = 0
start_time = time.time()
for message in consumer:
# 重新投递到处理队列
process_and_store(message.value)
processed += 1
if processed % 10000 == 0:
elapsed = time.time() - start_time
rate = processed / elapsed
print(f"已处理 {processed} 条, 当前速度: {rate:.0f} 条/秒")
print(f"补偿完成,共处理 {processed} 条消息")
补偿脚本运行了约4个小时,将积压的800万条消息全部处理完毕。
四、彻底修复
止损只是治标,治本才是关键。
4.1 推动TS SDK修复
我将问题反馈给了TS系统的研发团队。他们的修复方案是在SDK内部增加一个schema去重机制:
// TS SDK 修复后的版本
public class TSDatabaseClient {
private static final Cache<String, Schema> SCHEMA_CACHE =
CacheBuilder.newBuilder()
.maximumSize(5000)
.expireAfterWrite(10, TimeUnit.MINUTES) // 增加过期时间
.removalListener(notification -> {
logger.info("Schema cache evicted: {}", notification.getKey());
})
.build();
// 新增:对动态schema ID进行归一化处理
public boolean validateMetric(Metric metric) {
// 提取schema的核心内容作为缓存key,而不是使用动态ID
String canonicalKey = buildCanonicalKey(metric);
Schema schema = SCHEMA_CACHE.getIfPresent(canonicalKey);
if (schema == null) {
schema = parseSchema(metric);
SCHEMA_CACHE.put(canonicalKey, schema);
}
return schema.validate(metric);
}
// 归一化key:忽略动态部分,只保留schema结构
private String buildCanonicalKey(Metric metric) {
return metric.getSchemaVersion() + ":" + metric.getMetricName();
}
}
同时,研发团队也在内部推动产品侧优化:V3版本的Protobuf消息中,schema定义应该改为引用ID而非内联完整schema,这样能从源头减少消息体积和缓存压力。
4.2 消费端代码加固
在消费端,我也做了几处加固:
第一,增加熔断保护:
// 增加熔断器,防止TS SDK异常时拖垮消费端
private final CircuitBreaker validateCircuitBreaker =
CircuitBreaker.ofDefaults("tsValidate");
public void onMessage(MetricMessage message) {
MetricProto.Metric metric = MetricProto.Metric.parseFrom(message.getBody());
try {
// 带熔断的保护调用
validateCircuitBreaker.run(() ->
TSDatabaseClient.getInstance().validateMetric(metric)
);
} catch (CircuitBreakerOpenException e) {
// 熔断器打开时,降级处理:跳过校验,直接写入
logger.warn("TS validate circuit breaker open, skip validation");
} catch (Exception e) {
logger.error("TS validation failed, skip and continue", e);
}
businessDao.insert(metric);
}
第二,增加消息积压监控和自动扩缩容:
# Kubernetes HPA配置
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: ts-metrics-consumer
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: ts-metrics-consumer
minReplicas: 10
maxReplicas: 50
metrics:
- type: Pods
pods:
metric:
name: kafka_lag
target:
type: AverageValue
averageValue: "1000" # 每个pod平均积压1000条时触发扩容
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 75
第三,调整消费策略,避免一次性加载过多消息:
// 调整消费配置,限制每次拉取的消息数量和大小
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "broker1:9092,broker2:9092");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "100"); // 每次最多拉100条
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "300000"); // 最大处理超时5分钟
props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, String.valueOf(50 * 1024 * 1024)); // 50MB
4.3 压测验证
修复完成后,我在测试环境进行了一轮压测:
# 使用kafka-producer-perf-test.sh发送压力测试
kafka-producer-perf-test.sh \
--topic ts-metrics-topic \
--throughput 10000 \
--record-size 1200 \
--num-records 1000000 \
--props producer.properties \
--max-block-ms 60000
压测结果显示:
| 指标 | 修复前 | 修复后 |
|---|---|---|
| 消费端内存占用 | 持续上涨至OOM | 稳定在3.2GB |
| GC频率 | Full GC每15分钟 | G1 GC正常调度 |
| 吞吐量 | 200条/秒(退化后) | 4500条/秒 |
| 消息延迟 | 3小时+ | 秒 |
五、复盘与总结
这次故障虽然最终解决了,但让我后怕的是,如果监控告警再慢半小时,或者消费端没有及时发现OOM,整个TS数据链路可能会彻底断裂,到时候补偿的成本就更高了。
复盘这次故障,有几个关键点值得记住:
1. 版本升级要谨慎
TS系统从Protobuf V2升级到V3,看似只是协议版本的变更,但实际上schema的定义方式发生了根本变化。这种变更如果没有充分评估对下游的影响,很容易埋下隐患。以后类似的升级,必须做全链路的影响评估,包括消息大小、处理性能、缓存策略等。
2. 监控要覆盖关键指标
这次能够及时发现异常,靠的是MQ积压告警。但更完善的监控应该包括:
- 消费端JVM内存和GC频率
- TS SDK缓存命中率
- 单条消息大小的分布
3. 要有降级和熔断机制
消费端不应该强依赖下游SDK的健康状态。如果TS SDK校验失败或超时,应该有优雅的降级策略,而不是让整个消费链路挂掉。
4. 压测要接近真实场景
如果在上次升级时做了充分的压测,发现消息体积变大、缓存命中率下降的问题,可能就不会有这次线上故障了。
最后说一句,做运维和开发这些年,我越来越觉得:系统故障不可怕,可怕的是同样的错误犯第二次。 这次排查虽然折腾了三个小时,但把根因彻底摸清楚了,后续也补上了防护措施。希望这篇记录能帮到遇到类似问题的朋友,也欢迎各位在评论区交流讨论。
