【发布时间】: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<T, S>(其中T 和S 都可以使用反射进行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 来生成正确的架构,如果我尝试将T 或S 的架构推到fields 我得到:
java.lang.IllegalArgumentException: Unable to get field id from class null
2) 假设我知道如何编码MyOutput。我将什么传递给AvroIO.Write.withSchema?如果我通过 MyOutput.class 或架构,我会收到类型不匹配错误。
【问题讨论】: