【问题标题】:Dataflow Error with Apache Beam SDK 2.20.0Apache Beam SDK 2.20.0 的数据流错误
【发布时间】: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


【解决方案1】:

TL;DR:将 Apache Beam SDK 存档复制到可访问的路径中,并将路径作为变量提供。

我也在为这个设置而苦苦挣扎。最后我找到了一个解决方案——即使你的问题是在几天前提出的,这个答案也可能对其他人有所帮助。

可能有多种方法可以做到这一点,但以下两种方法非常简单。

作为先决条件,您需要创建 apache-beam-sdk 源存档,如下所示:

  1. 克隆Apache Beam GitHub

  2. 切换到所需的标签,例如。 v2.28.0

  3. cd 到beam/sdks/python

  4. 创建所需 beam_sdk 版本的 tar.gz 源存档,如下所示:

    python setup.py sdist 
    
  5. 现在您应该在路径beam/sdks/python/dist/ 中拥有源存档apache-beam-2.28.0.tar.gz

选项 1 - 使用 Flex 模板并在 Dockerfile 中复制 Apache_Beam_SDK
文档:Google Dataflow Documentation

  1. 创建一个 Dockerfile --> 你必须包含这个 COPY utils/apache-beam-2.28.0.tar.gz /tmp,因为这将是你可以在 SetupOptions 中设置的路径。
FROM gcr.io/dataflow-templates-base/python3-template-launcher-base
ARG WORKDIR=/dataflow/template
RUN mkdir -p ${WORKDIR}

WORKDIR ${WORKDIR}

# Due to a change in the Apache Beam base image in version 2.24, you must to install
# libffi-dev manually as a dependency. For more information:
# https://github.com/GoogleCloudPlatform/python-docs-samples/issues/4891

# update used packages
RUN apt-get update && apt-get install -y \
    libffi-dev \
 && rm -rf /var/lib/apt/lists/*


COPY setup.py .
COPY main.py .

COPY path_to_beam_archive/apache-beam-2.28.0.tar.gz /tmp

ENV FLEX_TEMPLATE_PYTHON_SETUP_FILE="${WORKDIR}/setup.py"
ENV FLEX_TEMPLATE_PYTHON_PY_FILE="${WORKDIR}/main.py"

RUN python -m pip install --user --upgrade pip setuptools wheel
  1. 将 sdk_location 设置为您已将 apache_beam_sdk.tar.gz 复制到的路径:
    options.view_as(SetupOptions).sdk_location = '/tmp/apache-beam-2.28.0.tar.gz'
  1. 使用 Cloud Build 构建 Docker 映像
    gcloud builds submit --tag $TEMPLATE_IMAGE .
  2. 创建 Flex 模板
gcloud dataflow flex-template build "gs://define-path-to-your-templates/your-flex-template-name.json" \
 --image=gcr.io/your-project-id/image-name:tag \
 --sdk-language=PYTHON \
 --metadata-file=metadata.json
  1. 在您的子网中运行生成的 flex-template(如果需要)
gcloud dataflow flex-template run "your-dataflow-job-name" \
--template-file-gcs-location="gs://define-path-to-your-templates/your-flex-template-name.json" \
--parameters staging_location="gs://your-bucket-path/staging/" \
--parameters temp_location="gs://your-bucket-path/temp/" \
--service-account-email="your-restricted-sa-dataflow@your-project-id.iam.gserviceaccount.com" \
--region="yourRegion" \
--max-workers=6 \
--subnetwork="https://www.googleapis.com/compute/v1/projects/your-project-id/regions/your-region/subnetworks/your-subnetwork" \
--disable-public-ips

选项 2 - 从 GCS 复制 sdk_location
根据 Beam 文档,您甚至应该能够直接为选项 sdk_location 提供 GCS / gs:// 路径,但它对我不起作用。但以下应该有效:

  1. 将之前生成的存档上传到您可以从您要执行的数据流作业中访问的存储桶。可能类似于gs://yourbucketname/beam_sdks/apache-beam-2.28.0.tar.gz
  2. 将源代码中的 apache-beam-sdk 复制到例如。 /tmp/apache-beam-2.28.0.tar.gz
# see: https://cloud.google.com/storage/docs/samples/storage-download-file
from google.cloud import storage

def download_blob(bucket_name, source_blob_name, destination_file_name):
    """Downloads a blob from the bucket."""
    # bucket_name = "your-bucket-name"
    # source_blob_name = "storage-object-name"
    # destination_file_name = "local/path/to/file"

    storage_client = storage.Client()
    bucket = storage_client.bucket("gs://your-bucket-name")

    # Construct a client side representation of a blob.
    # Note `Bucket.blob` differs from `Bucket.get_blob` as it doesn't retrieve
    # any content from Google Cloud Storage. As we don't need additional data,
    # using `Bucket.blob` is preferred here.
    blob = bucket.blob("gs://your-bucket-name/path/apache-beam-2.28.0.tar.gz")
    blob.download_to_filename("/tmp/apache-beam-2.28.0.tar.gz")

  1. 现在您可以将 sdk_location 设置为您已下载 sdk 存档的路径。
options.view_as(SetupOptions).sdk_location = '/tmp/apache-beam-2.28.0.tar.gz'
  1. 现在您的 Pipeline 应该能够在没有 Internet 中断的情况下运行。

【讨论】:

  • 解决方案 1 对我有用。谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-22
  • 1970-01-01
  • 2018-07-20
  • 2019-06-04
  • 2019-09-16
  • 1970-01-01
相关资源
最近更新 更多