【问题标题】:Sending messages to Kafka with header using KafkaProducer使用 KafkaProducer 将消息发送到带有标头的 Kafka
【发布时间】:2018-11-23 09:10:38
【问题描述】:

我正在尝试创建一个简单的实用程序来将消息发布到 Kafka,但需要与消息一起传递标头。 该实用程序在没有标头的情况下工作正常,但在尝试发送标头时出现错误。

下面是我正在使用的示例代码 -

public static void main(String[] args) throws Exception {

        Properties props = new Properties();

        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "host:port");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        KafkaProducer<String, String> producer = new KafkaProducer<>(props);

        List<Header> headers = Arrays.asList(new RecordHeader("sample_header", "sample_value".getBytes()));
        ProducerRecord<String, String> record = new ProducerRecord<>("TEST", 0, "key", "sample message", headers);
        Future<RecordMetadata> future = producer.send(record);
        System.out.println(future.get());

        producer.close();

    }

我得到的例外是 -

Exception in thread "main" java.lang.IllegalArgumentException: Magic v1 does not support record headers
    at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:385)
    at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:424)
    at org.apache.kafka.common.record.MemoryRecordsBuilder.append(MemoryRecordsBuilder.java:481)
    at org.apache.kafka.common.record.MemoryRecordsBuilder.append(MemoryRecordsBuilder.java:504)
    at org.apache.kafka.clients.producer.internals.ProducerBatch.tryAppend(ProducerBatch.java:106)
    at org.apache.kafka.clients.producer.internals.RecordAccumulator.append(RecordAccumulator.java:219)
    at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:791)
    at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:745)
    at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:634)
    at com.test.kafka.KafkaProducerApp.main(KafkaProducerApp.java:45)

【问题讨论】:

  • 哪个版本的 Kafka?
  • kafka_2.10-0.10.1.1是kafka broker的版本
  • @Paizo- kafka_2.10-0.10.1.1 是 kafka broker 的版本 .. 有什么建议吗??
  • 不抱歉,也许您使用的是不同的库版本?在这里您可以找到一些信息 stackoverflow.com/q/47953901/3224238 为什么使用标题而不是键?

标签: kafka-producer-api


【解决方案1】:

通过不同的方式(Spring Cloud Streams Kafka 生产者)遇到相同的错误。我的问题是由于将inter.broker.protocol.version 设置为低于代理本身的版本see Kafka's broker configs 它允许的:

指定将使用哪个版本的代理间协议。这通常会在所有代理都升级到新版本后出现。

对于我的部署,代理是 v1.1,但 inter.broker.protocol.versionv0.10.1。当v1.1 代理收到带有标头的消息时,它很好,直到它使用不支持标头的v0.10.1 复制消息(产生"Magic v1" 的错误)。

【讨论】:

    猜你喜欢
    • 2016-09-21
    • 1970-01-01
    • 1970-01-01
    • 2015-08-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-03-30
    • 1970-01-01
    相关资源
    最近更新 更多