【问题标题】:How to create dependency between tasks in Apache beam python如何在 Apache Beam python 中创建任务之间的依赖关系
【发布时间】:2018-03-17 10:12:11
【问题描述】:

我是 apache beam 的新手,正在探索 apache Beam 数据流的 python 版本。我想以特定顺序执行我的数据流任务,但它以并行模式执行所有任务。如何在 apache Beam python 中创建任务依赖?

示例代码:(在下面的代码中 sample.json 文件包含 5 行)

import apache_beam as beam
import logging
from apache_beam.options.pipeline_options import PipelineOptions

class Sample(beam.PTransform):
    def __init__(self, index):
        self.index = index

    def expand(self, pcoll):
        logging.info(self.index)
        return pcoll

class LoadData(beam.DoFn):
    def process(self, context):
        logging.info("***")

if __name__ == '__main__':

    logging.getLogger().setLevel(logging.INFO)
    pipeline = beam.Pipeline(options=PipelineOptions())

    (pipeline
        | "one" >> Sample(1)
        | "two: Read" >> beam.io.ReadFromText('sample.json')
        | "three: show" >> beam.ParDo(LoadData())
        | "four: sample2" >> Sample(2)
    )
    pipeline.run().wait_until_finish()

我预计它将按照一、二、三、四的顺序执行。但它是以并行模式运行的。

以上代码的输出:

INFO:root:Missing pipeline option (runner). Executing pipeline using the 
default runner: DirectRunner.
INFO:root:1
INFO:root:2
INFO:root:Running pipeline with DirectRunner.
INFO:root:***
INFO:root:***
INFO:root:***
INFO:root:***
INFO:root:***

【问题讨论】:

  • 你想通过按顺序执行来完成什么?另外,我不确定您的“示例”转换应该做什么:在实施时,它什么也不做。还要记住,就像数据库查询计划一样,首先构建管道(当您看到来自 expand() 的日志记录时),然后由运行器优化并执行(当您看到 "* **").
  • @jkff 我想将数据从 biquery 加载到 elasticsearch。在我的示例转换中,我正在执行创建、重新索引、删除弹性搜索索引等操作。所以首先我需要创建一个临时索引,第二个加载数据和 ES 临时索引,第三个重新索引它,第四个删除我的临时索引。所以我想以有序的方式执行所有这些任务。但这里首先执行创建、重新索引和删除任务,最后执行加载数据。 (你可以看到最后显示的日志“*****”)

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


【解决方案1】:

根据Dataflow's documentation

当管道运行器为分布式构建您的实际管道时 执行时,可以优化管道。例如,可能更多 一起运行某些变换的计算效率很高,或者在一个 不同的顺序。数据流服务完全管理这方面的 管道的执行。

同样按照Apache Beam's documentation:

API 强调并行处理元素,这使得它 难以表达诸如“为每个人分配一个序列号”之类的动作 PCollection 中的元素”。这是故意的,因为这样的算法是 更有可能遭受可扩展性问题的困扰。处理所有 并行元素也有一些缺点。具体来说,它使 不可能批处理任何操作,例如将元素写入 处理期间的接收器或检查点进度

所以问题在于 Dataflow 和 Apache Beam 本质上是并行的;它们旨在处理令人尴尬的并行用例,如果您需要以特定顺序执行操作,它们可能不是最好的工具。正如@jkff 指出的那样,Dataflow 将优化 Pipeline,使其以最佳方式并行化操作。

如果您确实需要按连续顺序执行每个步骤,解决方法是使用blocking execution,而不是使用the waitUntilFinish() method,如其他Stack Overflow answer 中所述。但是,我的理解是这样的实现只能在批处理管道中工作,因为流管道会连续消耗数据,因此您不能阻止执行以在连续步骤上工作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-10-29
    • 2011-07-06
    • 1970-01-01
    • 2021-11-29
    • 1970-01-01
    • 2020-07-20
    • 2019-10-28
    相关资源
    最近更新 更多