【问题标题】:How can I convert to String all value in a generic record in Avro with Kafka?如何使用 Kafka 将 Avro 中通用记录中的所有值转换为字符串?
【发布时间】:2022-12-06 02:38:31
【问题描述】:

有什么方法可以使用 Kafka 将 Avro 记录(包括嵌套数组)中的所有值转换为字符串?

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    最简单的方法是使用这些记录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 文档。

    【讨论】:

      猜你喜欢
      • 2019-06-03
      • 1970-01-01
      • 2014-04-18
      • 1970-01-01
      • 2021-10-28
      • 2021-08-08
      • 1970-01-01
      • 1970-01-01
      • 2012-04-26
      相关资源
      最近更新 更多