【问题标题】:kafka broker failed to handle request due to heap OOM由于堆 OOM,kafka 代理无法处理请求
【发布时间】:2018-09-02 22:54:45
【问题描述】:

由于 1.0.0 存在 OOM 问题,我已更新到 1.0.1。

我设置了有四个代理的集群。

大约有 150 个主题,总共大约 4000 个分区,ReplicationFactor 为 2。
连接器用于向代理写入/读取数据。
connecotr 版本为 0.10.1。
平均消息大小为 500B,每秒大约 60000 条消息。
经纪人之一保持报告OOM,并且无法处理以下请求:

[2018-03-24 12:37:17,449] 错误 [KafkaApi-1001] 处理请求时出错 {replica_id=-1,max_wait_time=500,min_bytes=1,topics=[{topic=voltetraffica.data,partitions=[ {partition=16,fetch_offset=51198,max_bytes=60728640} ,{partition=12,fetch_offset=50984,max_bytes=60728640}]}]} (kafka.server.KafkaApis) java.lang.OutOfMemoryError:Java 堆空间 在 java.nio.HeapByteBuffer.(HeapByteBuffer.java:57) 在 java.nio.ByteBuffer.allocate(ByteBuffer.java:335) 在 org.apache.kafka.common.record.AbstractRecords.downConvert(AbstractRecords.java:101) 在 org.apache.kafka.common.record.FileRecords.downConvert(FileRecords.java:253) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1$$anonfun$apply$4.apply(KafkaApis.scala:525) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1$$anonfun$apply$4.apply(KafkaApis.scala:523) 在 scala.Option.map(Option.scala:146) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1.apply(KafkaApis.scala:523) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1.apply(KafkaApis.scala:513) 在 scala.Option.flatMap(Option.scala:171) 在 kafka.server.KafkaApis.kafka$server$KafkaApis$$convertedPartitionData$1(KafkaApis.scala:513) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$createResponse$2$1.apply(KafkaApis.scala:561) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$createResponse$2$1.apply(KafkaApis.scala:560) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:891) 在 scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1334) 在 scala.collection.IterableLike$class.foreach(Ite​​rableLike.scala:72) 在 scala.collection.AbstractIterable.foreach(Ite​​rable.scala:54) 在 kafka.server.KafkaApis.kafka$server$KafkaApis$$createResponse$2(KafkaApis.scala:560) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$fetchResponseCallback$1$1.apply(KafkaApis.scala:574) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$fetchResponseCallback$1$1.apply(KafkaApis.scala:574) 在 kafka.server.KafkaApis$$anonfun$sendResponseMaybeThrottle$1.apply$mcVI$sp(KafkaApis.scala:2041) 在 kafka.server.ClientRequestQuotaManager.maybeRecordAndThrottle(ClientRequestQuotaManager.scala:54) 在 kafka.server.KafkaApis.sendResponseMaybeThrottle(KafkaApis.scala:2040) 在 kafka.server.KafkaApis.kafka$server$KafkaApis$$fetchResponseCallback$1(KafkaApis.scala:574) 在 kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$processResponseCallback$1$1.apply$mcVI$sp(KafkaApis.scala:593) 在 kafka.server.ClientQuotaManager.maybeRecordAndThrottle(ClientQuotaManager.scala:176) 在 kafka.server.KafkaApis.kafka$server$KafkaApis$$processResponseCallback$1(KafkaApis.scala:592) 在 kafka.server.KafkaApis$$anonfun$handleFetchRequest$4.apply(KafkaApis.scala:609) 在 kafka.server.KafkaApis$$anonfun$handleFetchRequest$4.apply(KafkaApis.scala:609) 在 kafka.server.ReplicaManager.fetchMessages(ReplicaManager.scala:820) 在 kafka.server.KafkaApis.handleFetchRequest(KafkaApis.scala:601) 在 kafka.server.KafkaApis.handle(KafkaApis.scala:99)

然后大量缩减 ISR(这个经纪人是 1001)

