【问题标题】:Error in kafka consumer while writing data from input topic to output topic using kafka streams使用卡夫卡流将数据从输入主题写入输出主题时卡夫卡消费者出错
【发布时间】:2018-06-25 07:32:15
【问题描述】:

通过使用 kafka 连接器,我将 avro 格式的数据写入 kafka 主题,然后通过使用 kafka 流,我映射一些值并将输出写入其他主题:

Stream.to("output_topic");

我的数据正在写入输出主题,但我面临偏移问题。如果我的输入主题中有 25 条记录,它将所有 25 条记录写入我的输出主题但抛出错误:

[2018-06-25 12:42:50,243] ERROR [ConsumerFetcher consumerId=console-consumer-3500_kafka-connector-1529910768088-712e7106,leaderId=0, fetcherId=0]Error due to(kafka.consumer.ConsumerFetcherThread)

kafka.common.KafkaException: Error processing data for partition Stream-0 offset 25

这是我的完整错误:

> [2018-06-25 12:42:50,243] ERROR [ConsumerFetcher
> consumerId=console-consumer-3500_kafka-connector-1529910768088-712e7106,
> leaderId=0, fetcherId=0] Error due to
> (kafka.consumer.ConsumerFetcherThread) kafka.common.KafkaException:
> Error processing data for partition Stream-0 offset 25    at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2$$anonfun$apply$mcV$sp$1$$anonfun$apply$2.apply(AbstractFetcherThread.scala:204)
>   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2$$anonfun$apply$mcV$sp$1$$anonfun$apply$2.apply(AbstractFetcherThread.scala:169)
>   at scala.Option.foreach(Option.scala:257)   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2$$anonfun$apply$mcV$sp$1.apply(AbstractFetcherThread.scala:169)
>   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2$$anonfun$apply$mcV$sp$1.apply(AbstractFetcherThread.scala:166)
>   at scala.collection.Iterator$class.foreach(Iterator.scala:891)  at
> scala.collection.AbstractIterator.foreach(Iterator.scala:1334)    at
> scala.collection.IterableLike$class.foreach(IterableLike.scala:72)    at
> scala.collection.AbstractIterable.foreach(Iterable.scala:54)  at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2.apply$mcV$sp(AbstractFetcherThread.scala:166)
>   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2.apply(AbstractFetcherThread.scala:166)
>   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2.apply(AbstractFetcherThread.scala:166)
>   at kafka.utils.CoreUtils$.inLock(CoreUtils.scala:250)   at
> kafka.server.AbstractFetcherThread.processFetchRequest(AbstractFetcherThread.scala:164)
>   at
> kafka.server.AbstractFetcherThread.doWork(AbstractFetcherThread.scala:111)
>   at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:82)

> Caused by: java.lang.IllegalArgumentException: Illegal batch type
> class org.apache.kafka.common.record.DefaultRecordBatch. The older
> message format classes only support conversion from class
> org.apache.kafka.common.record.AbstractLegacyRecordBatch, which is
> used for magic v0 and v1  at
> kafka.message.MessageAndOffset$.fromRecordBatch(MessageAndOffset.scala:29)
>   at
> kafka.message.ByteBufferMessageSet$$anonfun$internalIterator$1.apply(ByteBufferMessageSet.scala:169)
>   at
> kafka.message.ByteBufferMessageSet$$anonfun$internalIterator$1.apply(ByteBufferMessageSet.scala:169)
>   at scala.collection.Iterator$$anon$11.next(Iterator.scala:410)  at
> scala.collection.Iterator$class.toStream(Iterator.scala:1320)     at
> scala.collection.AbstractIterator.toStream(Iterator.scala:1334)   at
> scala.collection.TraversableOnce$class.toSeq(TraversableOnce.scala:298)
>   at scala.collection.AbstractIterator.toSeq(Iterator.scala:1334)     at
> kafka.consumer.PartitionTopicInfo.enqueue(PartitionTopicInfo.scala:59)
>   at
> kafka.consumer.ConsumerFetcherThread.processPartitionData(ConsumerFetcherThread.scala:87)
>   at
> kafka.consumer.ConsumerFetcherThread.processPartitionData(ConsumerFetcherThread.scala:37)
>   at
> kafka.server.AbstractFetcherThread$$anonfun$processFetchRequest$2$$anonfun$apply$mcV$sp$1$$anonfun$apply$2.apply(AbstractFetcherThread.scala:183)
>   ... 15 more

【问题讨论】:

  • 听起来像是版本不匹配...您的代理、消息格式、Connect、Kafka Streams 版本是什么?
  • broker - confluent-4.1.0,消息格式 - Avro,connect - kafka-connect-jdbc,kafka Stream - 1.0.0-cp1
  • “消息格式”我不是指你的数据类型(那似乎是 Avro),而是 Kafka 消息格式:kafka.apache.org/documentation/#messageformat 有一个配置 message.format.versionlog.message.format.version。另外,您是否在某个时候升级了您的代理并使用升级前创建的主题?

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


【解决方案1】:

我在使用 kafka-consumer-console.sh 时遇到了同样的错误

问题在于--zookeeper 选项。 如果你给--zookeeper选项,老消费者默认启动,魔术选项将设置为默认v0或v1(当前kafka版本1.1使用v2) 这就是版本不匹配的原因。

你可以通过使用--bootstrap-server选项而不是--zookeeper来解决这个错误。(这意味着运行新版本的消费者)

当你给--bootstrap-server选项时,必须有broker的域(或ip)和端口号。 例如)—bootstrap-server kafka.domain:9092,kafka2.domain:9092

Broker(Kafka 服务器)默认端口为 9092,您可以在 kafka/config/server.properties 中更改端口。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-07-13
    • 2018-12-04
    • 1970-01-01
    • 2020-03-29
    • 2019-07-03
    • 2018-05-05
    相关资源
    最近更新 更多