【发布时间】: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