尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kafka消费者假活18小时:监控失守的根因与排查方案

Kafka消费者假活18小时:监控失守的根因与排查方案 1. 事故现场先看现象再谈原因1.1 凌晨的告警是一句订单同步是不是挂了凌晨 1 点多业务群里有人扔了一句订单同步是不是挂了我登录监控平台去看Kafka 集群的 CPU、内存、磁盘、网络吞吐全部正常消费者实例的负载也看不出任何异常。但业务方的反馈很直接整整 18 个小时核心订单 topic 的增量数据一条都没有同步下去。这是一件让人后脊发凉的事。故障持续了 18 小时监控大盘所有的指标都是绿的没有任何一条告警触发最终发现问题靠的是业务同学手动比对数据。也就是说我们的监控体系在最需要它的时候完全失明了。这个 case 非常适合拿出来复盘因为它涉及的并不是高深的分布式理论而是 Kafka 消费端很典型的假活状态进程在跑、连接在、消费组也在但业务处理就是没有推进。后面我会把排查思路、监控盲区、代码改造、参数调优全部展开给同样维护 Kafka 消费链路的同学一个可以直接照抄的排障方案。先说结论整个问题的链路是消费端在调用一个下游补数据接口时没有设置 read timeoutTCP 连接挂住收不到响应消费者的处理线程被占住后poll 循环也停了因为 max.poll.interval.ms 被调得很大心跳线程又一直正常broker 一直认为消费者还活着所以没有任何异常。这一套组合拳下来监控全绿就成了必然结果。1.2 快速确认问题范围是 broker 挂了还是消费端卡了遇到 Kafka 消费异常第一步不是看代码而是先确认问题到底在哪一层。是 producer 没发出去broker 吞了消息还是 consumer 拉不下来我习惯用三段式来判断。第一段看 broker 侧。到 Grafana 上看集群的 BytesInPerSec、BytesOutPerSec看 topic 的分区状态看 ISR 列表是否正常。这一般能排除 broker 层面的故障。我们当时看下来broker 端一切正常消息在生产侧是正常流入的topic 的 log end offset 一直在增长。第二段看消费端。用 Kafka 自带的命令行工具去看消费组的状态这是判断消费者到底有没有在消费最直接的方式bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092 --describe --group core-sync输出结果里会看到每个 partition 的 current-offset 和 log-end-offset。我们当时看到的情况是log-end-offset 已经涨到了几十万而 current-offset 停在了 18 小时前的位置。也就是说消费者一直在掉队lag 在持续增长。第三段确认消费组成员。使用--members --describe查看消费组里的活跃成员bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092 --describe --group core-sync --members如果成员数变成了 0说明消费者已经掉线broker 侧会触发 rebalance。如果我们当时看到的是成员还在、状态还是 Stable那就说明消费者进程活着但业务处理卡住了。这正是假活的典型特征。1.3 第一波排障CLI 三连与小工具定位确认了是消费端的问题后我连续做了三件事。第一看是不是有频繁的 rebalance。Kafka 的 rebalance 日志非常明显消费者端日志里会出现ConsumerRebalanceListener相关的输出或者Revoking previously assigned partitions这样的字样。如果 rebalance 频繁发生说明消费组的稳定性出了问题。我们当时看到了 rebalance 的日志但频率并不高大概每 30 分钟左右发生一次每次 rebalance 之后消费组又重新回到 Stable 状态。这个频率很隐蔽不会触发常见的rebalance 风暴告警。第二看有没有超时相关的异常。消费者日志里最常见的就是org.apache.kafka.clients.consumer.commit.KafkaCommitCoordinator抛出的Offset commit failed或者org.apache.kafka.clients.consumer.internals.ConsumerCoordinator的failed to poll within the configured max.poll.interval.ms。我们的日志里确实出现过后者的信息但因为是 WARN 级别而且不是每次都出现之前一直没引起重视。第三用可视化工具辅助看 lag。当时我们用了 kafka 系列的可视化工具比如 kafka-ui 这类开源项目界面上能看到每个 topic、每个 consumer group 的 lag 情况。工具本身没问题能看到 lag 在涨但这把双刃剑反而让我意识到一个更深层的问题这些工具能告诉你消费停了但告诉不了你为什么停。看完这三步问题范围已经锁定了消费者进程活着topic 在涨lag 在涨消费者偶尔触发 rebalance 但很快恢复。下一步就是抓线程栈看看处理线程到底卡在了哪里。2. 根因定位为什么 poll 正常、业务却停摆2.1 拉开线程栈找到真正的阻塞点定位 Java 进程卡死jstack 是永远的第一工具。先找到消费者进程的 PID然后连续抓几次线程栈jps -l # 找到消费者进程的 PID 以后 jstack PID stack_1.log sleep 10 jstack PID stack_2.log连续抓两次的意义在于对比如果线程栈里的位置在两次抓取之间完全没变那基本可以确认线程是真正卡住了。如果栈内容发生了变化可能是线程在正常跑只是处理得慢。我们当时在栈里看到了这样的调用链consumer-core-sync-1 #78 prio5 os_prio0 tid0x00007f8c58002800 nid0x4e23 runnable java.lang.Thread.State: RUNNABLE at java.net.SocketInputStream.socketRead0(Native Method) at java.net.SocketInputStream.socketRead(SocketInputStream.java:116) at java.net.SocketInputStream.read(SocketInputStream.java:171) at org.apache.http.impl.io.SessionInputBufferImpl.streamRead(SessionInputBufferImpl.java:137) at org.apache.http.impl.io.SessionInputBufferImpl.fillBuffer(SessionInputBufferImpl.java:153) at org.apache.http.impl.io.SessionInputBufferImpl.readLine(SessionInputBufferImpl.java:282) ... at com.example.sync.OrderSyncService.sync(OrderSyncService.java:88) at com.example.consumer.OrderEventListener.onMessage(OrderEventListener.java:45)线程停在socketRead0一个网络读操作上已经挂了很久。Stack 里的类名看得清清楚楚消费逻辑在处理订单同步时调用了一个外部 HTTP 接口这个调用卡住了。配合lsof看这个进程的网络连接lsof -p PID -i -n能找到一条发往下游补数服务的 TCP 连接状态是 ESTABLISHED。这就有意思了连接没有半开没有被 RST只是对端迟迟不返回数据。结合上下游的日志最终确认是下游服务响应超时而消费者侧使用的 HTTP 客户端没有配置 read timeout于是一直挂在那里等。2.2 阻塞点背后第三方接口超时与同步调用陷阱很多人看到这里会问一个 HTTP 接口超时怎么就能让 Kafka 消费者卡死 18 个小时这里的关键在于 Kafka 消费者默认的执行模型。在使用 spring-kafka 的KafkaListener时默认是单线程模型同一个线程既负责poll拉取消息又负责处理消息。如果处理消息时发生了阻塞poll 就不会发生。而 Kafka 的 client 有一个max.poll.interval.ms参数默认是 300 秒。也就是说如果两次 poll 之间的间隔超过了这个时间消费者会主动发起 LeaveGroup然后触发 rebalance重新分配分区。我们当时为了让消费者更稳定把这个参数调成了 1800000 毫秒也就是 30 分钟。为什么要调因为之前出现过下游服务偶发抖动导致消费者的单条消息处理超过了 max.poll.interval.ms被误踢出消费组触发了几次不必要的 rebalance。调大参数后这类误踢确实少了但副作用是灾难性的一旦处理线程真的阻塞消费者在 30 分钟内不会触发 rebalancebroker 侧会认为这个消费者仍然健康。配合心跳机制理解就更清晰了。Kafka 客户端有一个独立的心跳线程默认每隔heartbeat.interval.ms默认 3000 毫秒向 broker 发送心跳。心跳线程和业务处理线程是分离的只要 poll 间隔没有超过 max.poll.interval.ms心跳就会一直正常发送。broker 只看心跳心跳正常就认为消费者健康不会将分区重新分配给别人。所以业务线程卡死 18 小时broker 侧完全无感监控指标也毫无波动。这里有一个非常重要的认知心跳正常不等于消费正常。心跳只代表消费者进程还活着并不代表数据在推进。2.3 重新平衡 vs 卡死的区别心跳线程与处理线程我还想多说一句 rebalance 和卡死的区别因为这个概念搞不清楚排查方向就会偏。如果消费者进程直接宕机或者网络完全断开心跳会停session.timeout.ms默认 45000 毫秒到期后broker 会把消费者从消费组里移除触发 rebalance其他消费者会接管它的分区。这种情况在监控上表现为consumer group 的成员数下降lag 开始上涨但上涨的原因是因为没人消费了。这种问题相对好查。而卡死的表现完全不一样。消费者进程还在心跳线程还在向 broker 证明我活着消费组的状态是 Stable成员数没有变化。但处理线程卡在业务代码里poll 不再被调用offset 不再提交lag 只会持续上涨。这种状态从消费组层面看几乎没有任何异常迹象。我们当时的情况更隐蔽因为 max.poll.interval.ms 是 30 分钟所以每过 30 分钟消费者才会因为 poll 超时发一次 LeaveGroup触发一次 rebalance然后又重新加入消费组继续消费同一条卡住的消息再次阻塞。从日志看这个循环相当规律但如果不把日志的时间跨度拉到小时级别很难注意到这个 30 分钟一次的规律。3. 复盘监控失守问题出在指标没选对3.1 只监控 broker等于没监控消费链路事故处理完业务恢复了接下来才是真正有价值的工作搞清楚为什么监控全绿。我们当时的监控体系严格来说覆盖了两层。第一层是 Kafka broker 的 JMX 指标比如消息进出速率、请求处理时间、分区副本状态。第二层是消费者服务的 JVM 指标和机器指标比如 CPU、内存、GC、磁盘。这两层指标看起来已经很全了但它们有一个共同的盲区都没有回答一个核心问题——业务数据到底有没有在被消费broker 的吞吐指标只能说明消息有进有出不能说明消费端处理得怎么样。JVM 指标只能说明服务进程健康不能说明业务逻辑在推进。这就好比我们监控了一台冰箱的电压和压缩机功率却没有人看冰箱里的食物是不是真的变凉了。消费链路是否健康最终必须落到消费组层和业务层的数据新鲜度上。更关键的是当时消费者服务虽然是 Java 应用但我们的 Prometheus 采集端没有把 Kafka 客户端的 consumer 指标暴露出来。Spring Boot Actuator 默认暴露的是 JVM、HTTP、数据源这类指标Kafka 消费者的 lag、poll 次数、消费速率这些指标如果没有显式引入 micrometer 的 Kafka metrics binder是不会出现在监控面板上的。也就是说消费端在监控层面完全是一个黑盒。3.2 消费端三大关键指标lag、处理速率、最后处理时间那次事故之后我把消费链路监控分成了三个层级缺一不可。这里分享一张我自己整理的指标清单。层级指标说明消费组层consumer lag按 topic 和 partition最能直接反映消费是否跟得上的指标消费组层consumer group 成员数成员数下降基本意味着消费者掉线或 rebalance 异常消费组层rebalance 次数/频率短时间大量 rebalance 通常是代码或参数问题客户端层records consumed rate消费者每秒消费的消息条数客户端层单条消息处理耗时p99/p999反映业务处理是否异常变慢客户端层poll 调用间隔间接反映处理线程是否被阻塞业务层最后成功消费时间最有业务感知力的指标直接对 数据是否最新 负责这七个指标里前三个可以通过 kafka_exporter 轻松拿到中间三个需要在消费端代码里手动埋点最后一个是我个人认为最重要也最容易被忽略的指标。最后成功消费时间这个名字听起来不专业但它最实在。我给它起得直白一点last success consumed timestamp。每次消费并成功提交一条消息后用一个 gauge 记录下来。只要这个指标在持续更新说明消费链路是通的只要这个指标超过 N 分钟没更新不管消费组 lag 是 0 还是几千都说明业务处理已经停了。这个指标还有一个好处它非常容易设告警不需要复杂的判断逻辑只需要对比当前时间减去指标值是否超过阈值。3.3 告警阈值怎么定才不误报不漏报监控指标定了告警规则如果定得不好同样会失效。我们当时最怕的两种告警一种是噪音太大天天被开发同学屏蔽另一种是静默故障永远不触发。lag 绝对值告警是最容易踩坑的。如果业务有高峰期和低谷期直接用固定阈值比如 lag 5000 就告警很容易在业务高峰误报。因为高峰时段消费速率可能跟不上生产速率lag 短暂上涨是正常的只要最终能追平就没问题。更可靠的思路是两种策略叠加。第一lag 持续上升告警比如delta(kafka_consumergroup_lag_sum[15m]) 0意思是 lag 在 15 分钟内持续上涨配合for: 10m条件过滤瞬时抖动。第二绝对阈值告警只针对 lag 已经大到影响业务恢复时间 的场景比如 lag 超过 10 万条恢复需要几小时这种必须立刻通知。对于最后成功消费时间告警规则最简单- alert: KafkaConsumerLastProcessTimeTooOld expr: time() - kafka_consumer_last_process_time 300 for: 5m labels: severity: critical annotations: summary: 消费者组 {{ $labels.group }} 已超过 5 分钟没有成功处理消息5 分钟没有新消息就告警这个阈值在大多数订单同步场景下是合理的。如果你的业务本身就存在合法的空闲时间可以按业务低谷拉长阈值但哪怕是 30 分钟也比 18 小时没有人发现要好得多。4. 代码与参数双管齐下让消费端不再假活4.1 消费逻辑的工程化改造监控只能帮我们更快地发现问题真正让消费端不再假活还得靠代码和参数层面的硬性改造。先说代码。第一件事所有下游调用必须有显式超时。这次事故的直接原因是 HTTP 客户端没有配置 read timeout。很多人只配了 connect timeout忽略了 read timeout认为连上了就没事。实际上connect timeout 只管建立连接建立连接后服务端迟迟不返回数据read timeout 才会起作用。像这次的情况就是因为连接建立了但响应一直没回来所以 connect timeout 完全帮不上忙。JDK 11 的 HttpClient 可以这样配置HttpClient client HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(3)) .build(); HttpRequest request HttpRequest.newBuilder(uri) .timeout(Duration.ofSeconds(5)) .GET() .build();Request.timeout设置的是整个请求的时长上限包括连接、发送、等待响应比单独设置 read timeout 更保险。如果用的是 OkHttp也有对应的 connectTimeout、readTimeout、writeTimeout。RPC 调用同理Dubbo 有 timeout 参数gRPC 有 deadline都必须显式配置。第二件事不要在 Kafka 消费线程里做不可控的同步操作。如果业务确实需要调外部接口保守的做法是给这些操作设置严格超时让异常快速抛出而不是无限期挂住。更进一步的方案是把异步化引入消费链路比如把消息丢到一个线程池里处理poll 线程只管拉取消息。但这里有一个需要注意的坑引入线程池之后offset 的提交时机和消息的处理完成顺序会变得复杂处理不当会造成消息丢失或乱序。需要配合处理完成后手动提交 offset和下游业务幂等来兜底。第三件事消费逻辑里必须捕获 Throwable 并记录成告警。很多消费者代码只 catch Exception遇到 Error 级别的异常比如 OOM 前的 OutOfMemoryError线程直接死掉连日志都没有。建议在消费入口做兜底KafkaListener(topics order-events, groupId core-sync) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { try { OrderEvent event objectMapper.readValue(record.value(), OrderEvent.class); orderSyncService.sync(event); ack.acknowledge(); } catch (Exception e) { // 记录完整消息内容方便排查 log.error(handle message failed, topic{}, partition{}, offset{}, value{}, record.topic(), record.partition(), record.offset(), record.value(), e); // 这里应该根据业务来决定是重试、丢弃还是进入死信队列 } }注意这里的 catch 分支不能把异常吞掉然后继续 ack否则消息就静默丢失了。比较稳妥的做法是记录日志后把这条消息投递到死信 topic或者使用带重试和退避的机制。4.2 关键参数怎么调才不埋雷参数调优是这次事故的另一条主线。先说结论不要为了追求稳定而无脑调大 max.poll.interval.ms。我来把这个参数讲透。max.poll.interval.ms控制的是两次 poll 之间的最大时间间隔默认 300 秒。如果单条消息处理时间稳定在几秒内这个默认值是够用的。只有当业务确实存在某些耗时长、但又是必要的操作时才需要调大。但调大意味着什么意味着消费者的假活检查周期变长了一旦处理线程真正阻塞系统要等更长的时间才能通过超时机制触发 rebalance 来恢复。正确的调优思路是先优化处理逻辑让单条消息的处理时间可控再回头设置这个参数。比如业务上要求单条消息处理不超过 60 秒那就把 max.poll.interval.ms 设置为 120 秒或 180 秒留一定的余量。这样既不会因为瞬时抖动被踢出消费组也不会让问题拖太久才暴露。max.poll.records这个参数也很重要。它决定一次 poll 最多返回多少条消息。默认值是 500但如果单条消息处理逻辑较重500 条累加起来的处理时间很容易超过 max.poll.interval.ms。调小一些比如 200甚至 100对吞吐的影响其实不大但能显著降低处理超时被踢出组的概率。还有两个基础参数需要确认是否正确。session.timeout.ms是 broker 判断消费者是否存活的超时时间默认 45000 毫秒heartbeat.interval.ms是心跳发送间隔默认 3000 毫秒原则上是 session.timeout.ms 的三分之一左右。这两个参数保持默认通常就够了不建议动它们来解决问题。spring-kafka 的完整配置参考spring: kafka: bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092 consumer: group-id: core-sync enable-auto-commit: false auto-offset-reset: latest max-poll-records: 200 properties: max.poll.interval.ms: 120000 session.timeout.ms: 10000 heartbeat.interval.ms: 3000 listener: ack-mode: manual concurrency: 3这里的关键是把enable-auto-commit关掉换成手动 ack配合 ack-mode 来精确控制消息确认时机。concurrency设置的是监听容器的并发消费者数如果消费逻辑里有不可控的外部调用这个值不宜开得过大否则会同时占住多个线程去等外部接口反而加剧下游压力。4.3 制造一个卡死现场本地复现验证改造完成之后最好的验证方式不是直接上生产而是本地故意制造一个卡死场景确认监控能告警、代码能兜底、参数能触发恢复。这个过程本身也是一个很好的排障演练。我当时在本地搭了一个最小复现环境。用一个模拟的KafkaListener消息处理逻辑里故意调一个只会 sleep 不返回的假接口KafkaListener(topics test-slow-topic, groupId test-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { try { // 模拟外部接口卡死 Thread.sleep(600_000); ack.acknowledge(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }这个测试充分证明了前面说的两个知识点第一在单线程模型下消息处理 sleep 期间 poll 完全停止offset 不提交lag 上涨第二把 max.poll.interval.ms 设置为很小的值比如 5000 毫秒进行测试时消费者会因为超过 poll 间隔触发 rebalance日志里会出现The consumer has not polled within the configured max.poll.interval.ms的警告。而把参数调大后这个警告会完全消失消费者进入假活状态。测试完行为之后再验证本地是否有指标能看到这条链路。给项目加上 micrometer 的 Kafka 消费者指标启动后在 Prometheus 里能看到kafka_consumer_fetch_manager_records_consumed_total和kafka_consumer_coordinator_commit_total这些指标。当消费者卡住时这些指标的速率归零配合最后成功消费时间的埋点告警就能在几分钟内触发。5. 监控体系重构从指标好看到业务可用5.1 采集链路搭建kafka_exporter Prometheus Grafana监控体系的重构第一步是补齐消费组层面的指标采集。我们用了一个现成的开源组件kafka_exporter。它通过连接到 Kafka broker主动拉取 consumer group 的 offset、topic 的 log end offset、broker 基本信息并以标准格式暴露成 Prometheus 指标。安装部署非常简单一个二进制文件就是一个 exporter。指定 Kafka 地址后启动./kafka_exporter --kafka.serverkafka1:9092 --kafka.serverkafka2:9092 --kafka.serverkafka3:9092 --web.listen-address:9308启动之后在 Prometheus 的 scrape 配置里加一个 jobscrape_configs: - job_name: kafka-exporter static_configs: - targets: - 10.0.0.20:9308采集到的核心指标包括kafka_consumergroup_lag每个 consumer group 在某个 topic 某个 partition 上的 lagkafka_consumergroup_lag_sumconsumer group 的 lag 总和kafka_consumergroup_membersconsumer group 的活跃成员数有了这三个指标消费组层面的监控盲区就补上了。这套技术栈就是业内很流行的 Prometheus 监控部署方案也是 Kafka 监控最实用的组合。5.2 一条从卡死到发现的完整告警路径指标有了告警规则需要针对我们这次事故的重点来设计。我最终保留了两条核心规则一条覆盖lag 持续上涨的场景一条覆盖lag 绝对值过大的场景。groups: - name: kafka_consumer_alerts.rules rules: - alert: KafkaConsumerLagIncreasing expr: delta(kafka_consumergroup_lag_sum[15m]) 0 for: 10m labels: severity: warning annotations: summary: 消费组 {{ $labels.consumergroup }} 的 lag 在持续上涨 - alert: KafkaConsumerLagTooHigh expr: kafka_consumergroup_lag_sum 50000 for: 10m labels: severity: critical annotations: summary: 消费组 {{ $labels.consumergroup }} 的 lag 超过 50000第一条规则针对的是本次事故这种缓慢恶化的场景。即使 lag 的总量还没到危险值但只要它持续上涨超过 10 分钟就说明消费者跟不上生产速率值得关注。第二条规则针对的是已经严重堆积的场景一旦 lag 超过业务可容忍的恢复时间必须立即介入。但还有一个问题这两条规则都依赖 lag 指标。如果消费者进程直接挂了或者服务所在的机器无法上报指标这些规则会陷入无数据状态。Prometheus 默认对无数据的告警规则不做处理不会触发告警。这需要额外配置for之外的空数据检测。更稳妥的方案是利用最后成功消费时间的埋点指标在消费端代码里暴露一个 gaugeTimed(kafka_consumer_process_time) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { ... // 消息处理成功后更新 kafkaConsumerMetrics.setLastProcessTime(Instant.now().getEpochSecond()); }这条指标的意义在于它不只是反映消费有没有在推进而是反映业务消息有没有被真正处理掉。lag 上涨只说明消费跟不上但消费跟上了也不代表业务处理成功了。比如消费者把消息拉下来后直接丢弃lag 可能是 0但业务一样是停滞的。最后成功消费时间才能覆盖这种场景。5.3 可视化、巡检与事故报告告警是出问题时能通知人可视化和巡检是平时就能发现问题苗头。Grafana 面板我建议至少放三个视图。第一个是总览视图展示每个 consumer group 的 lag 总和用表格形式列出一眼能看到哪些 group 在堆积。第二个是趋势视图展示某个重点 group 的 lag 曲线和消费速率曲线观察两者之间的关系。第三个是健康视图展示 consumer group 的成员数、rebalance 次数、最后成功消费时间这些指标反映的是链路是否健康与 lag 反映的是否跟不上是两个维度。关于 Grafana 的使用有一个比较实用的功能是生成 PDF 监控报告。Grafana 官方提供了 Image Renderer 插件配合 Grafana Reporting 或者简单的定时任务可以把 Dashboard 渲染成图片或 PDF 定期发送到邮件。这个功能在事故复盘时特别有用。我在这次事故的复盘会上直接导出了过去 24 小时的监控面板 PDF把故障期间监控全绿这个事实摆到所有人面前比口头解释一百遍都有效。巡检方面除了告警规则还可以写一个简单的定时脚本对每个 consumer group 做一次健康检查。检查消费组是否还存在检查最后成功消费时间是否在正常范围内检查成员数是否符合预期。如果业务上允许强烈建议把消费者断流演练做成定期事项。每个月挑一个低峰期人为把一个消费者的处理线程 suspend 掉验证告警能不能在 10 分钟内触发。有些监控体系看着指标齐全真到了故障时才发现告警路由没人接收、通知渠道失效这种问题只有在演练中才能暴露。6. 排障工具箱与三个习惯6.1 排障命令速查每次排 Kafka 消费问题我都会从下面这张表开始。这里整理成速查表方便直接复制使用。场景命令作用查看消费组 lagbin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --describe --group core-sync展示每个 partition 的 current-offset、log-end-offset 和 lag查看消费组成员上面的命令加--members --describe展示活跃消费者实例、分配的分区查看消费组状态上面的命令加--state展示消费组是 Stable、PreparingRebalance 还是 Empty抓线程栈jstack PID stack.log分析 Java 进程卡点连续抓两次对比查看 GC 情况jstat -gcutil PID 1000 20观察 GC 频率排除 GC 长时间停顿导致的问题查看进程网络连接lsof -p PID -i -n查看该进程建立了哪些 TCP 连接查看 Kafka 服务端日志关注kafka.log.LogCleaner和kafka.server.ReplicaFetcherThread排除 broker 侧异常排障的时候有一个顺序原则要记住先消费组、再线程、再外部依赖。不要一上来就翻代码先用命令行确认问题到底出在哪一层再决定是否看代码。6.2 jstack 典型迹象解读jstack 的输出对新手来说很容易看晕但只要抓住几个典型的栈顶特征就能快速判断卡点的类型。如果栈顶在sun.nio.ch.SocketDispatcher.read或者java.net.SocketInputStream.socketRead0说明这个线程阻塞在 Socket 读操作上最常见的场景就是调用外部 HTTP 接口或 RPC 接口对端没有及时返回。这是我们这次事故的典型特征。如果栈顶在java.util.concurrent.ThreadPoolExecutor.getTask或java.util.concurrent.LinkedBlockingQueue.take说明线程在等待任务通常不是问题。真正的问题是线程池里的核心线程全部被占满新提交的任务全部排队这种情况往往伴随业务超时和任务堆积。如果栈顶在java.lang.Thread.sleep则需要结合代码上下文判断。如果是人为的重试退避问题不大如果是一种死循环式的 sleep那就要警惕了。连续抓两次线程栈是一个很有用的技巧两次栈内容完全一致说明线程是真的卡住了栈内容在变化可能只是负载高、处理慢。我当时抓了三次每次间隔 10 秒卡在 socketRead0 上的线程纹丝不动这基本就锁定了问题。6.3 这次事故留给我的三个习惯复盘到最后我想分享这次事故给我留下的三个习惯。这三个习惯不是技术方案但它们的价值可能比任何一项技术改造都大。第一个习惯任何消费者服务上线前必须确认最后成功消费时间这个业务级指标已经接入监控。这是一个很朴素的判断标准不管 Kafka 内部状态多么正常只要业务数据没有在推进系统就是有问题的。宁可少几个花哨的中间件指标也要守住这条业务底线。第二个习惯所有外部同步调用必须显式设置超时时间。connect timeout 只是保护了连不上的场景read timeout 才保护连上了但不回应的场景。在消费者链路里任何不可控的同步操作都可能变成消费卡死的源头。我现在写代码看到第三方调用没有设置超时会直接打回去。第三个习惯每半年做一次断流演练。人为制造一次消费者停顿验证告警是否真的能响通知是否能到达责任人。监控体系是越用越可靠的如果只在故障时才发现监控失效代价太大了。我们的监控经过这次重构后做法很简单每个季度挑一个低峰期暂停一个非核心消费者组确认 10 分钟内告警能触发然后恢复消费。演练的成本很低但带来的安全感是实打实的。那次事故之后我再看到监控面板全绿第一反应不再是安心而是反问自己绿色代表的是什么是进程活着还是数据在推进这两个概念之间的距离可能就是一个 18 小时的静默故障。
返回列表