☰
使用 Python 客户端库操作 Cloud Dataproc API:集群管理与 PySpark 任务实战指南
2026/10/3 2:25:03 网站建设 项目流程
  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

本篇技术指南以当前仓库 dataproc/snippets 目录下的 Cloud Dataproc API 示例程序为核心,带你掌握如何用google-cloud-dataprocPython 客户端库完成集群列举、集群创建、PySpark 任务提交、结果回读与集群销毁的完整闭环。读完本文,你将能直接运行仓库中的命令行示例,理解其底层 API 调用链(ClusterController / JobController / WorkflowTemplateService),并掌握内联工作流模板等进阶用法,从而在自己的 GCP 项目中复现"一键提交 Spark 任务"的自动化方案。

概述:仓库中提供了哪些 Dataproc 示例

该目录下的示例程序(dataproc/snippets)是一组可直接运行的命令行工具,用于与 Cloud Dataproc API 交互,核心脚本及职责如下:

脚本功能说明
list_clusters.py简单命令行程序,演示连接 Cloud Dataproc API 并列出指定区域内的集群
submit_job_to_cluster.py创建集群、提交pyspark_sort.py任务、从 Google Cloud Storage 下载任务输出并打印结果
single_job_workflow.py使用 Cloud Dataproc InstantiateInlineWorkflowTemplate API,用一次 API 请求完成"创建临时集群 → 运行任务 → 删除集群"
pyspark_sort.py被上述脚本上传并在集群内运行的 PySpark 排序任务
pyspark_sort_gcs.py与pyspark_sort.py功能相同,但演示如何从 GCS bucket 读取数据
instantiate_inline_workflow_template.py独立的 Teragen/Terasort 内联工作流示例,演示多步骤任务编排
submit_job.py向既有集群提交任务的精简示例,配套测试 submit_job_test.py

此外,仓库还提供了完整的测试用例(create_cluster_test.py、submit_job_test.py、instantiate_inline_workflow_template_test.py 等)以及一份面向 Cloud Shell 的交互式演练文档 python-api-walkthrough.md,后者详细记录了从创建 bucket 到查看任务输出的端到端流程。

需要说明的是,本目录示例虽然通过 Dataproc API 实现功能,但同样的能力也可以使用 Cloud Console 或gcloudCLI 完成——选择 Python 客户端库的意义在于流程可编程、可自动化、可复用。

本地运行前置条件

在本地运行这些示例之前,需要准备以下环境:

  1. 安装 pip(Python 包管理器),并建议使用 virtualenv 创建隔离的 Python 虚拟环境。
  2. 启用 Dataproc API:进入 Google Cloud Console 的 API Manager(API 与服务),搜索 "Google Cloud Dataproc API" 并启用。若走完整演练流程,还需要启用 Compute Engine 与 Cloud Storage 相关 API,walkthrough 文档中给出了对应的 gcloud 命令(python-api-walkthrough.md):
gcloud services enable dataproc.googleapis.com \ compute.googleapis.com \ storage-component.googleapis.com \ --project=<your-project-id>

安装依赖

在虚拟环境中执行:

pip install -r requirements.txt

依赖清单见 requirements.txt,其中核心依赖为google-cloud-dataproc==5.20.0(提供 Dataproc v1 客户端),此外还包含google-cloud-storage==2.9.0(用于上传/下载 PySpark 文件与任务输出)、google-auth==2.38.0(认证)、grpcio==1.74.0(gRPC 传输层)等。测试依赖见 requirements-test.txt,包含pytest==9.0.3与pytest-xdist==3.3.0。

认证方式

所有示例都通过 Google Cloud 认证 机制访问 API。推荐方式是为运行环境配置一个带有 JSON 密钥的服务账号(Service Account),并通过环境变量GOOGLE_APPLICATION_CREDENTIALS指向该 JSON 密钥文件;也可以先在本地执行gcloud auth application-default login获取应用默认凭据。客户端库会自动拾取这些凭据,示例代码本身无需显式传入任何密钥。

环境变量

仓库示例依赖以下环境变量(值按需替换):

