【问题标题】:Running function once Dataflow Batch-Job step has completed数据流批处理作业步骤完成后运行函数
【发布时间】:2020-05-28 13:54:45
【问题描述】:

我有一个数据流作业,它有一个扇出步骤,每个步骤都将结果写入 GCS 上的不同文件夹。在批处理作业执行期间,每个文件夹都会写入数百个文件。

我想确定 FileIO 步骤何时完成,以便运行将文件夹的全部内容加载到 BigQuery 表的 java 代码。

我知道我可以使用 Cloud Functions 和 PubSub 通知为每个书面文件执行此操作,但我更喜欢仅在整个文件夹完成时执行一次。

谢谢!

【问题讨论】:

  • 您是否有理由不直接从 Dataflow 管道将结果与 GCS 同时写入 BigQuery?
  • 是的。有一些逻辑和技术原因。对我来说,将 AVRO 文件存储在 GCS 上要简单得多,然后一次将它们加载到 BQ

标签: java google-cloud-dataflow


【解决方案1】:

有两种方法可以做到这一点:

在你的管道之后执行它。

运行您的管道并在您的管道结果上调用waitUntilFinish(Python 中为wait_until_finish)以延迟执行直到您的管道完成后,如下所示:

pipeline.run().waitUntilFinish();

您可以根据 waitUntilFinish 的结果验证管道是否成功完成,然后您可以将文件夹的内容加载到 BigQuery。这种方法唯一需要注意的是,您的代码不是 Dataflow 管道的一部分,因此如果您在该步骤中依赖管道中的元素,它将更加困难。

在 FileIO.Write 之后添加转换

FileIO.Write 转换的结果是WriteFilesResult,它允许您通过调用getPerDestinationOutputFilenames 来获取包含写入文件的所有文件名的 PCollection。从那里,您可以使用可以将所有这些文件写入 BigQuery 的转换继续您的管道。下面是一个 Java 示例:

WriteFilesResult<DestinationT> result = files.apply(FileIO.write()...)
result.getPerDestinationOutputFilenames().apply(...)

Python 中的等效项似乎被称为 FileResult,但我找不到该文件的好文档。

【讨论】:

    【解决方案2】:

    @Daniel Oliveira 建议了一种您可以遵循的方法,但在我看来,这不是最好的方法。

    我与他不同的两个原因:

    1. 处理作业失败的范围狭窄:考虑这样一种情况,即您的 Dataflow 作业成功但加载到 Big Query 作业失败。由于这种紧密耦合,您将无法重新运行第二个作业。
    2. 第二个作业的性能将成为瓶颈:在文件大小会增长的生产场景中,您的加载作业将成为其他依赖进程的瓶颈

    正如您已经提到的,您不能在同一份工作中直接写信给 BQ。我会建议你以下方法:

    1. 创建另一个梁作业以将所有文件加载到 BQ。您可以参考this阅读beam中的多个文件。
    2. 使用 Dataflow Java OperatorDataflow Template Operator 使用 Cloud Composer 编排代码。设置气流触发规则为'all_sucess'并设置job1.setUpstream(job2)。请参考气流文档here

    希望对你有帮助

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-12-07
      • 1970-01-01
      • 2023-03-24
      • 1970-01-01
      • 2017-07-06
      相关资源
      最近更新 更多