【问题标题】:Issue with Kafka Broker - UnknownServerExceptionKafka 代理问题 - UnknownServerException
【发布时间】:2019-07-06 23:38:41
【问题描述】:

我们的应用程序使用 springBootVersion = 2.0.4.RELEASE 以及 compile('io.projectreactor.kafka:reactor-kafka:1.0.1.RELEASE') 依赖项。

我们拥有的 Kafka Broker 版本为 1.0.1

当我们通过创建reactor.kafka.sender.SenderRecord 将消息发送到 Kafka 时间歇性地发送消息,并在寻找 reactor.kafka.sender.SenderResult.exception() 时响应 Kafka

java.lang.RuntimeException: org.apache.kafka.common.errors.UnknownServerException: The server experienced an unexpected error when processing the request 填充在异常中。

重试几次后,消息成功通过。

在代理日志中,以下错误被多次打印,没有任何堆栈跟踪

[2019-02-08 15:43:07,501] ERROR [ReplicaManager broker=3] Error processing append operation on partition price-promotions-local-event-0 (kafka.server.ReplicaManager)

price-promotions-local-event 是我们的主题。

我在网上查看过,但没有明确的解决方案或方法来分类此问题,非常感谢您提供的任何帮助。

【问题讨论】:

  • 你有没有机会使用 snappy?
  • @AsierAranbarri,是的,reactor-kafka:1.0.1.RELEASE 的编译时间依赖于 snappy-java:1.1.4,请参见下文 -------------------- -------------------------------------------------- -------------------------- +--- io.projectreactor.kafka:reactor-kafka:1.0.1.RELEASE | +--- io.projectreactor:reactor-core:3.1.8.RELEASE (*) | \--- org.apache.kafka:kafka-clients:1.0.2 | +--- org.lz4:lz4-java:1.4 | +--- org.xerial.snappy:snappy-java:1.1.4 | \--- org.slf4j:slf4j-api:1.7.25
  • @AsierAranbarri 据我了解,Kafka 0.10.0 遭受了“处理分区上的附加操作时出错”的问题,这是由于 snappy-java 中的一个错误导致解析 MAGIC HEADER 处理不正确。 snappy-java-1.1.2.6 已发布以解决此问题。 reactor-kafka 的最新版本是 1.1.0.RELEASE,它依赖于 kafka-clients 2.0.0,而后者又依赖于 snappy-java 1.1.7.1,因此升级我们的 reactor-kafka 依赖项可以解决问题吗?还是会有其他方法?

标签: java spring spring-boot apache-kafka reactive-kafka


【解决方案1】:

在进一步调查中,我们可以将代理日志上的堆栈跟踪作为

