【发布时间】:2022-12-06 02:38:31
【问题描述】:
有什么方法可以使用 Kafka 将 Avro 记录(包括嵌套数组)中的所有值转换为字符串?
【问题讨论】:
标签: apache-kafka
有什么方法可以使用 Kafka 将 Avro 记录(包括嵌套数组)中的所有值转换为字符串?
【问题讨论】:
标签: apache-kafka
最简单的方法是使用这些记录KafkaAvro解串器.
您可以使用一个简单的应用程序使用该主题,并根据需要处理每个反序列化的消息。为了反序列化 Avro 消息,您还需要将架构传递给消费者。
这是一个使用 Confluent Schema Registry 的工作示例:
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.avro.generic.GenericRecord;
import java.io.FileInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.Arrays;
import java.util.Properties;
import java.util.Random;
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "group1");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "io.confluent.kafka.serializers.KafkaAvroDeserializer");
props.put("schema.registry.url", "http://localhost:8081");
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
String topic = "topic1";
final Consumer<String, GenericRecord> consumer = new KafkaConsumer<String, GenericRecord>(props);
consumer.subscribe(Arrays.asList(topic));
try {
while (true) {
ConsumerRecords<String, GenericRecord> records = consumer.poll(100);
for (ConsumerRecord<String, GenericRecord> record : records) {
System.out.printf("offset = %d, key = %s, value = %s
", record.offset(), record.key(), record.value());
}
}
} finally {
consumer.close();
}
如果您需要将解码后的数据发送到一个新的主题,只需将反序列化的记录发送到一个新的主题Kafka生产者在同一个进程中,将值编码为字符串。也有可能出于同样的目的运行 Kafka Streams 应用程序。
我还鼓励您查看 this link 以了解有关此主题的 Confluent 文档。
【讨论】: