【问题标题】:Can't Consume Messages When Using Kafka v.0.10.0.x使用 Kafka v.0.10.0.x 时无法消费消息
【发布时间】:2019-04-18 01:52:03
【问题描述】:

我在集群上使用 kafka v.0.10.2。

我可以使用 v.0.8.x 和 v0.10.2 很好地生成消息

但是在使用客户端 v0.10.0.x 使用消息时,我遇到以下错误;

警告 [ConsumerFetcherThread-console-consumer-myconsumer-0-1002],获取 kafka.consumer.ConsumerFetcherThread$FetchRequest@16090d7a 时出错。可能原因:java.nio.BufferUnderflowException(kafka.consumer.ConsumerFetcherThread)

好的,现在我的 kafka.clien 是 v.0.8.x 但我有一个新问题

    6 15:07:13 WARN scheduler.TaskSetManager: Lost task 0.0 in stage 0.0 (TID 0, hadoop11, executor 10): org.apache.spark.SparkException: Task failed while writing rows
        at org.apache.spark.internal.io.SparkHadoopWriter$.org$apache$spark$internal$io$SparkHadoopWriter$$executeTask(SparkHadoopWriter.scala:151)
        at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$3.apply(SparkHadoopWriter.scala:79)
        at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$3.apply(SparkHadoopWriter.scala:78)
        at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
        at org.apache.spark.scheduler.Task.run(Task.scala:109)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
Caused by: kafka.common.KafkaException: String exceeds the maximum size of 32767.
        at kafka.api.ApiUtils$.shortStringLength(ApiUtils.scala:73)
        at kafka.api.TopicData$.headerSize(FetchResponse.scala:107)
        at kafka.api.TopicData.<init>(FetchResponse.scala:113)
        at kafka.api.TopicData$.readFrom(FetchResponse.scala:103)
        at kafka.api.FetchResponse$$anonfun$4.apply(FetchResponse.scala:170)
        at kafka.api.FetchResponse$$anonfun$4.apply(FetchResponse.scala:169)
        at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
        at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
        at scala.collection.immutable.Range.foreach(Range.scala:160)
        at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
        at scala.collection.AbstractTraversable.flatMap(Traversable.scala:104)
        at kafka.api.FetchResponse$.readFrom(FetchResponse.scala:169)
        at kafka.consumer.SimpleConsumer.fetch(SimpleConsumer.scala:135)
        at org.apache.spark.streaming.kafka.KafkaRDD$KafkaRDDIterator.fetchBatch(KafkaRDD.scala:196)
        at org.apache.spark.streaming.kafka.KafkaRDD$KafkaRDDIterator.getNext(KafkaRDD.scala:212)
        at org.apache.spark.util.NextIterator.hasNext(NextIterator.scala:73)
        at scala.collection.Iterator$GroupedIterator.fill(Iterator.scala:1126)
        at scala.collection.Iterator$GroupedIterator.hasNext(Iterator.scala:1132)
        at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
        at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
        at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$4.apply(SparkHadoopWriter.scala:124)
        at org.apache.spark.internal.io.SparkHadoopWriter$$anonfun$4.apply(SparkHadoopWriter.scala:123)
        at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1411)
        at org.apache.spark.internal.io.SparkHadoopWriter$.org$apache$spark$internal$io$SparkHadoopWriter$$executeTask(SparkHadoopWriter.scala:135)
        ... 8 more

我做什么节目 字符串超过了 32767 的最大大小。

【问题讨论】:

  • 请显示所有使用的配置、代码和命令的minimal reproducible example
  • 我很确定最好的情况是使用与代理版本相同的客户端版本
  • 是的,我同意你的观点,但关键是版本不一致。
  • kafka.apache.org/0102/documentation.html 表示版本 0.10.2 的代理支持 0.8.x 和更新的客户端。但是您遇到的错误通常是由不匹配的版本引起的。
  • 字符串超过了最大大小32767,我该怎么办

标签: apache-kafka


【解决方案1】:

告诉我我的最终计划 我将版本升级到 0.10.2

【讨论】:

    猜你喜欢
    • 2019-09-23
    • 2018-11-29
    • 2020-02-22
    • 2018-08-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-11
    相关资源
    最近更新 更多