- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
本篇技术指南以当前仓库 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 客户端库的意义在于流程可编程、可自动化、可复用。
本地运行前置条件
在本地运行这些示例之前,需要准备以下环境:
- 安装 pip(Python 包管理器),并建议使用 virtualenv 创建隔离的 Python 虚拟环境。
- 启用 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
相关推荐
mpv Lua脚本三步上手指南:4个官方插件驯服黑边、音量与续播
mpv Lua脚本三步上手指南:4个官方插件驯服黑边、音量与续播 黑边裁不掉、播完不接着播、大音量吓人、暂停时窗口挡视线,这些小事都能被mpv自带的官方Lua脚
音视频视频音频VideoRAG核心技术揭秘:双渠道架构如何颠覆视频理解
VideoRAG核心技术揭秘:双渠道架构如何颠覆视频理解 在信息爆炸的时代,视频已成为知识传递的重要载体,但传统视频理解技术往往受限于上下文长度和多模态信息融合
人工智能RAG多模态视频知识图谱AI 应用本地部署KafkaJS管理客户端实战:全面掌握集群运维操作
KafkaJS管理客户端实战:全面掌握集群运维操作 KafkaJS是一个现代化的Apache Kafka客户端库,专为Node.js环境设计,提供强大的管理客户
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考