GOOGLE_CLOUD_PROJECT=your-project-id REGION=us-central1 # or your region CLUSTER_NAME=waprin-spark7 ZONE=us-central1-b

其中REGION对应集群所在的区域(Dataproc 集群位于某个区域或global),ZONE为具体可用区(如us-central1-b),CLUSTER_NAME为示例中的默认集群名(实际运行时通常替换为你自己的集群名)。

运行 list_clusters.py:列举区域内的集群

list_clusters.py是最简单的入门示例,用于演示连接 Dataproc API 并列出指定区域内的集群。运行方式:

python list_clusters.py --project_id=$GOOGLE_CLOUD_PROJECT --region=$REGION

从源码看(list_clusters.py),脚本核心逻辑是构建dataproc_v1.ClusterControllerClient,且有一个关键设计——区域端点(regional endpoint)选择:

  • 当--region=global时,使用默认的 gRPC 全局端点,直接实例化ClusterControllerClient();
  • 其他区域时,则通过client_options={"api_endpoint": f"{region}-dataproc.googleapis.com:443"}绑定到对应区域端点。

随后调用客户端的list_clusters方法(list_clusters.py):

def list_clusters(dataproc, project, region): """List the details of clusters in the region.""" for cluster in dataproc.list_clusters( request={"project_id": project, "region": region} ): print(f"{cluster.cluster_name} - {cluster.status.state.name}")

request字典中只需提供project_id与region两个字段,返回的集群对象包含cluster_name与status.state.name(如RUNNING、STOPPED、ERROR等集群状态),逐条打印即完成列举。

运行 submit_job_to_cluster.py:完整的任务提交闭环

submit_job_to_cluster.py是目录中最具代表性的示例,它演示了 Dataproc 数据处理的完整链路:创建集群 → 上传 PySpark 脚本到 GCS → 提交任务 → 从 GCS 回读任务输出 → 删除集群。

前置准备:创建 GCS bucket

Cloud Dataproc 依赖 GCS bucket 来暂存文件,请先在 Cloud Console 或通过 gsutil 创建:

gcloud storage buckets create gs://<your-staging-bucket-name>

然后设置环境变量:

BUCKET=your-staging-bucket CLUSTER=your-cluster-name

方式一:使用已有集群

python submit_job_to_cluster.py --project_id=$GOOGLE_CLOUD_PROJECT --zone=us-central1-b --cluster_name=$CLUSTER --gcs_bucket=$BUCKET

使用已有集群前,可用 Cloud Console 创建,或运行 gcloud 命令:

gcloud dataproc clusters create your-cluster-name

方式二:脚本自动创建新集群(推荐)

python submit_job_to_cluster.py --project_id=$GOOGLE_CLOUD_PROJECT --zone=us-central1-b --cluster_name=$CLUSTER --gcs_bucket=$BUCKET --create_new_cluster

脚本会完成以下动作:搭建集群 → 上传 PySpark 文件 → 提交任务 → 打印结果 → 如果集群由脚本创建,则在任务结束后删除集群。

可选参数--pyspark_file可将默认的pyspark_sort.py替换为你自己的脚本(见 submit_job_to_cluster.py)。

源码级拆解:脚本内部的关键调用链

① 创建集群(submit_job_to_cluster.py):脚本构建一个 cluster 配置字典,指定主节点与工作节点的实例数与机型:

cluster = { "project_id": project_id, "cluster_name": cluster_name, "config": { "master_config": {"num_instances": 1, "machine_type_uri": "n1-standard-2"}, "worker_config": {"num_instances": 2, "machine_type_uri": "n1-standard-2"}, }, } operation = cluster_client.create_cluster( request={"project_id": project_id, "region": region, "cluster": cluster} ) result = operation.result()

这里主节点 1 台、工作节点 2 台、机型均为n1-standard-2。create_cluster是一个长期运行操作(LRO),operation.result()会阻塞直到集群创建完成。

② 上传 PySpark 文件(submit_job_to_cluster.py):使用google.cloud.storage客户端,将本地 PySpark 文件以 blob 形式上传到指定 bucket;文件默认取自脚本同目录下的pyspark_sort.py(常量DEFAULT_FILENAME,L35)。

