【问题标题】:Avro schema Java deepcopy issue with field order字段顺序的 Avro 模式 Java deepcopy 问题
【发布时间】:2019-10-05 23:12:55
【问题描述】:

我目前正在研究使用 Java 处理特定 AVRO 模式演变场景时出现意外行为的解决方案,并在消费者中进行深度复制以将 GenericRecord 类解析为从 AVRO 模式生成的特定类。

为了解释发生了什么,我将使用一个简化的模式示例:

{
  "name":"SimpleEvent",
  "type":"record",
  "namespace":"com.simple.schemas",
  "fields":[
     {
        "name":"firstfield",
        "type":"string",
        "default":""
     },
     {
        "name":"secondfield",
        "type":"string",
        "default":""
     },
     {
        "name":"thirdfield",
        "type":"string",
        "default":""
     }
  ]
}

这只是一个包含三个字符串字段的简单模式,所有字段都是可选的,因为它们具有默认值。假设在某个时候我想添加另一个字符串字段,并删除一个字段,因为它不再需要,你最终会得到这个:

{
  "name":"SimpleEvent",
  "type":"record",
  "namespace":"com.simple.schemas",
  "fields":[
     {
        "name":"firstfield",
        "type":"string",
        "default":""
     },
     {
        "name":"secondfield",
        "type":"string",
        "default":""
     },
     {
        "name":"newfield",
        "type":"string",
        "default":""
     }
  ]
}

根据架构演变规则,这不应破坏更改。但是,当生产者开始使用更新的模式生成事件时,下游消费者会发生一些奇怪的事情。

原来生成的Java类(我使用Gradle avro插件生成类,但是maven插件和avro工具命令行代码生成产生相同的输出)只看字段顺序,而不要'不根据名称映射字段。

意味着字段“newfield”的值被下游消费者映射到“thirdfield”,这些消费者使用旧版本的架构来读取数据。

我发现了一些基于名称执行manual mapping 的工作,但是这不适用于嵌套对象。

通过一些本地实验,我还发现了另一种可以正确解决架构差异的方法:

    Schema readerSchema = SimpleEvent.getClassSchema();
    Schema writerSchema = request.getSchema();

    if (readerSchema.equals(writerSchema)){
        return (SimpleEvent)SpecificData.get().deepCopy(writerSchema, request);
    }

    DatumWriter<GenericRecord> writer = new SpecificDatumWriter<>(writerSchema);
    BinaryEncoder encoder = null;
    ByteArrayOutputStream stream = new ByteArrayOutputStream();
    encoder = EncoderFactory.get().binaryEncoder(stream, encoder);

    writer.write(request, encoder);
    encoder.flush();

    byte[] recordBytes = stream.toByteArray();

    Decoder decoder = DecoderFactory.get().binaryDecoder(recordBytes, null);

    SpecificDatumReader<SimpleEvent> specificDatumReader = new SpecificDatumReader(writerSchema, readerSchema);
    SimpleEvent result = specificDatumReader.read(null, decoder);
    return result;

但是,这似乎是一种相当浪费/不优雅的方法,因为您首先必须将 GenericRecord 转换为 byteArray,然后使用 SpecificDatumReader 再次读取它。

deepcopy 和 datumreader 类之间的区别在于,datumReader 类似乎适用于编写器架构与读取器架构不同的场景。

我觉得应该/可以有更好、更优雅的方式来处理这个问题。我真的很感激任何帮助/提示。

提前致谢:)

奥斯卡

【问题讨论】:

    标签: java kafka-consumer-api avro confluent-schema-registry


    【解决方案1】:

    在深入研究并查看了我们之前在侦听器中使用的 KafkaAvroDeserializer 之后,我注意到 AbstractKafkaAvroDeserializer 具有反序列化功能,可以在您可以传入读取器模式的位置进行反序列化。看起来不错,但确实有效!

    package com.oskar.generic.consumer.demo;
    
    import com.simple.schemas;
    
    import io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer;
    import io.confluent.kafka.serializers.KafkaAvroDeserializerConfig;
    import org.apache.kafka.common.serialization.Deserializer;
    
    import java.util.Map;
    
    public class SimpleEventDeserializer extends AbstractKafkaAvroDeserializer implements Deserializer<Object> {
    
    private boolean isKey;
    
    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        this.isKey = isKey;
        configure(new KafkaAvroDeserializerConfig(configs));
    }
    
    @Override
    public Object deserialize(String s, byte[] bytes) {
        return super.deserialize(bytes, SimpleEvent.getClassSchema());
    }
    
    @Override
    public void close() {
    
    }
    }
    

    然后像这样在消费者工厂中使用:

    @Bean
    public ConsumerFactory<String, GenericRecord> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:29095");
        props.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "one");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, SimpleEventDeserializer.class);
    
        return new DefaultKafkaConsumerFactory<>(props);
    }
    

    监听器代码本身如下所示:

     @KafkaListener(topics = "my-topic")
    public GenericRecord listen(@Payload GenericRecord request, @Headers MessageHeaders headers) throws IOException {
        SimpleEvent event = (SimpleEvent) SpecificData.get().deepCopy(request.getSchema(), request);
        return request;
    }
    

    【讨论】:

      猜你喜欢
      • 2018-02-01
      • 2023-02-18
      • 2018-12-16
      • 1970-01-01
      • 2022-12-06
      • 2021-02-15
      • 2016-06-17
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多