【发布时间】:2018-07-04 09:21:36
【问题描述】:
我们将数据作为字符串发送给 Kafka 生产者,消费者的最终输出是 Avro Schema 格式。
我需要使用 avro 模式解码最终输出。有人可以分享示例 java 代码来执行此操作吗?
【问题讨论】:
标签: apache-kafka kafka-consumer-api
我们将数据作为字符串发送给 Kafka 生产者,消费者的最终输出是 Avro Schema 格式。
我需要使用 avro 模式解码最终输出。有人可以分享示例 java 代码来执行此操作吗?
【问题讨论】:
标签: apache-kafka kafka-consumer-api
按照以下步骤-
1.从 avro 模式创建对象
java -jar /path/to/avro-tools-1.8.2.jar compile schema <schema file> <destination>
eg.
java -jar /path/to/avro-tools-1.8.2.jar compile schema user.avsc .
这将根据架构的命名空间在包中生成适当的源文件
private static void deserialize() {
try {
// Deserialize Users from disk
DatumReader<User> userDatumReader = new SpecificDatumReader<>(User.class);
DataFileReader<User> dataFileReader = new DataFileReader<User>(new File("users.avro"), userDatumReader);
User user = null;
while (dataFileReader.hasNext()) {
// Reuse user object by passing it to next(). This saves us from
// allocating and garbage collecting many objects for files with
// many items.
user = dataFileReader.next(user);
System.out.println("deserialized : "+user);
}
} catch (IOException e) {
e.printStackTrace();
}
}
【讨论】:
这是我的解决方法:此代码将以 Avro 模式格式打印消费者消息。
props.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer","io.confluent.kafka.serializers.KafkaAvroDeserializer");
props.put("schema.registry.url", "http://kafka-.XXX");
KafkaConsumer<String, avro_schema> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("topic"));
//infinite poll loop
try {
while (true) {
ConsumerRecords<String, avro_schema> records = consumer.poll(200L);
for (ConsumerRecord<String, avro_schema> record : records) {
System.out.println(record.value());
}
}
}
【讨论】: