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

如何在Google Cloud Shell提交PySpark作业、传递并读取文件与参数

在Google Cloud Shell中处理PySpark作业:提交、传参与文件读取

1. 在Google Cloud Shell中提交PySpark作业

Google Cloud Shell内置gcloud命令行工具,可直接提交PySpark作业到Dataproc集群(需提前创建集群):

基础提交命令

  • 脚本存储在Cloud Shell本地:
gcloud dataproc jobs submit pyspark ./your_script.py \
    --cluster=your-cluster-name \
    --region=your-region
  • 脚本存储在GCS(Google Cloud Storage):
gcloud dataproc jobs submit pyspark gs://your-bucket/path/to/script.py \
    --cluster=your-cluster-name \
    --region=your-region

前置:创建Dataproc集群(若未创建)

gcloud dataproc clusters create your-cluster-name \
    --region=your-region \
    --num-workers=2

2. 在PySpark提交命令中传递文件与参数

传递依赖文件

  • 用--files传递单个/多个普通文件(如配置、辅助脚本),多文件用逗号分隔:
gcloud dataproc jobs submit pyspark ./your_script.py \
    --cluster=your-cluster-name \
    --region=your-region \
    --files ./config.json,./helper.py
  • 用--py-files传递Python模块包(.zip/.egg格式):
gcloud dataproc jobs submit pyspark ./your_script.py \
    --cluster=your-cluster-name \
    --region=your-region \
    --py-files ./utils.zip

传递命令行参数

用--分隔gcloud自身参数与作业参数,之后直接写入参数(支持键值对或纯参数):

gcloud dataproc jobs submit pyspark ./your_script.py \
    --cluster=your-cluster-name \
    --region=your-region \
    -- --input=gs://your-bucket/input --output=gs://your-bucket/output batch-size=1000

3. 在PySpark代码中读取文件与参数

读取传递的文件

--files/--py-files传递的文件会自动分发到每个节点的工作目录,直接用相对路径访问:

# 读取配置文件
import json
with open("config.json", "r") as f:
    config = json.load(f)

# 导入传递的Python模块
from helper import calculate_metrics

读取命令行参数

用Python标准库sys.argv获取参数,跳过第一个元素(脚本自身路径):

import sys
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("JobExample").getOrCreate()

# 解析参数
args = sys.argv[1:]
input_path = None
output_path = None
batch_size = 100

for arg in args:
    if arg.startswith("--input="):
        input_path = arg.split("=")[1]
    elif arg.startswith("--output="):
        output_path = arg.split("=")[1]
    elif "batch-size=" in arg:
        batch_size = int(arg.split("=")[1])

# 使用参数处理数据
df = spark.read.parquet(input_path)
processed_df = df.limit(batch_size)
processed_df.write.parquet(output_path)

内容的提问来源于stack exchange,提问作者Rakesh Reddy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 08:56:06