③ 提交任务(submit_job_to_cluster.py):创建JobControllerClient,构建 PySpark 任务描述并调用submit_job_as_operation:

job = { "placement": {"cluster_name": cluster_name}, "pyspark_job": {"main_python_file_uri": f"gs://{gcs_bucket}/{spark_filename}"}, } operation = job_client.submit_job_as_operation( request={"project_id": project_id, "region": region, "job": job} ) response = operation.result()

任务输出会写入 Dataproc 分配给任务输出的 GCS 路径,脚本用正则re.match("gs://(.*?)/(.*)", response.driver_output_resource_uri)解析出 bucket 与 blob 前缀,再拼接.000000000后缀下载 driver 输出内容(L131-L139)。

④ 删除集群(submit_job_to_cluster.py):调用delete_cluster并等待操作完成,输出Cluster xxx successfully deleted.。

预期运行输出

参照 walkthrough 文档(python-api-walkthrough.md),一次成功的运行会在终端依次打印集群创建成功、任务完成、排序结果、集群删除成功:

Cluster created successfully: cluster-name. ... Job finished successfully. ... ['Hello,', 'dog', 'elephant', 'panther', 'world!'] ... Cluster cluster-name successfully deleted.

PySpark 任务脚本:pyspark_sort.py 与 pyspark_sort_gcs.py

这两个文件是被上传到集群、在 PySpark 环境中执行的任务脚本(注意:它们不应在本地直接运行,必须在 PySpark 运行时中执行)。

pyspark_sort.py 的核心逻辑极简,创建一个SparkContext,并行化一个单词列表,排序并收集结果打印:

sc = pyspark.SparkContext() rdd = sc.parallelize(["Hello,", "world!", "dog", "elephant", "panther"]) words = sorted(rdd.collect()) print(words)

其输出即前文示例中的['Hello,', 'dog', 'elephant', 'panther', 'world!']。

pyspark_sort_gcs.py 则演示了从 GCS 读取数据的另一种用法,通过sc.textFile("gs://path-to-your-GCS-file")直接以 GCS URI 读取文件内容后再排序:

sc = pyspark.SparkContext() rdd = sc.textFile("gs://path-to-your-GCS-file") print(sorted(rdd.collect()))

两者结合使用--pyspark_file参数即可切换,适用于"脚本内嵌数据"与"数据存于 GCS"两种场景。

single_job_workflow.py:用内联工作流模板一次请求完成全流程

single_job_workflow.py演示 Cloud Dataproc 的InstantiateInlineWorkflowTemplateAPI——将集群配置、任务定义打包进一个内联工作流模板,一次 API 请求即可自动完成"创建临时集群 → 运行任务 → 删除集群"。

运行方式(single_job_workflow.py):

python single_job_workflow.py --project_id=$PROJECT --gcs_bucket=$BUCKET \ --cluster_name=$CLUSTER --zone=$ZONE

若集群位于全局区域(global region),可附加--global_region参数;默认情况下脚本会从--zone推导所在区域(get_region_from_zone取 zone 的前缀段,如us-central1-b→us-central1)。

从源码看(single_job_workflow.py),工作流模板数据由两部分组成:

  • placement.managed_cluster:声明由工作流托管的临时集群,包含zone_uri、主节点 1 台n1-standard-1、工作节点 2 台n1-standard-1;
  • jobs:一个 PySpark 任务列表,main_python_file_uri指向已上传到 GCS 的脚本,step_id为pyspark-job。

然后调用:

workflow = dataproc.instantiate_inline_workflow_template( request={"parent": parent, "template": workflow_data} ) workflow.add_done_callback(callback)

脚本通过add_done_callback注册完成回调,并用全局标志waiting_callback配合轮询循环等待工作流结束(wait_for_workflow_end),同时提示可在 Cloud Console 的 Dataproc Workflows 页面查看工作流与任务进度。

更复杂的内联工作流:Teragen + Terasort 多步骤编排

