PySpark作业调度:Dataproc集群按需创建技术咨询
解决PySpark作业调度:自动创建+提交+销毁GCP Dataproc集群
嘿,我之前也碰到过一模一样的需求——每月跑一次PySpark作业,总让Dataproc集群一直运行太浪费钱了!刚好可以用GCP的Python客户端库把「创建集群→提交作业→销毁集群」整个流程串起来,给你一步步拆解:
第一步:准备依赖和权限
首先确保你安装了GCP Dataproc的Python客户端:
pip install google-cloud-dataproc google-auth
另外,你的服务账号需要拥有Dataproc Editor(或者更细粒度的权限:dataproc.clusters.create、dataproc.jobs.submit、dataproc.clusters.delete)以及GCS的读写权限(如果你的作业脚本/依赖存在GCS上)。
第二步:编写完整流程代码
下面是整合了集群创建、作业提交、集群销毁的完整示例脚本,你可以根据自己的需求调整参数:
from google.cloud import dataproc_v1 from google.cloud.dataproc_v1.gapic.transports import ( cluster_controller_grpc_transport, job_controller_grpc_transport, ) import time # 配置你的GCP参数 PROJECT_ID = "your-project-id" REGION = "us-central1" # 换成你的集群区域 CLUSTER_NAME = f"monthly-pyspark-cluster-{int(time.time())}" # 用时间戳保证集群名唯一 GCS_PYSPARK_SCRIPT = "gs://your-bucket/path/to/your/script.py" # 你的PySpark脚本路径 JOB_ARGS = ["arg1", "arg2"] # 脚本需要的参数(可选) def create_dataproc_cluster(): # 初始化集群控制器 transport = cluster_controller_grpc_transport.ClusterControllerGrpcTransport( address=f"{REGION}-dataproc.googleapis.com:443" ) client = dataproc_v1.ClusterControllerClient(transport=transport) # 定义集群配置(按需调整机器类型、数量) cluster_config = { "master_config": { "num_instances": 1, "machine_type_uri": "n1-standard-2" }, "worker_config": { "num_instances": 2, "machine_type_uri": "n1-standard-2" }, "software_config": { "image_version": "2.1-debian11" # 选合适的Dataproc镜像版本 } } # 发送创建集群请求 operation = client.create_cluster( request={ "project_id": PROJECT_ID, "region": REGION, "cluster_name": CLUSTER_NAME, "cluster": cluster_config } ) print(f"正在创建集群 {CLUSTER_NAME}...") # 等待集群创建完成 operation.result() print(f"集群 {CLUSTER_NAME} 创建完成!") def submit_pyspark_job(): # 初始化作业控制器 transport = job_controller_grpc_transport.JobControllerGrpcTransport( address=f"{REGION}-dataproc.googleapis.com:443" ) client = dataproc_v1.JobControllerClient(transport=transport) # 定义PySpark作业配置 job = { "placement": {"cluster_name": CLUSTER_NAME}, "pyspark_job": { "main_python_file_uri": GCS_PYSPARK_SCRIPT, "args": JOB_ARGS } } # 提交作业 operation = client.submit_job( request={ "project_id": PROJECT_ID, "region": REGION, "job": job } ) print("PySpark作业已提交,等待执行完成...") # 等待作业完成 response = operation.result() if response.status.state == dataproc_v1.JobStatus.State.SUCCESS: print("作业执行成功!") else: print(f"作业执行失败:{response.status.details}") # 这里可以根据需求抛出异常或者做其他处理 raise Exception(f"Job failed: {response.status.details}") def delete_dataproc_cluster(): # 初始化集群控制器 transport = cluster_controller_grpc_transport.ClusterControllerGrpcTransport( address=f"{REGION}-dataproc.googleapis.com:443" ) client = dataproc_v1.ClusterControllerClient(transport=transport) # 发送删除集群请求 operation = client.delete_cluster( request={ "project_id": PROJECT_ID, "region": REGION, "cluster_name": CLUSTER_NAME } ) print(f"正在销毁集群 {CLUSTER_NAME}...") operation.result() print(f"集群 {CLUSTER_NAME} 已销毁!") if __name__ == "__main__": try: # 执行完整流程 create_dataproc_cluster() submit_pyspark_job() finally: # 不管作业成功还是失败,都销毁集群 delete_dataproc_cluster()
第三步:关键注意事项
- 集群名称唯一性:用时间戳生成集群名,避免重复创建导致报错。
- 异常处理:用
try-finally确保不管作业成功还是失败,集群都会被销毁,不会遗留资源。 - 镜像版本:选择和你的PySpark版本兼容的Dataproc镜像,比如
2.1-debian11对应PySpark 3.3。 - 资源配置:根据你的作业规模调整机器类型和数量,避免过度配置浪费钱。
第四步:触发调度
既然作业是每月执行一次,你可以用Cloud Scheduler来触发App Engine上的这个脚本:
- 创建一个Cloud Scheduler任务,选择HTTP目标,指向你的App Engine服务URL。
- 设置调度频率为
0 0 1 * *(每月1号凌晨执行,可按需调整)。
这样整个流程就完全自动化了——每月自动创建集群、跑作业、销毁集群,不用手动干预,还能省不少资源费用!
内容的提问来源于stack exchange,提问作者CARREAU Clément
相关产品推荐
相关产品推荐

