【问题标题】:Schema update while writing to Avro files写入 Avro 文件时更新架构
【发布时间】:2020-01-24 20:22:29
【问题描述】:

背景: 我们有一个 Dataflow 作业,它将 PubSub 消息转换为 Avro GenericRecords 并将它们作为“.avro”写入 GCS。 PubSub 消息和 GenericRecords 之间的转换需要一个模式。此架构每周更改一次,仅添加字段。我们希望能够在不更新 Dataflow 作业的情况下更新字段。

我们做了什么: 我们听取了this post 的建议,并创建了一个每分钟刷新一次内容的 Guava 缓存。刷新功能将从 GCS 中提取模式。然后我们让 FileIO.write 查询 Guava Cache 以获取最新的模式,并将具有模式的元素转换为 GenericRecord。我们还有 FileIO.write 输出到 Avro 接收器,它也是使用模式创建的。

代码如下:

genericRecordsAsByteArrays.apply(FileIO.<byte[]>write()
    .via(fn((input, c) -> {
          Map<String, Object> schemaInfo = cache.get("");
          Descriptors.Descriptor paymentRecordFd =
              (Descriptors.Descriptor) schemaInfo.get(DESCRIPTOR_KEY);
          DynamicMessage paymentRecordMsg = DynamicMessage.parseFrom(paymentRecordFd, input);
          Schema schema = (Schema) schemaInfo.get(SCHEMA_KEY);

          //From concrete PaymentRecord bytes to DynamicMessage
          try (ByteArrayOutputStream output = new ByteArrayOutputStream()) {
            BinaryEncoder encoder = EncoderFactory.get().directBinaryEncoder(output, null);
            ProtobufDatumWriter<DynamicMessage> pbWriter = new ProtobufDatumWriter<>(schema);
            pbWriter.write(paymentRecordMsg, encoder);
            encoder.flush();

            // From dynamic message to GenericRecord
            byte[] avroContents = output.toByteArray();
            DatumReader<GenericRecord> reader = new GenericDatumReader<>(schema);
            BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(avroContents, null);
            return reader.read(null, decoder);
          }
        }, requiresSideInputs()),
        fn((output, c) -> {
          Map<String, Object> schemaInfo = cache.get("");
          Schema schema = (Schema) schemaInfo.get(SCHEMA_KEY);
          return AvroIO.sink(schema).withCodec(CodecFactory.snappyCodec());
        }, requiresSideInputs()))
    .withNumShards(5)
    .withNaming(new PerWindowFilenames(baseDir, ".avro"))
    .to(baseDir.toString()));

我的问题:

  1. 当我们正在写入一个 Avro 文件时会发生什么,但突然发生架构更新,现在我们正在将新架构写入使用旧架构创建的 Avro 文件中?
  2. Dataflow 在看到新架构时是否会启动新文件?
  3. 在创建新文件之前,Dataflow 是否会忽略新架构和其他字段?

每个 Avro 文件在文件的开头都有自己的架构,所以我不确定预期的行为是什么。

【问题讨论】:

    标签: java google-cloud-dataflow avro apache-beam dataflow


    【解决方案1】:

    现在我们正在将新架构写入使用旧架构创建的 Avro 文件中

    这是不可能的。每个 Avro 文件只有一个模式。如果它发生变化,根据定义,您将写入一个新文件。

    我怀疑 Dataflow 会忽略字段。

    【讨论】:

    • 我觉得将 Protobuf 转换为 Avro 很奇怪
    • 是的.....由于遗留原因,我们需要对 Avro 进行 Proto。我们可以通过更改发布者来避免这种情况,但这需要更多的工作。
    • 所以你的意思是:如果我在“AvroIO.sink(schema).withCodec(CodecFactory.snappyCodec());”处更改 Avro 模式,它只会创建一个新文件来写入到?我能够确认它不会删除任何记录,但是检查它是否忽略了新字段是很棘手的
    • 它不应该忽略新字段。我也不确定你要写到哪里,但文件不会被更新或附加
    猜你喜欢
    • 2014-01-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多