【问题标题】:Dataflow reading from PubSub works at GCP, can't run locally从 PubSub 读取的数据流在 GCP 上工作,无法在本地运行
【发布时间】:2021-07-17 01:19:05
【问题描述】:

我有一个小型测试数据流作业,它只是从 PubSub 订阅中读取并丢弃消息,我们正在使用它来开始一些概念验证工作。

它在 GCP 上运行良好,但在本地失败。我的期望是相同的代码应该以任何一种方式工作,只需切换 Dataflow 运行器,但也许情况并非如此?代码如下:

import os
from datetime import datetime
import logging

from apache_beam import Map, io, Pipeline
from apache_beam.options.pipeline_options import PipelineOptions

def noop(element):
    pass

def run(input_subscription, pipeline_args=None):
    pipeline_options = PipelineOptions(
        pipeline_args, streaming=True, save_main_session=True
    )

    with Pipeline(options=pipeline_options) as pipeline:
        (
            pipeline
            | "Read from Pub/Sub" >> io.ReadFromPubSub(subscription=input_subscription, with_attributes=True)
            | "noop" >> Map(noop)
        )


if __name__ == "__main__":
    logging.getLogger().setLevel(logging.INFO)

    run(
        os.environ['INPUT_SUBSCRIPTION'],
        [
            '--runner', os.getenv('RUNNER', 'DirectRunner'),
            '--project', os.getenv('PROJECT'),
            '--region', os.getenv('REGION'),
            '--temp_location', os.getenv('TEMP_LOCATION'),
            '--service_account_email', os.getenv('SERVICE_ACCOUNT_EMAIL'),
            '--network', os.getenv('NETWORK'),
            '--subnetwork', os.getenv('SUBNETWORK'),
            '--num_workers', os.getenv('NUM_WORKERS'),
        ]
    )

如果我使用这个命令行运行它,它会在 Google Cloud 中创建并运行该作业:

INPUT_SUBSCRIPTION=subscriptionname \
RUNNER=DataflowRunner \
PROJECT=project \
REGION=region \
TEMP_LOCATION=gs://somewhere/temp \
SERVICE_ACCOUNT_EMAIL=serviceaccount@project.iam.gserviceaccount.com \
NETWORK=network \
SUBNETWORK=https://www.googleapis.com/compute/v1/projects/project/regions/region/subnetworks/subnetwork \
NUM_WORKERS=3 \
python read-pubsub-with-dataflow.py

如果我省略了RUNNER 选项,那么它使用DirectRunner

INPUT_SUBSCRIPTION=subscriptionname \
PROJECT=project \
REGION=region \
TEMP_LOCATION=gs://somewhere/temp \
SERVICE_ACCOUNT_EMAIL=serviceaccount@project.iam.gserviceaccount.com \
NETWORK=network \
SUBNETWORK=https://www.googleapis.com/compute/v1/projects/project/regions/region/subnetworks/subnetwork \
NUM_WORKERS=3 \
python read-pubsub-with-dataflow.py

它失败并显示大量错误消息,但我只包含第一个(我认为其余的只是级联):

INFO:apache_beam.runners.direct.direct_runner:Running pipeline with DirectRunner.
/Users/denis/redacted/env/lib/python3.6/site-packages/google/auth/_default.py:70: UserWarning: Your application has authenticated using end user credentials from Google Cloud SDK without a quota project. You might receive a "quota exceeded" or "API not enabled" error. We recommend you rerun `gcloud auth application-default login` and make sure a quota project is added. Or you can use service accounts instead. For more information about service accounts, see https://cloud.google.com/docs/authentication/
  warnings.warn(_CLOUD_SDK_CREDENTIALS_WARNING)
ERROR:apache_beam.runners.direct.executor:Exception at bundle <apache_beam.runners.direct.bundle_factory._Bundle object at 0x7fed3e368448>, due to an exception.
 Traceback (most recent call last):
  File "/Users/denis/redacted/env/lib/python3.6/site-packages/apache_beam/runners/direct/transform_evaluator.py", line 694, in _read_from_pubsub
    self._sub_name, max_messages=10, return_immediately=True)
  File "/Users/denis/redacted/env/lib/python3.6/site-packages/google/cloud/pubsub_v1/_gapic.py", line 40, in <lambda>
    fx = lambda self, *a, **kw: wrapped_fx(self.api, *a, **kw)  # noqa
  File "/Users/denis/redacted/env/lib/python3.6/site-packages/google/pubsub_v1/services/subscriber/client.py", line 1106, in pull
    "If the `request` argument is set, then none of "
ValueError: If the `request` argument is set, then none of the individual field arguments should be set.

During handling of the above exception, another exception occurred:
...etc...

我怀疑这可能与凭据有关?还是我们的项目配置?也许我应该尝试一个新的空白项目。

【问题讨论】:

  • 错误说明了一切,它与凭据以及管道如何使用 GCP 进行身份验证有关。我建议您改用服务帐户,并为该服务帐户提供所有需要的权限。

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


【解决方案1】:

结果证明这是不兼容的软件包版本。我的requirements.txt 是:

apache_beam[gcp]
google_apitools
google-cloud-pubsub

但那是在安装破坏apache_beamgoogle-cloud-pubsub 软件包版本。我将requirements.txt 更改为:

apache_beam[gcp]
google_apitools

现在一切正常!

对于它的价值,使用DirectRunner 在本地运行我显然不需要DataflowRunner 所需的很多选项。这足够了:

GOOGLE_APPLICATION_CREDENTIALS=/path/to/credentials.json \
RUNNER=DirectRunner \
INPUT_SUBSCRIPTION=projects/mytopic/subscriptions/mysubscription \
python read-pubsub-with-dataflow.py

【讨论】:

    猜你喜欢
    • 2020-06-09
    • 1970-01-01
    • 1970-01-01
    • 2019-06-15
    • 1970-01-01
    • 2017-12-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多