仓库还提供了 instantiate_inline_workflow_template.py,展示一个包含两个 Hadoop 步骤的内联工作流:先用teragen生成 1000 条数据到 HDFS,再用terasort对其排序。其关键点在于步骤依赖关系——terasort通过prerequisite_step_ids: ["teragen"]声明依赖前一步(L54-L62),并且同样采用managed_cluster自动托管的集群(zone_uri留空即自动选择可用区)。运行命令:

python instantiate_inline_workflow_template.py <PROJECT_ID> <REGION>

成功时输出Workflow ran successfully.。这个示例说明内联工作流不仅适用于 PySpark 单任务,也支持 Hadoop 多步骤的复杂编排,是理解 Dataproc Workflow 编排能力的最佳入口。

区域端点与 global 区域的选择

综合上述示例源码可以发现一个共性设计:Dataproc Python 客户端对区域端点非常敏感。

  • list_clusters.py:--region=global走默认全局端点,其他区域拼装{region}-dataproc.googleapis.com:443;
  • submit_job_to_cluster.py:直接要求传入--region,客户端始终使用区域端点;
  • single_job_workflow.py:提供--global_region开关,开启时使用WorkflowTemplateServiceClient()默认全局端点,关闭时通过WorkflowTemplateServiceGrpcTransport指定区域地址(L143-L156)。

实际使用时应根据集群所在区域选择匹配的端点;如果集群位于某个具体区域而客户端使用 global 端点,可能无法列出或访问该集群。测试用例(submit_job_test.py)也印证了这一点——fixture 统一使用us-central1-dataproc.googleapis.com:443区域端点创建客户端。

测试与验证:仓库如何保障示例可用

仓库为上述示例配备了集成测试,既是质量保障,也是理解 API 调用细节的辅助材料。以 submit_job_test.py 为例:

  • 测试自动创建集群(配置中还加入了boot_disk_size_gb: 100的启动盘设置,L40-L55);
  • 提交任务前通过get_cluster断言集群处于ClusterStatus.State.RUNNING状态(L109-L112),避免在集群错误状态下重试导致无效调用;
  • 使用backoff库对ServiceUnavailable、InternalServerError等瞬态错误做指数退避重试(最多 5 次);
  • 断言脚本输出包含"Job finished successfully",并在 finally 中清理集群。

测试的 Python 版本约束与项目环境变量配置见 noxfile_config.py:默认忽略 3.8/3.9/3.11/3.12/3.13 版本(即主要在 Python 3.10 上运行),项目 ID 通过GOOGLE_CLOUD_PROJECT环境变量注入。

清理资源与收尾建议

  • 使用--create_new_cluster模式时,脚本会自动删除自己创建的集群;使用已有集群模式则不会删除,任务结束后可按需手动清理。
  • 如果为演练创建了独立的 GCS bucket,可待 bucket 为空后删除:
gcloud storage buckets delete gs://$BUCKET

如需连同其中对象一并删除且确认数据可丢弃:

gcloud storage rm --recursive gs://$BUCKET
  • 任务详情(如执行日志、耗时、资源配置)可在 Cloud Console 的 Dataproc Jobs 页面按 PySpark 任务名查看。

小结

本目录示例完整覆盖了 Dataproc Python 客户端库的主要使用模式:ClusterControllerClient负责集群生命周期(创建/列举/删除),JobControllerClient负责任务提交与输出回读,WorkflowTemplateServiceClient通过内联工作流模板实现临时集群的自动化编排。配合 quickstart 与 python-api-walkthrough.md 的端到端演练,你可以从零开始,在几分钟内跑通"创建集群 → 提交 PySpark 任务 → 获取结果 → 释放资源"的完整流程,并将其沉淀为可复用的自动化脚本。

  • 示例工程

【免费下载链接】python-docs-samples

Code samples used on cloud.google.com

项目地址:https://gitcode.com/GitHub_Trending/py/python-docs-samples
点击查看免费下载

相关推荐

上一篇:VC++ 运行库修复:一条命令搞定 2005 到 2022 全版本
下一篇:dep init 遇到损坏 glide.yaml 的容错处理:配置导入失败后的依赖求解回退路径解析

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询