TL;DR:将 Apache Beam SDK 存档复制到可访问的路径中,并在 Dataflow 管道中将路径作为 SetupOption sdk_location 变量提供。
我也在这个设置上苦苦挣扎了很长时间。最后我找到了一个在执行时不需要互联网访问的解决方案。
可能有多种方法可以做到这一点,但以下两种方法相当简单。
作为先决条件,您需要创建 apache-beam-sdk 源存档,如下所示:
-
克隆Apache Beam GitHub
-
切换到所需的标签,例如。 v2.28.0
-
cd 到beam/sdks/python
-
创建所需 beam_sdk 版本的 tar.gz 源存档,如下所示:
python setup.py sdist
-
现在您应该在路径beam/sdks/python/dist/ 中拥有源存档apache-beam-2.28.0.tar.gz
选项 1 - 使用 Flex 模板并在 Dockerfile 中复制 Apache_Beam_SDK
文档:Google Dataflow Documentation
- 创建一个 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
- 将 sdk_location 设置为您已将 apache_beam_sdk.tar.gz 复制到的路径:
options.view_as(SetupOptions).sdk_location = '/tmp/apache-beam-2.28.0.tar.gz'
- 使用 Cloud Build 构建 Docker 映像
gcloud builds submit --tag $TEMPLATE_IMAGE .
- 创建 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
- 在您的子网中运行生成的 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:// 路径,但它对我不起作用。但以下应该有效:
- 将之前生成的存档上传到您可以从您要执行的数据流作业中访问的存储桶。可能类似于
gs://yourbucketname/beam_sdks/apache-beam-2.28.0.tar.gz
- 将源代码中的 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")
- 现在您可以将 sdk_location 设置为您已下载 sdk 存档的路径。
options.view_as(SetupOptions).sdk_location = '/tmp/apache-beam-2.28.0.tar.gz'
- 现在您的 Pipeline 应该能够在没有 Internet 中断的情况下运行。