【问题标题】:What does an Apache Beam Dataflow job do locally?Apache Beam 数据流作业在本地做什么?
【发布时间】:2018-04-27 17:31:49
【问题描述】:

我在使用 Apache Beam Python SDK 定义的数据流时遇到了一些问题。如果我单步执行我的代码,它会到达 pipeline.run() 步骤,我认为这意味着执行图已成功定义。但是,该作业从未在 Dataflow 监控工具上注册,this 让我认为它永远不会到达管道验证步骤。

我想详细了解这两个步骤之间发生的情况,以帮助调试问题。我看到输出表明我的 requirements.txtapache-beam 中的包正在安装 pip,并且在发送到 Google 的服务器之前似乎有些东西被腌制了。这是为什么?如果我已经下载了 apache-beam,为什么还要重新下载呢?究竟是什么被腌制了?

我不是在这里寻找解决我的问题的方法,只是想更好地理解这个过程。

【问题讨论】:

  • 对于 Apache Beam 的 pip 安装,您的机器上是否安装了最新版本?安装有没有报错?
  • 我将最新版本安装到虚拟环境中。出于某种原因,第二个 apache-beam 下载会抱怨找不到 nose 并退出。将nose 安装到虚拟环境中修复了该问题。但是,将其放入 requirements.txt 不会。不知道为什么它首先需要一个测试包。

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


【解决方案1】:

在图构建期间,Dataflow 会检查管道中的错误和任何非法操作。检查成功后,执行图将转换为 JSON 并传输到 Dataflow 服务。在 Dataflow 服务中,JSON 图形经过验证并成为一项工作。 但是,如果管道在本地执行,则图形不会转换为 JSON 或传输到 Dataflow 服务。因此,该图不会在监控工具中显示为作业,它将在本地计算机上运行 [1]。您可以按照文档配置本地机器 [2]。

[1]https://cloud.google.com/dataflow/service/dataflow-service-desc#pipeline-lifecycle-from-pipeline-code-to-dataflow-job

[2]https://cloud.google.com/dataflow/pipelines/specifying-exec-params#configuring-pipelineoptions-for-local-execution

【讨论】:

  • 此计算设置为使用 DataflowRunner,管道不应在本地执行。自从询问以来,我已经深入研究了代码。转换为 JSON 是 pickling 的用武之地。管道操作使用dill 进行序列化,并作为 JSON 请求的一部分发送到 Google 的服务器。我仍然不确定为什么我看到第二个 apache-beam 副本在本地下载。
  • “sdk_location”字段的默认管道选项允许 Beam SDK 在工作流提交期间从指定位置下载或复制 SDK。删除“sdk_location”字段的“默认”字符串值并将其留空将导致 SDK 未复制 [1]。 [1]beam.apache.org/documentation/runners/dataflow/…
猜你喜欢
  • 2018-02-07
  • 2020-02-21
  • 1970-01-01
  • 2018-07-20
  • 2019-06-04
  • 2023-02-12
  • 2019-09-16
  • 2018-12-27
  • 1970-01-01
相关资源
最近更新 更多