【问题标题】:Dataflow output parameterized type to avro file数据流将参数化类型输出到 avro 文件
【发布时间】:2023-04-10 20:27:02
【问题描述】:

我有一个成功输出 Avro 文件的管道,如下所示:

@DefaultCoder(AvroCoder.class)
class MyOutput_T_S {
  T foo;
  S bar;
  Boolean baz;
  public MyOutput_T_S() {}
}

@DefaultCoder(AvroCoder.class)
class T {
  String id;
  public T() {}
}

@DefaultCoder(AvroCoder.class)
class S {
  String id;
  public S() {}
}
...
PCollection<MyOutput_T_S> output = input.apply(myTransform);
output.apply(AvroIO.Write.to("/out").withSchema(MyOutput_T_S.class));

除了参数化输出MyOutput&lt;T, S&gt;(其中TS 都可以使用反射进行Avro 编码)之外,我如何重现这种确切的行为。

主要问题是 Avro 反射不适用于参数化类型。所以根据这些回复:

1) 我想我需要编写一个自定义的CoderFactory,但是我很难弄清楚它是如何工作的(我很难找到示例)。奇怪的是,一个完全幼稚的编码器工厂似乎让我运行管道并使用 DataflowAssert 检查正确的输出:

cr.RegisterCoder(MyOutput.class, new CoderFactory() {
  @Override
  public Coder<?> create(List<? excents Coder<?>> componentCoders) {
    Schema schema = new Schema.Parser().parse("{\"type\":\"record\,"
      + "\"name\":\"MyOutput\","
      + "\"namespace\":\"mypackage"\","
      + "\"fields\":[]}"
    return AvroCoder.of(MyOutput.class, schema);
  }
  @Override
  public List<Object> getInstanceComponents(Object value) {
    MyOutput<Object, Object> myOutput = (MyOutput<Object, Object>) value;
    List components = new ArrayList();
    return components;
  }

虽然我现在可以成功地断言输出,但我希望这不会因为写入文件而中断。我还没有弄清楚我应该如何使用提供的componentCoders 来生成正确的架构,如果我尝试将TS 的架构推到fields 我得到:

java.lang.IllegalArgumentException: Unable to get field id from class null

2) 假设我知道如何编码MyOutput。我将什么传递给AvroIO.Write.withSchema?如果我通过 MyOutput.class 或架构,我会收到类型不匹配错误。

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    我认为有两个问题(如果我错了,请纠正我):

    1. 如何启用编码器注册表,以便为MyOutput&lt;T, S&gt; 的各种参数化提供编码器?
    2. 如何使用AvroIO.WriteMyOutput&lt;T, S&gt; 的值添加到文件中。

    第一个问题是通过在您找到的链接问题中注册CoderFactory 来解决。

    您的幼稚编码器可能允许您毫无问题地运行管道,因为序列化正在被优化掉。当然,没有字段的 Avro 模式会导致这些字段在序列化+反序列化往返过程中被丢弃。

    但假设您使用字段填写架构,您对CoderFactory#create 的方法看起来是正确的。我不知道消息java.lang.IllegalArgumentException: Unable to get field id from class null 的确切原因,但是对于适当组装的schema,对AvroCoder.of(MyOutput.class, schema) 的调用应该可以工作。如果这有问题,更多详细信息(例如堆栈跟踪的其余部分)会有所帮助。

    但是,您对 CoderFactory#getInstanceComponents 的覆盖应该返回一个值列表,每个类型参数 MyOutput 一个。像这样:

    @Override
    public List<Object> getInstanceComponents(Object value) {
      MyOutput<Object, Object> myOutput = (MyOutput<Object, Object>) value;
      return ImmutableList.of(myOutput.foo, myOutput.bar);
    }
    

    第二个问题可以使用一些与第一个相同的支持代码来回答,但在其他方面是独立的。 AvroIO.Write.withSchema 始终明确使用提供的架构。它确实在后台使用了AvroCoder,但这实际上是一个实现细节。提供一个兼容的架构是所有必要的 - 必须为您想要输出 MyOutput&lt;T, S&gt;TS 的每个值组成这样的架构。

    【讨论】:

    • 感谢您澄清我走在正确的轨道上。如果我继续遇到麻烦,我会再试一次并尝试更具体地说明错误。
    猜你喜欢
    • 2014-06-07
    • 1970-01-01
    • 2023-01-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-04
    • 2013-07-15
    • 1970-01-01
    相关资源
    最近更新 更多