【问题标题】:Bijection - Java Avro Serialization双射 - Java Avro 序列化
【发布时间】:2016-08-05 10:58:54
【问题描述】:

我正在寻找一个示例来对 Avro SpecificRecordBase 对象进行双射,类似于 GenericRecordBase,或者是否有更简单的方法将 AvroSerializer 类用作 Kafka 键和值序列化程序。

Injection<GenericRecord, byte[]> genericRecordInjection =
                                        GenericAvroCodecs.toBinary(schema);
byte[] bytes = genericRecordInjection.apply(type);

【问题讨论】:

    标签: serialization apache-kafka avro bijection


    【解决方案1】:

    https://github.com/miguno/kafka-storm-starter 提供了这样的示例代码。

    例如,参见AvroDecoderBolt。来自它的 javadocs:

    这个螺栓需要 Avro 编码的二进制格式的传入数据,根据 T 的 Avro 模式进行序列化。它将传入的数据反序列化为T pojo,并将此 pojo 发送给下游消费者。因此,这个螺栓可以被认为是 Twitter Bijection 的 Injection.invert[T, Array[Byte]](bytes) 用于 Avro 数据的 Storm 等价物。

    在哪里

    T:Avro 记录的类型(例如 Tweet)基于正在使用的底层 Avro 模式。必须是 Avro 的 SpecificRecordBase 的子类。

    代码的关键部分是(我把代码折叠成这个sn-p):

    // With T <: SpecificRecordBase
    
    implicit val specificAvroBinaryInjection: Injection[T, Array[Byte]] =
    SpecificAvroCodecs.toBinary[T]
    
    val bytes: Array[Byte] = ...; // the Avro-encoded data
    val decodeTry: Try[T] = Injection.invert(bytes)
    decodeTry match {
      case Success(pojo) =>
        System.out.println("Binary data decoded into pojo: " + pojo)
      case Failure(e) => log.error("Could not decode binary data: " + Throwables.getStackTraceAsString(e))
    }
    

    【讨论】:

      【解决方案2】:
      Schema.Parser parser = new Schema.Parser();
                  Schema schema = parser.parse(new File("/Users/.../schema.avsc"));
                  Injection<Command, byte[]> objectInjection = SpecificAvroCodecs.toBinary(schema);
                  byte[] bytes = objectInjection.apply(c);
      

      【讨论】:

      • 我是否正确假设架构仍然是对象本身的一部分?因为(通用)记录上有一个方法 .getSchema() 可用。在我看来,这似乎违背了拥有单独架构的全部目的
      猜你喜欢
      • 2013-03-19
      • 1970-01-01
      • 2012-09-23
      • 2020-03-04
      • 1970-01-01
      • 1970-01-01
      • 2019-08-30
      • 2020-05-21
      • 1970-01-01
      相关资源
      最近更新 更多