【问题标题】:Decode kafka consumer msg from string to avro using avro schema使用 avro 模式将 kafka 消费者 msg 从字符串解码为 avro
【发布时间】:2018-07-04 09:21:36
【问题描述】:

我们将数据作为字符串发送给 Kafka 生产者,消费者的最终输出是 Avro Schema 格式。

我需要使用 avro 模式解码最终输出。有人可以分享示例 java 代码来执行此操作吗?

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    按照以下步骤-

    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 .
    

    这将根据架构的命名空间在包中生成适当的源文件

    1. 使用 avro 架构和上面生成的类进行反序列化 例如,如果上面的命令将源类创建为 User.class 然后反序列化数据如下

    
    
    
    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();
            }
        }
    

    【讨论】:

    • 谢谢。我必须读取消费者流数据并解码每个偏移量。
    • 我们可以从文件中读取模式吗? Schema schema = new Schema.Parser().parse(new File("src/avrto_schema.avsc"));
    【解决方案2】:

    这是我的解决方法:此代码将以 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());
    
               }
          }
     }
    

    【讨论】:

      猜你喜欢
      • 2016-12-09
      • 2021-05-09
      • 1970-01-01
      • 2021-02-13
      • 2019-08-30
      • 2021-11-05
      • 1970-01-01
      • 2016-08-28
      • 2020-05-03
      相关资源
      最近更新 更多