【问题标题】:How to best handle SerializationException from KafkaConsumer poll method如何最好地处理来自 KafkaConsumer 轮询方法的 SerializationException
【发布时间】:2019-12-11 22:39:29
【问题描述】:

在 Kafka 的消费者轮询循环中,当 poll 方法抛出 SerializationException 时,有没有办法跳过这条消息(又名“毒丸”)并继续消费主题中的下一个事件?

我可以捕获异常并使用 consumer.seek() 方法将偏移量移动到下一条消息,但该方法需要分区和偏移量作为输入参数。有没有办法获得这些值?

我在 github 存储库中有示例代码。运行示例:

$ git clone https://github.com/bjornhjelle/kafka-streams-examples-gradle.git
$ cd kafka-streams-examples-gradle
$ ./gradlew build -x test
$ ./gradlew test --tests no.test.SerializationExceptionExample

该示例向 Kafka 生成三个事件。第二个事件导致 SerializationException。捕获并记录异常。在这一点上,我想将偏移量移过这个事件。而是在轮询循环中再次抛出。因此不消耗第三个事件,因此测试失败。

我知道关于同一主题的这个未解决的问题,但它提到了 Kafka 客户端版本 https://issues.apache.org/jira/browse/KAFKA-4740

我也知道我可能可以通过使用 Kafka Streams 和在那里处理毒丸的新功能来解决这个问题 (KIP-161: streams deserialization exception handlers)

首先让我对此进行调查的是这个异常: (示例代码会导致不同的 SerializationException,因为我无法重新创建这个)

Exception in thread "SimpleAsyncTaskExecutor-3"         org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition apcitxpt-1 at offset 339798. If needed, please seek past the record to continue consumption.
Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id -1
Caused by: java.nio.BufferUnderflowException
at java.nio.Buffer.nextGetIndex(Buffer.java:500)
at java.nio.HeapByteBuffer.get(HeapByteBuffer.java:135)
at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.getByteBuffer(AbstractKafkaAvroDeserializer.java:77)
at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:119)
at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserialize(AbstractKafkaAvroDeserializer.java:93)
at io.confluent.kafka.serializers.KafkaAvroDeserializer.deserialize(KafkaAvroDeserializer.java:55)
at org.apache.kafka.common.serialization.ExtendedDeserializer$Wrapper.deserialize(ExtendedDeserializer.java:65)
at org.apache.kafka.common.serialization.ExtendedDeserializer$Wrapper.deserialize(ExtendedDeserializer.java:55)
at org.apache.kafka.clients.consumer.internals.Fetcher.parseRecord(Fetcher.java:923)
at org.apache.kafka.clients.consumer.internals.Fetcher.access$2600(Fetcher.java:93)
at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.fetchRecords(Fetcher.java:1100)
at org.apache.kafka.clients.consumer.internals.Fetcher$PartitionRecords.access$1200(Fetcher.java:949)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:570)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:531)
at org.apache.kafka.clients.consumer.KafkaConsumer.pollOnce(KafkaConsumer.java:1170)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1103)
at no.ruter.nextgen.kafkaConsumerRunners.ApcInitRunner.run(ApcInitRunner.java:63)
at java.lang.Thread.run(Thread.java:745)

【问题讨论】:

  • 如果您使用消费者组并捕获异常,当您继续轮询时消费者不会继续前进吗?
  • 我使用消费者组,但我看不出这有什么帮助。一个分区只分配给消费者组中的一个消费者,因此错误消息仍会阻止从该分区消费更多消息,直到我能够以某种方式将偏移量移动到当前偏移量之后。如果我因为未捕获异常而让进程失败,则该分区将被重新分配给组中的另一个消费者,然后该消费者将收到与错误消息相同的问题。 (如果我错了,请纠正我......)
  • > “直到我能够以某种方式移动偏移量”。如果你捕捉到异常并且消费者没有崩溃,它就不会继续读取和处理?对于消费者组,消费者“移动偏移量”。如果你也抓住了,那不会发生吗?
  • 不,如果我捕捉到异常,则不会提交偏移量,因为异常来自轮询方法,因此消息没有被正确使用。如果我没有捕捉到异常,但让应用程序退出,同样的事情。消息一直准备好再次被使用,直到它过期,因为达到了保留期。因此,在此之前,该分区上的所有其他消费都被“阻止”。
  • 我明白了。抱歉,我从来没有处理过这个问题。是的,您应该能够手动搜索/提交吗?正如您在问题中所说,这需要您知道偏移量/分区,您可以随着消费者的进展跟踪它们。当您遇到异常时,您会知道最后一条消息是好的。

标签: java apache-kafka kafka-consumer-api


【解决方案1】:

我找到的解决方案是解析序列化异常消息以获取所需的数据。主题名称、分区和偏移量:

catch (SerializationException se) {
                String s = se.getMessage().split("Error deserializing key/value for partition ")[1].split(". If needed, please seek past the record to continue consumption.")[0];
                String topic = s.split("-")[0];
                int offset = parseInt(s.split("offset ")[1]);
                int partition = parseInt(s.split("-")[1].split(" at")[0]);

                TopicPartition topicPartition = new TopicPartition(topic, partition);
                logger.debug("Skipping {}-{} offset {}", topic, partition , offset);
                consumer.seek(topicPartition, offset + 1L);}

【讨论】:

  • 这可行,但不是一种可靠的方法,因为它取决于来自异常的实际消息。因此,您可以在这里看到更好的选择:link 在反序列化器本身中捕获异常,然后在 consumerRecords 循环中相应地处理结果。
  • 不推荐这样做。它可能非常脆弱,例如如果主题名称中有“-”。
猜你喜欢
  • 2014-02-15
  • 1970-01-01
  • 2020-10-03
  • 1970-01-01
  • 2021-08-10
  • 2013-06-13
  • 2014-01-28
  • 1970-01-01
  • 2021-02-03
相关资源
最近更新 更多