【问题标题】:what is the entry point of the workers in apache beam? (what methods are called?)apache Beam 中工作人员的切入点是什么? (调用什么方法?)
【发布时间】:2021-04-27 20:02:19
【问题描述】:

Apache Beam 有大量出色的文档,但我看不到为创建管道而运行的代码与工作人员运行的代码。我想我看到这段代码将运行一次,但它也会由每个启动的工作人员运行..

    public static void main(String[] args) {
        // Create the pipeline.
        PipelineOptions options =
            PipelineOptionsFactory.fromArgs(args).create();
        Pipeline p = Pipeline.create(options);

        // Create the PCollection 'lines' by applying a 'Read' transform.
        PCollection<String> lines = p.apply(
          "ReadMyFile", TextIO.read().from("gs://some/inputData.txt"));
    }

【问题讨论】:

    标签: google-cloud-dataflow apache-beam


    【解决方案1】:

    这是一个很好的问题,也是 Apache Beam 的核心。

    tl;dr 没有 用户定义的 工作人员启动时调用的入口点。

    长答案

    当您使用 Apache Beam SDK(使用应用程序等)进行编码时,您真正要做的是在后台创建一个包含所有应用转换的图表,请参阅文档 here。因此,一旦p.run() 被调用,图就会被发送给worker 执行。然后将图上的转换分解为组件并按顺序执行。

    至于您在问题中编写的代码,只会运行一次。当您执行 jar 时,该代码只运行一次。但是,图表中的转换是针对数据中的每个元素运行的(或更多或更少,具体取决于您的图表)。

    如果您对如何执行转换和用户定义函数 (ParDos) 感到好奇,那么入口点位于 Apache Beam SDK here

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-05-17
      • 2020-10-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多