【发布时间】:2020-05-17 23:43:56
【问题描述】:
我正在尝试在 Python 3.7 中使用 Beam sdk 版本 2.20.0 构建 Apache Beam 管道,该管道已成功部署在 Dataflow 上,但似乎没有做任何事情。在worker日志中可以看到重复报如下错误信息
Error syncing pod xxxxxxxxxxx(), skipping: Failed to start container worker log
我已经尽我所能,但这个错误非常顽固,我的管道看起来像这样。
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import GoogleCloudOptions
from apache_beam.options.pipeline_options import StandardOptions
from apache_beam.options.pipeline_options import WorkerOptions
from apache_beam.options.pipeline_options import SetupOptions
from apache_beam.options.pipeline_options import DebugOptions
options = PipelineOptions()
options.view_as(GoogleCloudOptions).project = PROJECT
options.view_as(GoogleCloudOptions).job_name = job_name
options.view_as(GoogleCloudOptions).region = region
options.view_as(GoogleCloudOptions).staging_location = staging_location
options.view_as(GoogleCloudOptions).temp_location = temp_location
options.view_as(WorkerOptions).zone = zone
options.view_as(WorkerOptions).network = network
options.view_as(WorkerOptions).subnetwork = sub_network
options.view_as(WorkerOptions).use_public_ips = False
options.view_as(StandardOptions).runner = 'DataflowRunner'
options.view_as(StandardOptions).streaming = True
options.view_as(SetupOptions).sdk_location = ''
options.view_as(SetupOptions).save_main_session = True
options.view_as(DebugOptions).experiments = []
print('running pipeline...')
with beam.Pipeline(options=options) as pipeline:
(
pipeline
| 'ReadFromPubSub' >> beam.io.ReadFromPubSub(topic=topic_name).with_output_types(bytes)
| 'ProcessMessage' >> beam.ParDo(Split())
| 'WriteToBigQuery' >> beam.io.WriteToBigQuery(table=bq_table_name,
schema=bq_schema,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
)
result = pipeline.run()
我尝试使用 sdk_location 参数从计算实例中提供一个梁 sdk 2.20.0.tar.gz,但这也不起作用。我不能使用 sdk_location = default,因为它会触发从 pypi.org 下载。我在离线环境中工作,无法连接到互联网。任何帮助将不胜感激。
管道本身部署在容器上,所有与 apache beam 2.20.0 一起使用的库都在 requirements.txt 文件中指定,docker image 安装所有库。
【问题讨论】:
-
您好! 1. 你能添加一个包含所有导入的完整代码吗? 2. 请查看有关specifying subnetwork 的官方文档 - 您指定正确了吗?configuring options 和另一个SO thread
-
是的,我已经正确指定了网络和子网,相同的网络和子网可以正常工作于默认数据流作业。
-
所以,澄清一下,工作没有'sdk_location'标志但没有它? “默认数据流作业”是什么意思?该错误可能意味着 Dataflow 运行时环境无权下载执行作业所需的容器,或者容器不可用(如果您指定的是客户容器)。
-
默认数据流作业我的意思是数据流中可用的默认模板,我也为我自己的模板使用相同的网络和子网络值,所以网络设置应该没有问题。但是是的,没有连接到互联网,使用默认模板创建的作业可以正常工作,但使用 Python 设计的自定义作业会失败。
-
听起来您正在尝试为 Dataflow 指定自定义容器。我认为这还没有得到完全支持/记录。这里有一些信息stackoverflow.com/questions/44465818/…。但如果您需要更多具体信息,我建议您联系 Google Cloud 支持。
标签: python-3.x google-cloud-dataflow dataflow apache-beam