【问题标题】:Error reading field 'topic_metadata': Error reading array of size 1139567, only 45 bytes available读取字段“topic_metadata”时出错:读取大小为 1139567 的数组时出错,只有 45 个字节可用
【发布时间】:2016-10-02 04:57:05
【问题描述】:

--消费者

Properties props = new Properties();
        String groupId = "consumer-tutorial-group";
        List<String> topics = Arrays.asList("consumer-tutorial");
        props.put("bootstrap.servers", "192.168.1.75:9092");
        props.put("group.id", groupId);
        props.put("enable.auto.commit", "true");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(props);
        try {
            consumer.subscribe(topics);
            while (true) {

                ConsumerRecords<String, String> records = consumer.poll(Long.MAX_VALUE);
                for (ConsumerRecord<String, String> record : records)
                    System.out.printf("offset = %d, key = %s, value = %s", record.offset(), record.key(), record.value());


            }
        } catch (Exception e) {
            System.out.println(e.toString());
        } finally {
            consumer.close();
        }
    }

我正在尝试编写运行上面的代码,它是一个简单的消费者代码,它试图从一个主题中读取,但我遇到了一个奇怪的异常,我无法处理它。

org.apache.kafka.common.protocol.types.SchemaException: Error reading field 'topic_metadata': Error reading array of size 1139567, only 45 bytes available

我引用你也是我的生产者代码

--生产者

Properties props = new Properties();
        props.put("bootstrap.servers", "192.168.1.7:9092");
        props.put("acks", "all");
        props.put("retries", 0);
        props.put("batch.size", 16384);
        props.put("linger.ms", 1);
        props.put("buffer.memory", 33554432);
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

        Producer<String, String> producer = new KafkaProducer<String, String>(props);
        for(int i = 0; i < 100; i++)
            producer.send(new ProducerRecord<String, String>("consumer-tutorial", Integer.toString(i), Integer.toString(i)));

        producer.close();

这里是 kafka 配置

--启动zookeeper

bin/zookeeper-server-start.sh config/zookeeper.properties

--启动Kafka服务器

bin/kafka-server-start.sh config/server.properties

-- 创建话题

bin/kafka-topics.sh --create --topic consumer-tutorial --replication-factor 1 --partitions 3 --zookeeper 192.168.1.75:2181

--Kafka 0.10.0

<dependency>
           <groupId>org.apache.kafka</groupId>
           <artifactId>kafka-clients</artifactId>
           <version>0.10.0.0</version>
   </dependency>
   <dependency>
           <groupId>org.apache.kafka</groupId>
           <artifactId>kafka_2.11</artifactId>
           <version>0.10.0.0</version>
   </dependency>

【问题讨论】:

  • 您使用的是什么客户端版本?
  • 我正在使用 Kafka 0.10.0
  • 哎呀,我忘了问你经纪人了!也在使用 kafka 0.10 吗?我遇到了同样的错误,因为 kafka 客户端 0.10 与代理 0.9 不兼容。
  • 我降级到客户端 0.9 并删除了异常但代码仍然无法正常工作
  • 尝试在 0.9 上使用代理和在 0.10 上使用客户端时遇到了类似的问题。我升级了我的融合,一切都解决了。

标签: apache-kafka kafka-consumer-api


【解决方案1】:

在使用版本为 0.10.0.0 的 kafka_2.11 工件时,我也遇到了同样的问题。但是,一旦我将 kafka 服务器更改为 0.10.0.0,这个问题就得到了解决。早些时候我指的是0.9.0.1。看起来服务器和你的 pom 版本应该是同步的。

【讨论】:

    【解决方案2】:

    我通过降级到 kafka 0.9.0 解决了我的问题,但这对我来说仍然不是一个有效的解决方案。如果有人知道如何在 kafka 0.10.0 版本中解决此问题的有效方法,请随时发布。在那之前这是我的解决方案

    <dependency>
               <groupId>org.apache.kafka</groupId>
               <artifactId>kafka-clients</artifactId>
               <version>0.9.0.0</version>
       </dependency>
       <dependency>
               <groupId>org.apache.kafka</groupId>
               <artifactId>kafka_2.11</artifactId>
               <version>0.9.0.0</version>
       </dependency>
    

    【讨论】:

      【解决方案3】:

      我有同样的问题。客户端 jar 兼容性问题,因为我使用的是 Kafka 服务器 9.0.0 和 Kafka 客户端 10.0.0。基本上 Kafka 0.10.0 引入了一种新的消息格式,并且无法从旧版本中读取主题元数据版本。

      <dependency>
       <groupId>org.springframework.kafka</groupId>
       <artifactId>spring-kafka</artifactId>
       <version>1.0.0.RELEASE</version> <!-- changed due lower version of the kafka server -->
      </dependency>
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-09-26
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多