ERROR [ReplicaManager broker=1] Error processing append operation on partition price-promotions-local-event-0 (kafka.server.ReplicaManager)
java.lang.IllegalArgumentException: Magic v1 does not support record headers
    at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:403)
    at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:442)
    at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:595)
    at kafka.log.LogValidator$.$anonfun$convertAndAssignOffsetsNonCompressed$2(LogValidator.scala:138)
    at kafka.log.LogValidator$.$anonfun$convertAndAssignOffsetsNonCompressed$2$adapted(LogValidator.scala:136)
    at scala.collection.Iterator.foreach(Iterator.scala:929)
    at scala.collection.Iterator.foreach$(Iterator.scala:929)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1417)
    at scala.collection.IterableLike.foreach(IterableLike.scala:71)
    at scala.collection.IterableLike.foreach$(IterableLike.scala:70)
    at scala.collection.AbstractIterable.foreach(Iterable.scala:54)
    at kafka.log.LogValidator$.$anonfun$convertAndAssignOffsetsNonCompressed$1(LogValidator.scala:136)
    at kafka.log.LogValidator$.$anonfun$convertAndAssignOffsetsNonCompressed$1$adapted(LogValidator.scala:133)
    at scala.collection.Iterator.foreach(Iterator.scala:929)
    at scala.collection.Iterator.foreach$(Iterator.scala:929)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1417)
    at scala.collection.IterableLike.foreach(IterableLike.scala:71)
    at scala.collection.IterableLike.foreach$(IterableLike.scala:70)
    at scala.collection.AbstractIterable.foreach(Iterable.scala:54)
    at kafka.log.LogValidator$.convertAndAssignOffsetsNonCompressed(LogValidator.scala:133)
    at kafka.log.LogValidator$.validateMessagesAndAssignOffsets(LogValidator.scala:64)
    at kafka.log.Log.liftedTree1$1(Log.scala:654)
    at kafka.log.Log.$anonfun$append$2(Log.scala:642)
    at kafka.log.Log.maybeHandleIOException(Log.scala:1669)
    at kafka.log.Log.append(Log.scala:624)
    at kafka.log.Log.appendAsLeader(Log.scala:597)
    at kafka.cluster.Partition.$anonfun$appendRecordsToLeader$1(Partition.scala:499)
    at kafka.utils.CoreUtils$.inLock(CoreUtils.scala:217)
    at kafka.utils.CoreUtils$.inReadLock(CoreUtils.scala:223)
    at kafka.cluster.Partition.appendRecordsToLeader(Partition.scala:487)
    at kafka.server.ReplicaManager.$anonfun$appendToLocalLog$2(ReplicaManager.scala:724)
    at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:234)
    at scala.collection.mutable.HashMap.$anonfun$foreach$1(HashMap.scala:138)
    at scala.collection.mutable.HashTable.foreachEntry(HashTable.scala:236)
    at scala.collection.mutable.HashTable.foreachEntry$(HashTable.scala:229)
    at scala.collection.mutable.HashMap.foreachEntry(HashMap.scala:40)
    at scala.collection.mutable.HashMap.foreach(HashMap.scala:138)
    at scala.collection.TraversableLike.map(TraversableLike.scala:234)
    at scala.collection.TraversableLike.map$(TraversableLike.scala:227)
    at scala.collection.AbstractTraversable.map(Traversable.scala:104)
    at kafka.server.ReplicaManager.appendToLocalLog(ReplicaManager.scala:708)
    at kafka.server.ReplicaManager.appendRecords(ReplicaManager.scala:459)
    at kafka.server.KafkaApis.handleProduceRequest(KafkaApis.scala:465)
    at kafka.server.KafkaApis.handle(KafkaApis.scala:98)
    at kafka.server.KafkaRequestHandler.run(KafkaRequestHandler.scala:65)
    at java.lang.Thread.run(Thread.java:748)

org.apache.kafka:kafka-clients:1.0.2 中可用的类文件MemoryRecordsBuilder 中,我们有下面的IllegalArgumentException 被抛出的地方。


if (magic < RecordBatch.MAGIC_VALUE_V2 && headers != null && headers.length > 0)
  throw new IllegalArgumentException("Magic v" + magic + " does not support record headers");

因此,ProducerRecord 中设置了导致问题的标头,在打印 ProducerRecord 时,我们发现标头是由 AppDynamics 添加的——“singularityheader”被添加到 Kafka Produced 记录中。

c.t.a.p.i.m.i.KafkaProducerInterceptor   : The kafka Interceptor ProducerRecord header:: RecordHeader(key = singularityheader, value = [110, 111, 116, 120, 100, 101, 116, 101, 99, 116, 61, 116, 114, 117, 101, 42, 99, 116, 114, 108, 103, 117, 105, 100, 61, 49, 53, 53, 49, 51, 55, 51, 54, 57, 49, 42, 97, 112, 112, 73, 100, 61, 55, 49, 48, 51, 50, 42, 110, 111, 100, 101, 105, 100, 61, 49, 51, 53, 55, 53, 51, 53])

更多阅读https://developer.ibm.com/messaging/2018/07/10/additional-rfh-header-added-appdynamics-monitor-agent-tool/

所以我们在拦截器中明确地将标头设置为空,这已经解决了问题。

【讨论】:

    猜你喜欢
    • 2019-04-22
    • 2019-07-10
    • 2021-08-12
    • 2020-02-12
    • 2017-05-18
    • 1970-01-01
    • 2017-07-09
    • 2020-02-20
    • 1970-01-01
    相关资源
    最近更新 更多