【问题标题】:external api call in apache beam dataflowApache Beam 数据流中的外部 api 调用
【发布时间】:2019-11-17 17:28:51
【问题描述】:

我有一个用例,我读入存储在谷歌云存储中的换行 json 元素并开始处理每个 json。在处理每个 json 时,无论之前是否发现了该 json 元素,我都必须调用一个外部 API 来进行重复数据删除。我在每个 json 上使用 DoFnDoFn

我还没有看到任何在线教程说明如何从 apache beam DoFn Dataflow 调用外部 API 端点。

我正在使用 Beam 的 JAVA SDK。我研究的一些教程解释了使用startBundleFinishBundle,但我不清楚如何使用它

【问题讨论】:

  • 这是流式管道还是批处理管道?
  • 这是一个批处理管道

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


【解决方案1】:

如果您需要在外部存储中检查每条 JSON 记录的重复项,那么您仍然可以使用 DoFn。有几个注解,如@Setup@StartBundle@FinishBundle 等,可用于注解DoFn 中的方法。

例如,如果您需要实例化客户端对象以向外部数据库发送请求,那么您可能希望在 @Setup 方法中执行此操作(如 POJO 构造函数),然后在您的 @ProcessElement 中利用此客户端对象方法。

让我们考虑一个简单的例子:

static class MyDoFn extends DoFn<Record, Record> {

    static transient MyClient client;

    @Setup
    public void setup() {
        client = new MyClient("host");
    }

    @ProcessElement
    public void processElement(ProcessContext c) {
        // process your records
        Record r = c.element();
        // check record ID for duplicates
        if (!client.isRecordExist(r.id()) {
            c.output(r);
        }
    }

    @Teardown
    public void teardown() {
        if (client != null) {
            client.close();
            client = null;
        }
    }
}

此外,为了避免对每条记录进行远程调用,您可以将记录批处理到内部缓冲区(将输入数据拆分为捆绑包)并以批处理模式检查重复项(如果您的客户端支持此操作)。为此,您可以使用 @StartBundle@FinishBundle 注释方法,这些方法将在相应地处理 Beam 包之前和之后调用。

对于更复杂的示例,我建议查看不同 Beam IO 中的 Sink 实现,例如 KinesisIO

【讨论】:

    【解决方案2】:

    以下博文中有一个使用有状态 DoFn 批量调用外部系统的示例:https://beam.apache.org/blog/2017/08/28/timely-processing.html,可能会有所帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-07-20
      • 2019-06-04
      • 2018-12-27
      • 2019-04-23
      • 2019-09-16
      • 2022-08-16
      • 1970-01-01
      • 2021-06-27
      相关资源
      最近更新 更多