【问题标题】:Does Apache Beam need internet to run GCP Dataflow jobsApache Beam 是否需要互联网来运行 GCP Dataflow 作业
【发布时间】:2019-05-17 20:19:19
【问题描述】:

我正在尝试在可以访问 GCP 资源但无法访问互联网的 GCP 虚拟机上部署数据流作业。当我尝试运行作业时,出现连接超时错误,如果我尝试连接到 Internet,这将是有意义的。代码中断是因为正在代表 apache-beam 尝试 http 连接。

Python 设置: 在切断 VM 之前,我使用 pip 和 requirements.txt 安装了所有必要的包。这似乎有效,因为代码的其他部分工作正常。

以下是我运行代码时收到的错误消息。

Retrying (Retry(total=0, connect=None, read=None, redirect=None, status=None)) 
after connection broken by 'ConnectTimeoutError(
<pip._vendor.urllib3.connection.VerifiedHTTPSConnection object at foo>, 
'Connection to pypi.org timed out. (connect timeout=15)')': /simple/apache-beam/

Could not find a version that satisfies the requirement apache-beam==2.9.0 (from versions: )

No matching distribution found for apache-beam==2.9.0

如果您正在运行 python 作业,是否需要连接到 pypi?这有什么技巧吗?

【问题讨论】:

  • 它是一个 GCP 虚拟机,那么为什么它不能访问互联网?我不确定我是否正确理解了这种情况。
  • 您的系统正在安装 Apache Beam 2.9.0。这是通过互联网进行的。如果您安装了不同版本的 Beam,请指定该版本。
  • @JohnHanley - 我在梁上安装了正确的版本。出于某种原因,它想尝试重新安装这是有问题的,因为我没有互联网。 apache-beam 代码中有什么东西告诉它寻找更新或补丁吗?如果是这样,我可以直接将其关闭吗?
  • Apache Beam 只是一个 Python 程序。在 requirements.txt 中指定您已安装的版本,以便匹配。
  • 对。是的,版本匹配。这就是为什么我不明白它为什么要重新安装它。也许我应该尝试更新软件包

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


【解决方案1】:

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

我也在这个设置上苦苦挣扎了很长时间。最后我找到了一个在执行时不需要互联网访问的解决方案。

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

作为先决条件,您需要创建 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 中断的情况下运行。

【讨论】:

    【解决方案2】:

    如果您在私有 Cloud Composer 中运行 DataflowPythonOperator,则作业需要访问 Internet 以从图像 projects/dataflow-service-producer-prod 下载一组包。但在私有集群中,VM 和 GKE 无法访问互联网。

    要解决这个问题,你需要创建一个Cloud NAT和一个路由器:https://cloud.google.com/nat/docs/gke-example#step_6_create_a_nat_configuration_using

    这将允许您的实例将数据包发送到互联网并接收入站流量。

    【讨论】:

      【解决方案3】:

      当我们使用启用了私有 ip 的谷歌云作曲家时,我们无法访问互联网。

      解决这个问题:

      • 创建 GKE 集群并创建一个新的节点池名称“default-pool”(使用相同的名称)。
      • 在网络标签中:添加“private”。
      • 在安全方面:勾选允许访问所有云 API。

      【讨论】:

      • 您好,您能详细说明一下这个解决方案吗?
      • 更新了评论@KevinEid
      猜你喜欢
      • 2017-09-30
      • 2021-04-29
      • 2018-01-03
      • 1970-01-01
      • 2016-03-29
      • 2013-11-03
      • 1970-01-01
      • 1970-01-01
      • 2015-09-28
      相关资源
      最近更新 更多