【发布时间】:2016-09-25 19:52:30
【问题描述】:
我们将在我们的项目中部署 Apache Kafka 2.10,并通过 JSON 对象在生产者和消费者之间进行通信。
到目前为止,我想我需要:
- 实现自定义序列化程序以将 JSON 转换为字节数组
- 实现自定义反序列化器将字节数组转换为 JSON 对象
- 生成消息
- 在 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