【问题标题】:Howto catch exceptions thrown by built-in transform in Apache Beam (in this case JSON Parsing)如何捕获 Apache Beam 中内置转换引发的异常(在本例中为 JSON 解析)
【发布时间】: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


【解决方案1】:

不幸的是,我不认为他们有任何干净的方式来做到这一点。但是,如果您打算编写自己的转换,那么如果您可以增强 ParseJsons 以添加无效 json 的可选输出流,那就太好了。一般来说,它可能很有用。

【讨论】:

    【解决方案2】:

    我只是想补充一下 Ankur 之前所说的参考,BEAM-5638 中已经完成了一些工作以添加异常处理,但尚未完全完成/合并 JSON 转换。

    编辑:Exceptions handling for Json transforms 最近是merged,它将很快在 Beam 2.17 版本中可用。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-06-18
      • 1970-01-01
      • 1970-01-01
      • 2019-09-13
      • 1970-01-01
      • 1970-01-01
      • 2015-06-12
      相关资源
      最近更新 更多