You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 02:32:30