【问题标题】:Submit a Python project to Dataproc job将 Python 项目提交到 Dataproc 作业
【发布时间】:2020-08-06 17:36:04
【问题描述】:

我有一个python项目,它的文件夹有结构

main_directory - lib - lib.py
               - run - script.py

script.py

from lib.lib import add_two
spark = SparkSession \
    .builder \
    .master('yarn') \
    .appName('script') \
    .getOrCreate()

print(add_two(1,2))

lib.py

def add_two(x,y):
    return x+y

我想在 GCP 中作为 Dataproc 作业启动。我在网上查过,但我不太明白怎么做。我正在尝试使用

启动脚本
gcloud dataproc jobs submit pyspark --cluster=$CLUSTER_NAME --region=$REGION \
  run/script.py

但我收到以下错误消息:

from lib.lib import add_two
ModuleNotFoundError: No module named 'lib.lib'

您能否帮助我了解如何在 Dataproc 上启动这项工作?我发现这样做的唯一方法是删除绝对路径,将此更改为script.py

 from lib import add_two

并以

的身份启动作业
gcloud dataproc jobs submit pyspark --cluster=$CLUSTER_NAME --region=$REGION \
  --files /lib/lib.py \
  /run/script.py

但是,我想避免每次手动列出文件的繁琐过程。

按照@Igor 的建议,打包成一个zip 文件,我发现

zip -j --update -r libpack.zip /projectfolder/* && spark-submit --py-files libpack.zip /projectfolder/run/script.py

有效。但是,这会将所有文件放在 libpack.zip 中的同一个根文件夹中,因此如果子文件夹中有同名文件,这将不起作用。

有什么建议吗?

【问题讨论】:

    标签: python pyspark google-cloud-dataproc


    【解决方案1】:

    如果您想在提交 Dataroc 作业时保留项目结构,那么您应该将您的项目打包成一个 .zip 文件,并在提交作业时在 --py-files 参数中指定:

    gcloud dataproc jobs submit pyspark --cluster=$CLUSTER_NAME --region=$REGION \
      --py-files lib.zip \
      run/script.py
    

    要创建 zip 存档,您需要运行脚本:

    cd main_directory/
    zip -x run/script.py -r libs.zip .
    

    请参阅this blog post,了解有关如何将依赖项打包到 PySpark 作业的 zip 存档中的更多详细信息。

    【讨论】:

    • 这是一个不错的建议。但是,如何从script.py文件中调用add_two函数呢?
    • 我认为您需要在创建 zip 存档时添加 -r 选项,它将保留项目结构 - 请参阅更新的答案。
    【解决方案2】:

    压缩依赖 -

    cd base-path-to-python-modules
    zip -qr deps.zip ./* -x script.py
    

    将 deps.zip 复制到 hdfs/gs。提交作业时使用uri,如下所示。

    使用 Dataproc 的 Python 连接器提交一个 Python 项目 (pyspark)

    from google.cloud import dataproc_v1
    from google.cloud.dataproc_v1.gapic.transports import (
        job_controller_grpc_transport)
    
    region = <cluster region>
    cluster_name = <your cluster name>
    project_id = <gcp-project-id>
    
    job_transport = (
        job_controller_grpc_transport.JobControllerGrpcTransport(
            address='{}-dataproc.googleapis.com:443'.format(region)))
    dataproc_job_client = dataproc_v1.JobControllerClient(job_transport)
    
    job_file = <gs://bucket/path/to/main.py or hdfs://file/path/to/main/job.py>
    
    # command line for the main job file
    args = ['args1', 'arg2']
    
    # required only if main python job file has imports from other modules
    # can be one of .py, .zip, or .egg. 
    addtional_python_files = ['hdfs://path/to/deps.zip', 'gs://path/to/moredeps.zip']
    
    job_details = {
        'placement': {
            'cluster_name': cluster_name
        },
        'pyspark_job': {
            'main_python_file_uri': job_file,
            'args': args,
            'python_file_uris': addtional_python_files
        }
    }
    
    res = dataproc_job_client.submit_job(project_id=project_id,
                                         region=region, 
                                         job=job_details)
    job_id = res.reference.job_id
    
    print(f'Submitted dataproc job id: {job_id}')
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-12-31
      • 1970-01-01
      • 1970-01-01
      • 2020-02-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多