【问题标题】:Kafka Java Producer API unable to serialize key as Long or IntKafka Java Producer API 无法将密钥序列化为 Long 或 Int
【发布时间】:2020-09-23 14:22:54
【问题描述】:

这是在 Kafka 中生成数据的 Java 代码:

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.LongSerializer;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class ExampleClass {
  private final static String TOPIC = "my-example-topic";
  private final static String BOOTSTRAP_SERVERS = "confbroker:9092";

  private static Producer<Long, String> createProducer() {
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName());
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    return new KafkaProducer<>(props);
  }
  private static void runProducer() throws Exception {
    final Producer<Long, String> producer = createProducer();
    long sensorId = 1001L;
    try {
      for (long index = sensorId; index < sensorId + 5; index++) {
        final ProducerRecord<Long, String> record = new ProducerRecord<>(TOPIC, index, "This is sensor no: " + index);
        RecordMetadata metadata = producer.send(record).get();
        System.out.printf("sent record(key=%s value=%s) " + "meta(partition=%d, offset=%d)\n", record.key(),
            record.value(), metadata.partition(), metadata.offset());
      }
    } finally {
      producer.flush();
      producer.close();
    }
  }
  public static void main(String... args) throws Exception {
      runProducer();
  }
}

Confluent 5.4.0 中运行控制台使用者时,我得到的结果为:

关键是乱码。

如何生成 IntLong 类型的 Key

PS:

=> Confluent 5.5 中的结果也相同。

=> 与 IntegerSerializer 的结果相同。

【问题讨论】:

    标签: apache-kafka kafka-producer-api confluent-platform


    【解决方案1】:

    控制台使用者使用 StringDeserialisers 作为键和值的默认值。如果您想将密钥反序列化为Long,则必须在控制台消费者命令中明确提及:

    --property key.deserializer org.apache.kafka.common.serialization.LongDeserializer
    

    【讨论】:

    • 感谢@mike 的快速回复.....只是跟进,流也一样,因为我从这个主题创建了一个流,并且 rowkey 是一样的乱码。
    • 是的,将是一样的。在 Kafka 中,所有内容都存储为字节,而 Kafka 本身不知道这些数据是如何序列化的。您的消费者需要这些知识才能正确读取数据。
    猜你喜欢
    • 2017-03-22
    • 2021-03-22
    • 1970-01-01
    • 1970-01-01
    • 2021-12-27
    • 2020-11-26
    • 2017-12-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多