【问题标题】:Producing and consuming JSON in KAFKA在 KAFKA 中生产和消费 JSON
【发布时间】:2016-09-25 19:52:30
【问题描述】:

我们将在我们的项目中部署 Apache Kafka 2.10,并通过 JSON 对象在生产者和消费者之间进行通信。

到目前为止,我想我需要:

  1. 实现自定义序列化程序以将 JSON 转换为字节数组
  2. 实现自定义反序列化器将字节数组转换为 JSON 对象
  3. 生成消息
  4. 在 Consumer 类中读取消息

关于第一点,我认为应该是这样的:

@Override
public byte[] serialize(String topic, T data) {
    if (data == null)
        return null;
    try {
        ObjectMapper mapper = new ObjectMapper();
        return mapper.writeValueAsBytes(data);
    } catch (Exception e) {
        throw new SerializationException("Error serializing JSON message", e);
    }
}

其中T data 可以作为字符串"{\"key\" : \"value\"}" 传递。

但是,现在 2-4 点存在问题。我在自定义反序列化器中试过这个:

@Override
public JsonNode deserialize(String topic, byte[] bytes) {
    if (bytes == null)
        return null;

    JsonNode data;
    try {
        objectMapper = new ObjectMapper();
        data = objectMapper.readTree(bytes);
    } catch (Exception e) {
        throw new SerializationException(e);
    }
    return data;
}

在我的消费者中我尝试过:

    KafkaConsumer<String, TextNode> consumer = new KafkaConsumer<String, TextNode>(messageConsumer.properties);
    consumer.subscribe(Arrays.asList(messageConsumer.topicName));
    int i = 0;
    while (true) {
        ConsumerRecords<String, TextNode> records = consumer.poll(100);
        for (ConsumerRecord<String, TextNode> record : records) {
            System.out.printf("offset = %d, key = %s, value = %s\n", record.offset(), record.key(), record.value().asText());
        }
    }

我认为这会产生一个正确的原始 json 字符串,但我调用 record.value().asText() 得到的只是一些哈希字符串 "IntcImtleVwiIDogXCJ2YWx1ZVwiIH0i"

非常感谢任何在 kafka 中通过 JSON 进行通信的建议或示例。

【问题讨论】:

  • 您查看过 Gson 库吗?不确定它是否适合 Kafka(因此这是一条评论),但 Gson 会自动为您将 Java 对象转换为 JSON,因此您一点也不麻烦。
  • 那实际上是哪个 Kafka 版本? “2.10”指的是Scala版本(Kafka是用Scala实现的),不是Kafka自己的版本。

标签: java json apache-kafka


【解决方案1】:

我建议您使用 UTF-8 编码作为字符串 JSON 序列化器:
1. Producer 获取数据为 JSON 字符串 ("{\"key\" : \"value\"}")
2. Producer 使用 UTF-8 将 JSON 字符串序列化为字节 (jsonString.getBytes(StandardCharsets.UTF_8);)
3. Producer 将此字节发送给 Kafka 4.消费者从Kafka读取字节 5. 消费者使用 UTF-8 (new String(consumedByteArray, StandardCharsets.UTF_8);) 将字节反序列化为 JSON 字符串
6. 消费者使用 JSON 字符串做任何需要的事情

我故意没有使用你的代码,所以流程是可以理解的,我认为你可以很容易地将这个例子应用到你的项目中:)

【讨论】:

    【解决方案2】:

    Apache Kafka 有一个内置的 JsonSerializer。我不确定 2.10 对某个版本意味着什么,但当前版本 0.10.0.1 绝对可以为您进行序列化。只需查找 JsonSerializer 类。

    【讨论】:

    • 2.10 后缀是用于编译二进制文件的 Scala 版本。
    猜你喜欢
    • 2020-11-11
    • 2023-01-21
    • 1970-01-01
    • 2020-08-13
    • 2019-01-15
    • 1970-01-01
    • 1970-01-01
    • 2016-12-02
    • 2018-01-07
    相关资源
    最近更新 更多