018-03-24 13:43:00,285] INFO [分区 gnup.source.offset.storage.topic-5 broker=1001] 将 ISR 从 1001,1002 缩小到 1001 (kafka.cluster.Partition)
018-03-24 13:43:00,286] 信息 [分区 s1mme.data-72 代理 = 1001] 将 ISR 从 1001,1002 缩小到 1001 (kafka.cluster.Partition)
018-03-24 13:43:00,286] INFO [分区 gnup.sink.status.storage.topic-17 broker=1001] 将 ISR 从 1001,1002 缩小到 1001 (kafka.cluster.Partition)
018-03-24 13:43:00,287] INFO [Partition probessgsniups.sink.offset.storage.topic-4 broker=1001] 将 ISR 从 1001,1002 缩小到 1001 (kafka.cluster.Partition)
018-03-24 13:43:01,447] INFO [GroupCoordinator 1001]:稳定的组连接-VOICE_1_SINK_CONN 第 26 代 (__consumer_offsets-18) (kafka.coordinator.group.GroupCoordinator)
我每次运行时都无法转储堆:
[root@sslave1 kafka]# jcmd 55409 GC.heap_dump /home/ngdb/heap_dump175.hprof
55409: com.sun.tools.attach.AttachNotSupportedException:无法打开套接字文件:目标进程没有响应或 HotSpot VM 未加载 在 sun.tools.attach.LinuxVirtualMachine.(LinuxVirtualMachine.java:106) 在 sun.tools.attach.LinuxAttachProvider.attachVirtualMachine(LinuxAttachProvider.java:63) 在 com.sun.tools.attach.VirtualMachine.attach(VirtualMachine.java:208) 在 sun.tools.jcmd.JCmd.executeCommandForPid(JCmd.java:147) 在 sun.tools.jcmd.JCmd.main(JCmd.java:131)

JVM参数为:

-XX:+ExplicitGCInvokesConcurrent -XX:GCLogFileSize=104857600 -XX:InitialHeapSize=2147483648 -XX:InitiatingHeapOccupancyPercent=35 -XX:+ManagementServer -XX:MaxGCPauseMillis=20 -XX:MaxHeapSize=4294967296 -XX:NumberOfGCLogFiles=10 -XX:+ PrintGC -XX:+PrintGCDateStamps -XX:+PrintGCDetails -XX:+PrintGCTimeStamps -XX:+UseCompressedClassPointers -XX:+UseCompressedOops -XX:+UseG1GC -XX:+UseGCLogFileRotation

当我使用 -XX:mx=2G 时,四个经纪人报告 OOM,
我把它创建到4G后,只有一个经纪人报告了OOM。
Ticker 也在https://issues.apache.org/jira/browse/KAFKA-6709 中提出。

【问题讨论】:

  • 当我在0.10.2.2下测试时,同样的吞吐量,堆大小为-XX:mx=1G,-XX:ms=1G,没有报告OOM。

标签: java apache-kafka


【解决方案1】:

在 0.10.X 和 >= 0.11.X Kafka 版本之间,消息格式发生了变化。

因此,当使用较旧的客户端 (= 0.11) 时,代理必须在将消息发送回客户端之前进行下转换。这记录在升级说明中:http://kafka.apache.org/documentation/#upgrade_11_message_format

您可以在堆栈跟踪中看到这确实发生了:

at java.nio.HeapByteBuffer.(HeapByteBuffer.java:57)
at java.nio.ByteBuffer.allocate(ByteBuffer.java:335)
at org.apache.kafka.common.record.AbstractRecords.downConvert(AbstractRecords.java:101)
at org.apache.kafka.common.record.FileRecords.downConvert(FileRecords.java:253)

这会带来性能损失,并且还会增加所需的内存量,因为代理需要分配新的缓冲区来创建向下转换的消息。

您应该尝试将您的客户端升级到与代理相同的版本。还要考虑到您当前的堆有多小(4GB),增加它可能会有所帮助。

另一种选择是强制较新的代理使用较旧的消息格式(使用log.message.format.version),但这会阻止您使用一些较新的功能。

【讨论】:

  • 谢谢。我会试试的。知道为什么只有一个经纪人有这样的问题吗?它可能是由整个爆裂引起的吗?还是代理之间的流量不平衡?
【解决方案2】:

我遇到了同样的问题,我通过重启kafka进程解决了这个问题。

【讨论】:

    猜你喜欢
    • 2018-12-27
    • 2014-01-31
    • 2021-07-26
    • 2013-06-13
    • 1970-01-01
    • 2021-12-18
    • 2012-08-31
    • 2017-11-27
    • 2020-11-11
    相关资源
    最近更新 更多