ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka消息堆积实战排查:从lag监控到参数调优与架构设计

Kafka消息堆积实战排查:从lag监控到参数调优与架构设计 Kafka消息堆积这事干过大数据的兄弟应该都懂那是真能把人逼疯的。平时跑得好好的集群某天突然发现某个topic的lag蹭蹭往上飙消费者那边肉眼可见地延迟上游数据还在疯狂往里灌你这边要么加机器要么改代码手忙脚乱一顿操作。我这些年在大数据一线摸爬滚打踩过无数堆积的坑今天把实战里积累的排查方法、调优参数和架构设计思路一次性讲透希望能帮你少走几步弯路。这篇文章不是什么理论教材就是我自己处理线上故障的记录和经验总结适合正在和Kafka消息堆积做斗争的朋友无论是刚入门的新手还是已经有几年经验的老手都可以在这里面找到一些用得上的东西。1. 消息堆积的本质与核心影响1.1 什么是消息堆积它从哪来Kafka的消息堆积用一句大白话说就是生产者的写入速度大于消费者的处理速度消息在broker端越积越多。正常情况下topic里的消息被消费之后就可以被删除或压缩掉但如果消费端长期跟不上生产端的节奏新的消息不断进来旧的消息又一直没被消费日志段文件就会一直膨胀消费者位点和最新位点之间的距离也就是lag就会越来越大。这里的核心机制在于Kafka的分区模型。每个topic被拆成多个partition每个partition是一个有序的消息日志消费者通过维护自己的offset来记录已经处理到哪条消息。生产者往分区里写消息时是顺序追加速度极快但消费者是一个个拉取、逐条处理的一旦下游处理逻辑变重、数据库响应变慢、或者某个接口超时整个消费链路就会被拖住。结果就是生产端10条/秒进来消费端只处理得了3条/秒积压那是必然的。还有一个很容易被忽略的堆积来源就是消费者组重平衡。当一个消费者实例挂掉、新实例加入、或者分区分配发生变化时整个consumer group会暂停消费进入rebalance过程。这个过程如果频繁发生比如session超时时间设得太短、处理逻辑偶尔卡顿导致心跳发送延迟那么消费者实际上有一大半时间都在“暂停中”消费速率自然骤降消息堆积也就此产生。1.2 堆积不处理会带来什么后果消息堆积不是一个“等一会儿就好了”的问题它是实实在在的事故源头。最直接的后果就是数据延迟。在实时数仓、风控、推荐、监控告警这类场景里数据晚到一分钟都可能带来完全不同的结果。比如做实时风控的用户刷了一笔可疑交易结果这条消息在Kafka里积压了十分钟才被消费到钱早就出去了。这是业务事故。其次堆积会对Kafka集群本身造成存储压力。每个partition的消息都有自己的保留期限和大小上限消息一直消费不掉日志段文件就不会被清除磁盘空间会被大量占用。尤其是那些开了无限保留或者保留周期很长的topic一旦堆积磁盘可能几天之内被打满broker直接宕机。我见过一个集群因为堆积导致磁盘使用率达到95%触发了broker的自动下线策略整个集群进入不健康状态最后全部业务都受到牵连。还有更隐形的问题就是消费端内存溢出。Kafka消费者拉取消息时会先把消息放到本地内存缓冲区如果堆积严重、拉取到的消息量很大而下游处理又慢下一个pull循环又开始了内存里的消息越积越多最终触发OOM。我处理过一次线上消费者OOM的事故当时就是堆积导致消费者一次拉取了大量消息处理不过来又没来得及释放JVM直接堆溢出消费者进程崩溃lag再次飙升形成恶性循环。2. 快速定位堆积从监控到命令行的一线排查方法2.1 一套实用的监控指标体系排查消息堆积第一步不是去看业务日志而是先看监控大盘。一套成熟的Kafka监控体系至少要覆盖几个核心指标消息生产速率、消息消费速率、消费者lag、broker磁盘使用率、网络IO和CPU使用率。这些指标可以从Kafka内置的JMX指标里拿配合Prometheus和Grafana做可视化也可以直接用Kafka生态里现成的工具比如LinkedIn开源的Burrow。Burrow是我个人非常推荐的一个lag监控工具。和其他监控方案不同Burrow不是简单地拿到lag然后画条线就完事它会通过评估消费行为来判断消费者组的健康状态把lag值映射成OK、WARNING、ERROR三个等级然后通过HTTP接口暴露给外部报警系统。我之前有个项目就是用Burrow做核心topic的堆积监控配合Prometheus告警规则一旦某个消费者组的lag连续5分钟超过阈值就触发钉钉机器人报警效果非常好。具体部署方式不复杂Burrow本身是Go写的直接编译好扔到服务器上跑就行。除了监控工具还要在代码层面埋点。生产端的Send回调、消费端的poll处理耗时和offset提交成功率这些都能在代码里记成日志或者指标。有了这些数据你才能回答一个最关键的问题消费速率降低到底是网络慢、下游慢还是消费者本身的处理逻辑变慢了2.2 用命令行工具确认堆积点的实操记录监控告警只能告诉你“出事了”真正定位问题还得靠命令行直接看数据。最常用的就是kafka-consumer-groups脚本这是Kafka自带的工具语法如下# 查看所有消费组 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list # 查看指定消费组的消费进度和lag kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my_consumer_group执行之后会看到一张表里面列出了这个消费者组消费的每个topic、每个分区以及CURRENT-OFFSET当前消费到哪个位置、LOG-END-OFFSET该分区最新消息位置和LAG积压了多少条。这张表就是一线排查的核心依据。比如我看到某个topic有20个分区其中15个分区的lag都是0只有5个分区lag堆积了几百万那基本可以判定数据分布不均匀某个分区的数据量太大或者对应处理逻辑卡住了。我还经常用另一个命令来确认生产端有没有异常写入# 查看topic的详细情况包括分区数、副本分布等 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic my_topic # 查看topic的起始和最新offset kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list localhost:9092 --topic my_topic --time -1这些命令执行的时候要注意--bootstrap-server参数一定要配上正确的broker地址多个broker用逗号分隔。而且现在Kafka 2.x以后很多老命令已经被整合进kafka-consumer-groups了执行时如果提示找不到命令大概率是版本变了去libs目录下找对应的JAR包跑kafka-run-class.sh就行。2.3 判断堆积发生在“入口”还是“出口”排查堆积一定要先搞清楚积压的位置到底在入口还是出口。入口积压指的是生产端写入Kafka的速率本身就慢或者broker处理生产请求的能力下降导致客户端生产超时重试数据根本没完整落进集群。这种情况在监控上看是生产速率指标下降、broker端CPU或磁盘IO飙高消费者lag不一定会涨因为消息压根没进来多少。出口积压才是我们最常说的消息堆积也就是生产端正常甚至超量写入但消费者处理不过来。这时候生产速率曲线是平的甚至向上的消费者的处理速率曲线却是下降的lag曲线呈现快速爬升的形态。绝大多数堆积属于这一类。怎么快速区分这两种情况我惯用的套路是先看消费者的总lag趋势。如果所有分区lag都在同步上涨那是系统性的比如消费者组整体变慢、线程阻塞、下游依赖挂了如果只有部分分区lag上涨那就是局部问题比如某个key的消息特别多导致某个分区数据量异常、某台消费者所在的机器负载太高、或者某个消费者的处理线程卡死。然后去看生产速率和消费速率两条曲线确认是不是消费速率下降导致的lag上升如果是再去下游服务日志里找原因如果不是就得往上游查。3. 影响消息积压的核心因素与调优方向3.1 生产者端的限速与batch机制很多人在处理堆积时把注意力全放在消费者身上其实生产端也会成为堆积的根源。Kafka生产者在发送消息时两个batch参数对吞吐量的影响是决定性的batch.size和linger.ms。batch.size默认16KB也就是说生产者会尽量把多条消息攒成一个batch再发送减少网络往返次数linger.ms默认0即立即发送。如果业务场景对延迟不敏感、只重吞吐把linger.ms调到5到10msbatch.size适当调大到32KB或64KB单条消息发送请求的条数会成倍增加broker端处理生产请求的吞吐量也能明显提升。还有一个经常被忽略的坑是生产者的max.in.flight.requests.per.connection参数。默认值是5意思是同一个连接上最多可以有5个未确认的请求在途。当这个值大于1并且启用了重试时消息的发送顺序可能被打乱。但反过来如果把值设成1吞吐量会下降不少。我的建议是如果你的业务对消息顺序要求没那么严保持默认的5就行没必要为了微乎其微的顺序问题牺牲大量吞吐如果严格要求有序那就老老实实把max.in.flight设为1并关闭重试或者通过给消息加sequence字段在消费端做排序。生产端的限流也很重要。很多大数据项目在上游接的是埋点数据或日志数据流量天然有峰值比如业务高峰期、定时任务触发点、大促秒杀时刻。这种场景下如果不加流控瞬时大量请求打到Kafka上broker压力剧增反而会拖慢整个集群的处理能力。适当的做法是在生产端做本地缓冲或限速配合Kafka本身的背压机制让写入曲线尽量平滑。3.2 消费者端的消费能力瓶颈消费端的消费能力是整个链路里最容易出问题的环节。消费者拉取消息后要做反序列化、业务处理、结果写库每一步都可能成为瓶颈。最常见的坑有这几个单条消息处理耗时过长比如消费一条消息需要调用一个外部HTTP接口接口响应时间200毫秒那一个消费者线程一秒最多处理5条消息100个分区100万条积压算一下就知道了要处理非常久。消费者线程模型设计不合理默认的KafkaConsumer是单线程模型一个消费者实例里只跑一个poll循环即使你的服务器有32核CPU如果只靠单线程处理那CPU大部分都闲着消费速率自然是上不去的。反序列化和业务处理耦合在一起有些项目把数据的解密、校验、格式转换全堆在消费线程里做导致poll线程长时间阻塞心跳发送超时触发了消费者组重平衡性能进一步恶化。我之前接手过一个实时计算项目消费延迟一直在涨看代码发现消费者在poll之后直接在一个循环里做同步的SQL写入每条消息都要等数据库返回结果延迟当然压不住。后来改造的思路很简单poll出来的消息放到一个内存队列里由独立的线程池去异步处理poll线程保持轻快消费速率一下就上来了。这个思路后面会详细展开。3.3 分区数量与消费者数量的匹配关系这里必须把分区数、消费者数、消费并发度这三者的关系掰开揉碎讲清楚。Kafka的设计里有一个铁律同一个消费者组内一个分区只能被一个消费者实例消费。换句话说如果topic有10个分区你的消费者组里最多只有10个消费者能同时消费加了第11个也是闲着。反过来如果你只起了1个消费者那10个分区的消息就全压在这一个消费者身上并发度完全起不来。所以当你发现堆积是因为消费并发度不够时要做的事情很明确增加topic的分区数或者增加消费者实例数。但两者有本质区别。增加消费者数不改变数据分布逻辑只是让更多消费者平分现有分区增加分区数才能让每个现有消费者拿到更多的并行单位。举一个实际例子。之前我们有个订单topic一开始只有6个分区消费者组里也是6个实例处理得稳稳的。后来业务量上来了lag开始增长我们先是把消费者实例加到12个结果发现只有6个在干活另外6个完全闲置。原因就是分区数只有6个没法分配给更多消费者。后来把topic的分区数扩容到了24个让每个消费者处理4个分区问题才彻底解决。这里有个坑要提醒topic的分区数只能增加不能减少所以刚开始建topic时要结合未来两年的业务量预估来定分区数别省着用。4. 从消费端发力解决消息堆积的实战参数与代码调整4.1 调整fetch参数与大页块拉取解决了“多消费者并发”的问题接下来就是提升消费者单次拉取消息的效率。Kafka消费者每次poll会调用底层的fetch请求有几个参数直接影响拉取的行为fetch.min.bytesbroker返回给消费者一次拉取的最小字节数默认1字节。设大之后消费者会攒够一定数据量才返回减少请求次数提升吞吐。fetch.max.wait.msbroker最多等待多长时间再响应消费者请求默认500毫秒。吞吐优先的场景下可以把这个值调到1000毫秒甚至更多让broker多攒点数据再一次性返回。max.partition.fetch.bytes每个分区在单次fetch请求中最多返回的字节数默认1MB。如果你的消息体比较大比如每条几百KB这个值可能需要调大否则一次fetch返回不了几条消息浪费网络往返。我常用的“高吞吐消费端”配置组合大致是这样的Properties props new Properties(); props.put(bootstrap.servers, broker1:9092,broker2:9092); props.put(group.id, order_consumer_group); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); // 每次fetch至少拉取2MB数据 props.put(fetch.min.bytes, 2097152); // 最多等待800毫秒再返回 props.put(fetch.max.wait.ms, 800); // 每个分区最多拉取2MB props.put(max.partition.fetch.bytes, 20971520); // 自动提交关闭手动管理offset props.put(enable.auto.commit, false); props.put(auto.offset.reset, earliest); // 允许消费者自动创建topic生产环境建议关掉 props.put(allow.auto.create.topics, false);不要小看这几个参数的调整。在一个生产集群上我把fetch.min.bytes从默认值调到2MB之后消费者端的总体拉取次数减少了一半以上CPU和网络开销都有明显下降。对延迟要求不高的场景这套配置基本可以无脑用。4.2 关闭自动提交后的位移管理很多堆积问题的背后其实还藏着“重复消费”的隐患。Kafka消费者的offset提交方式有两种自动提交和手动提交。自动提交enable.auto.committrue默认每隔5秒提交一次消费位移但这个“自动”有个致命问题如果你poll了一批消息正在处理的时候进程挂了而自动提交的间隔又到了offset可能已经被提交了于是这批消息没处理完就丢了反过来如果处理完还没到提交间隔进程挂了重启后又会重复消费一批消息。在堆积场景下处理速度慢、处理时间长自动提交这种粗粒度的机制非常不可靠。我生产环境里一律关闭自动提交改成手动提交推荐两种模式。一种是每处理完一批消息就同步提交while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { process(record); // 真正地处理业务逻辑 } // 同步提交当前拉取的位移 consumer.commitSync(); }这种做法的优点是实现简单、提交可靠不会丢消息缺点是同步提交会阻塞poll循环如果处理慢、提交频繁会对吞吐有一定影响。另一种是异步提交加回调监控while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { process(record); } consumer.commitAsync((offsets, exception) - { if (exception ! null) { log.error(commit failed, offsets: {}, offsets, exception); } }); }异步提交不阻塞主循环性能更好但提交失败时不会自动重试需要在回调里做补偿。我一般建议消费逻辑简单、追求吞吐的场景用异步提交消费逻辑复杂、丢消息会出大事故的场景用同步提交。还有一种比较取巧的方式在进程退出时先调用commitSync()兜底确保退出前位移完整提交一遍。4.3 消费者线程模型改造配置参数只能优化到一定程度如果你的消费逻辑确实重单线程的poll模型瓶颈始终在这里。这里我得说一个很多教程都不愿意讲透的点KafkaConsumer不是线程安全的不能简单地用一个consumer实例启动多线程去并发poll。这是很多新手改造线程模型时踩过的坑。有两条路可以走。一条是一个consumer实例配一个线程池poll线程负责拉取消息然后把消息路由给线程池去处理poll线程和worker线程之间用内存队列解耦。核心代码如下ExecutorService workerPool Executors.newFixedThreadPool(20); BlockingQueueConsumerRecordString, String queue new LinkedBlockingQueue(10000); // poll线程 new Thread(() - { while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { queue.put(record); } // 手动提交为poll到的最大offset做异步准备 } }).start(); // worker线程 for (int i 0; i 20; i) { workerPool.submit(() - { while (running) { ConsumerRecordString, String record queue.poll(100, TimeUnit.MILLISECONDS); if (record ! null) { process(record); } } }); }这种模型的优点是改造简单、性能提升明显缺点是无法保证分区内消息的顺序因为同一个分区的消息可能被多个worker线程并发处理。如果业务要求分区级有序可以在路由时用record.partition()对worker数取模保证同一个分区的消息进同一个线程。另一种方式是一个分区一个消费者实例topic有多少个分区就启动多少个消费者线程每个消费者消费一个独立分区。这种方案隔离性好分区内有序天然保证但实例数太多线程和连接的开销也不小适合分区数不多、单分区数据量很大的场景。5. 从架构侧化解堆积削峰填谷与背压设计5.1 队列拆分与业务隔离如果你的系统里不同重要程度的消息混在同一个topic里那么一旦发生堆积轻则全部受到影响。一个典型的错误是把日志消息、业务消息、监控消息全部塞进一个topic高峰期日志量猛增直接把业务消息给“淹”了所有消费者都被日志消息拖住真正要紧的订单数据却处理不过来。我的建议是设计阶段就做好队列拆分。按照业务线或优先级建多个topic比如order_topic、log_topic、notification_topic每个topic独立设置分区数、副本数、保留策略。不同topic的消费者相互隔离一个topic堆积不影响另一个。注意这里说的“隔离”不只是逻辑上的独立broker物理层面最好也做隔离。比如关键业务topic放在专用broker上日志topic放到另一个broker上这样可以避免某个topic的突发流量打满整个集群的磁盘IO拖累其他业务。架构上还应该做“分类治理”高优业务用同步发送、低延迟消费者非核心业务用异步批量、大batch消费。这就像交通系统里的公交专用道高峰期公交车优先通行不会被私家车堵在路上整体的交通效率反而更高。5.2 临时扩容加消费者与加分区当堆积已经发生快速恢复业务是最重要的这时候最有效的操作就是扩容。扩容有两个方向横向加消费者实例增加topic分区。但就像前面讲的如果消费者组里已经有等于分区数的消费者实例了继续加消费者不会带来任何提升。这时候就必须扩容分区。也有一些细节需要提醒。在Kafka 2.2及之前版本中增加分区的操作本身可能引发rebalance影响正在进行的消费。如果你有严格的SLA要求建议在业务低峰期做分区扩容或者提前规划好分区数。另外扩容分区后消费者组的分区分配会自动重新计算旧消费者会尽量保持已有分区的分配新分区会分配到当前负载较低的消费者上。如果你的消费者代码里硬编码了“分区数必须等于某个值”的逻辑扩容后可能会产生Bug代码里不要对分区数做这种假设。实际扩容过程中还有一类特殊的“临时扩容”场景消费者进程需要处理的消息量太大即使加了足够多的消费者也因为下游数据库、接口的极限支撑不住而无法提速。这种情况下你可以让消费端暂时把消息落盘到本地文件或临时表等高峰期过了再补处理而不是硬扛实时消费。思路就是“先存下来保住不丢再慢慢消化”。5.3 持久化兜底与降级策略Kafka本身的机制是消息默认保留7天这个保留期内的消息是不会被物理删除的。所以堆积发生时消息本身并没有丢消费者有机会慢慢追上来。但如果堆积时间超过了保留期最早的消息就会被自动删除那才是真正的数据丢失事故。为了兜底可以在消费端设计一套持久化重放机制如果Kafka中的消息因为堆积太久被删除你还能从下游的备份存储里恢复数据。这个兜底思路说白了就是“双写”或者“旁路备份”。对于一些核心的、要保证不丢的数据消费端在正常处理之外同步把原始消息写入HBase、ES或OSS后续如果需要做历史数据回放、数据校验、补偿处理直接读备份存储不依赖Kafka本身的数据保留。这个方案会带来额外的存储成本但对核心数据链路来说这点成本是值的。降级策略同样重要。当堆积已经发生、消费者短时间内追不上时不要死磕实时链路而是主动降级。比如做实时数仓的项目如果Kafka堆积严重下游实时任务可以先切换到离线批处理模式用Spark或Flink读Kafka里积压的数据做批量清洗入库等lag降下来了再切回实时模式。这种流批一体的切换方案是我们治理堆积问题时的一个核心思路。Kafka天然的“可回放”特性为这种切换提供了很好的基础——消费者重置offset到最早位置就能把积压消息重新拉一遍。6. 常见问题与排查技巧实录6.1 一查到堆积就急着加消费者先查分区数这是我在工作中见过最多的一类问题。很多团队一发现consumer lag上升第一反应就是“加机器、加消费者实例”结果加到和分区数一样多之后发现lag还在涨才意识到分区数已经是瓶颈了。这个问题的排查顺序应该是先去kafka-consumer-groups --describe里看分区数和消费者实例数。如果消费者数已经大于等于分区数就不要再加消费者了要考虑扩容分区或者优化单条消费速度。还有一个小坑很多人以为消费者实例数越多越好其实Kafka里同一个consumer group内每个消费者至少负责一个分区消费者数量超过分区数时多出来的消费者会处于空闲状态不仅没有帮助还白白占用资源、增加rebalance时的心跳开销。所以消费者实例数量一般建议等于分区数的70%到100%之间留出一点扩展空间又不会浪费资源。6.2 堆积偶尔出现但很快恢复有一种堆积是很“诡异”的它不持续只会隔一段时间冒出来一次lag冲到几百上千过了一会儿又跌回0。这种瞬时堆积往往不是消费能力问题而是消费端触发了某些慢操作。比如消费线程每处理1000条消息就调一次外部刷新token的接口恰巧这个接口超时了10秒或者某个定时任务在整点执行抢占了数据库连接池导致消费端的SQL写入全部排队。处理这类问题关键是把“慢”的那一部分定位出来。我一般习惯在消费代码里做耗时打点用Metrics统计poll方法的耗时、process方法里每个子步骤的耗时再配合日志里的耗时分布基本能在几次堆积发生后找到规律。找到规律后解决手段就很多了把慢操作异步化、预加载、连接池隔离、熔断降级等等。6.3 关于“消息重复消费”与堆积的纠缠消息重复和消息堆积像是两个冤家经常结伴出现。最典型的情况是消费者处理消息超时了消费者组判定它失联并触发重平衡将分区分配给另一个消费者新消费者从上次提交的offset重新消费于是同一条消息被处理了两遍。如果业务对重复数据很敏感这就会造成下游数据的重复写入而重复处理又会占用消费端资源进一步加剧堆积。处理重复消费的标准思路是消费幂等。技术上可以在消息体里加一个唯一消息ID消费时先用Redis的setNx做去重已处理过的ID直接跳过也可以在下游数据库建唯一索引靠数据库天然的主键约束把重复数据拒掉。对于实时计算引擎像Flink和Spark Structured Streaming提供了exactly-once语义配合Kafka source自带的事务和offset管理可以在很大程度上解决重复消费问题。总之在堆积治理的过程中一定要把重复消费的策略一并考虑进去否则堆积追平之后面对的是一堆重复数据那才叫真正的灾难。6.4 一套稳用的监控预警组合前面提到过Burrow这里再展开讲一下我习惯的监控报警组合。监控数据的采集可以用JMX_exporter配合Prometheus来实现每个broker和每个消费者进程都暴露metrics端点然后在Prometheus里配置采集规则。lag数据的获取我除了用kafka-consumer-groups手工查也会用Burrow的/v3/kafka/consumer_group接口定期拉取把它转换成Prometheus指标。告警规则方面我的经验是不要用固定的lag阈值因为不同业务量级差异很大有的topic平时lag就是个位数有的GPU算力任务天天lag一万也不算事。更好的办法是设置“lag持续增长时长”和“lag超过最小offset delta”的双重条件。比如某个消费者组的lag在过去10分钟内持续增长且当前值大于5000才触发告警。这样能过滤掉大量瞬时波动把真正的问题暴露出来。如果你连Burrow都不想折腾图省事的话也有简单方案写一个定时脚本每5分钟跑一次kafka-consumer-groups --describe把结果发送到监控系统然后在Grafana里画出发送速率、消费速率和lag三条曲线也能覆盖大部分排查需求。关键是数据要有积累出现堆积时才能对比出是谁先“异常”的。写在最后的一点体会说实话消息堆积这个问题90%的根因都不是Kafka本身不行而是业务发展太快技术方案没跟上。很多项目在初期设计topic分区数和消费者并发度时往往只按当时的流量拍脑袋等业务量翻了几倍之后各种问题就集中爆发了。我的经验是建topic之初分区数尽量按未来一到两年流量的峰值来规划每个消费者的单分区处理能力按比较保守的值来估算比如单分区每秒200条简单SQL处理宁可多分几个分区把并发度预留出来也不要等到堆积了再去扩容。另外监控一定要前置告警一定要有人响应堆积不可怕最可怕的是堆积发生了没人发现。希望这篇文章里的排查思路和参数调整能帮你在下次“lag了”的时候冷静应对快速止血也祝你的大数据链路越来越稳少踩几个坑。
返回列表