【发布时间】:2019-08-23 12:33:02
【问题描述】:
我的管道如下(其中 StringToKVTransForm、kafkaoutput、kafkainput 是我在其他地方创建或配置的转换;这里的重点是 ParseJsons,因为它是一个内置的转换
try {
PCollection<MyClass> myObjects = p
.apply(kafkaInput.withoutMetadata())
.apply(Values.create())
.apply(ParseJsons.of(MyClass.class)).setCoder(SerializableCoder.of(MyClass.class))
.apply(AsJsons.of(MyClass.class))
.apply(new StringToKvTransform())
.apply(kafkaOutput);
} catch (Throwable e){
log.info("Unexpected error", e);
}
log.info("pipeline initialized");
p.run().waitUntilFinish();
}
这里的问题是,由于各种原因,我得到的数据可能并不总是正确的 json 格式;不幸的是,这会导致整个管道崩溃并出现异常
org.apache.beam.sdk.util.UserCodeException: java.lang.RuntimeException: 无法从 JSON 值解析 path.to.MyClass: { “myIncorrectJsonString” }
在这种情况下,我希望我的管道继续运行并忽略不正确的输入事件,但是,我不知道如何...
原因是这是一个内置的转换(ParseJsons),它似乎把错误抛到了我无法控制的地方,导致整个程序崩溃。
Allthe 我看过的教程建议在转换中捕获错误,这显然不是这里的选项。
我的 goto 解决方案是扩展 ParseJsons 类并捕获错误,但它有一个私有构造函数,因此无法扩展。
有什么想法,还是我必须编写自己的 ParseJsons 转换类?
【问题讨论】:
-
不幸的是,我不认为他们是任何干净的方式来做到这一点。但是,如果您打算编写自己的转换,那么如果您可以增强
ParseJsons以添加无效 json 的可选输出流,那就太好了。一般情况下它可能很有用。 -
谢谢!我也是这么想的。如果你把它变成一个答案,我会接受它:)
标签: json apache-beam