【发布时间】:2022-01-17 07:45:07
【问题描述】:
我的 Python 包具有以下结构,beam.py 是 Dataflow 的入口点脚本:
package_name\
__init__.py
tasks\
__init__.py
package_sum.py
utils\
__init__.py
beam.py
.gitignore
requirements.txt
setup.py
Dockerfile
beam.py:
import argparse
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from package_name.tasks import package_sum
def run(input_pubsub_topic, beam_args):
beam_pipeline_options = PipelineOptions(
beam_args,
save_main_session=True,
streaming=True
)
# Initialize the Beam pipeline
pipeline = beam.Pipeline(options=beam_pipeline_options)
pipeline | 'ReadFromPubSub' >> beam.io.ReadFromPubSub(input_pubsub_topic)
| 'Sum' >> beam.Map(package_sum)
# Run pipeline
pipeline.run().wait_until_finish()
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="test")
parser.add_argument("--input_pubsub_topic")
args, beam_args = parser.parse_known_args()
run(args.input_pubsub_topic, beam_args)
在我下面的Dockerfile 中,我安装了镜像中的包并将其下载到/tmp/dataflow-requirements-cache:
FROM gcr.io/dataflow-templates-base/python3-template-launcher-base:latest
ARG WORKDIR=/dataflow/template
RUN mkdir -p ${WORKDIR}
WORKDIR ${WORKDIR}
COPY requirements.txt .
COPY package_name package_name
# Do not include `apache-beam` in requirements.txt
ENV FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE="${WORKDIR}/requirements.txt"
ENV FLEX_TEMPLATE_PYTHON_PY_FILE="${WORKDIR}/package_name/utils/beam.py"
# Install apache-beam and other dependencies to launch the pipeline
RUN pip install --no-cache-dir --upgrade pip \
&& pip install --no-cache-dir apache-beam[gcp]==2.32.0
&& pip install --no-cache-dir -r $FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE \
&& pip install --no-cache-dir . \
# Download the requirements to speed up launching the Dataflow job.
&& pip download --no-cache-dir --dest /tmp/dataflow-requirements-cache -r $FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE \
&& pip download --no-cache-dir --dest /tmp/dataflow-requirements-cache .
# Since we already downloaded all the dependencies, there's no need to rebuild everything.
ENV PIP_NO_DEPS=True
当我启动 flex 模板作业时,它仍然会导致错误:ModuleNotFoundError: No module named 'package_name'。我该如何解决这个问题?
【问题讨论】:
-
你可以添加你的
setup.py的内容吗?
标签: google-cloud-dataflow apache-beam dataflow