news 2026/9/12 17:03:25

Data Engineering Zoomcamp 如何创建 GCP Dataproc 集群并提交 PySpark 作业

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Data Engineering Zoomcamp 如何创建 GCP Dataproc 集群并提交 PySpark 作业

Data Engineering Zoomcamp 如何创建 GCP Dataproc 集群并提交 PySpark 作业

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

这篇文章解决的任务是:把已经跑通的 PySpark 作业从本地 Spark 移到 Google Cloud 上,具体动作是——在 GCP 控制台创建一个 Dataproc 集群,然后把处理纽约出租车数据的 PySpark 脚本作为 PySpark 作业提交上去,先通过 Web UI 提交,再用gcloudSDK 从终端提交。内容取自 Data Engineering Zoomcamp 2027 届 batch 模块的第 15 单元(15-setting-up-a-dataproc-cluster.md)及其前置单元。

Dataproc 是 Google Cloud 的托管 Spark 服务:它负责创建 master 和 worker,你只负责向集群提交作业。适用前提是:

  • 已有一个 GCS bucket,里面存着 parquet 格式的出租车数据。课程示例使用的 bucket 是dtc_data_lake_de-zoomcamp-nytaxi,数据在pq/green/pq/yellow/下,区域为europe-west6(Zurich)。
  • 已有一个可直接运行的 PySpark 脚本,课程示例是 06_spark_sql.py:它接收--input_green--input_yellow--output三个参数,读取两份 parquet 数据,按月、区域、车型统计收入,再把结果写成 parquet 到--output指定的路径。
  • 终端上装有 Google Cloud SDK(gcloud)和gsutil

有一个重要的简化:Dataproc 集群天生能访问 GCS,不需要第 13 单元(13-connecting-to-google-cloud-storage.md)里那套 GCS connector for Hadoop 的 jar 和认证配置——那套配置只在你自己的机器或未经配置的虚拟机上跑 Spark 时才需要。

在 GCP 控制台创建集群

在 Google Cloud 控制台打开 Dataproc。第一次进入时会要求启用 API,点一下即可。然后点 Create cluster,按课程文档的配置填写:

  • Namede-zoomcamp-cluster
  • Region and zone:选你的 bucket 所在的区域,让集群离数据近。文档中的 bucket 在europe-west6(Zurich)。
  • Cluster type:实际生产中通常用 standard(一个 master 加若干 worker)。课程因为只是做实验且数据量不大,用 single node 就够了。
  • Additional components:勾选 Jupyter notebook 组件,可以在集群上直接做实验;再勾选 Docker,文档中后文的课程单元会用到。

master 和 worker 的机型保持默认,点 Create。几分钟后集群进入 running 状态。Dataproc 背后为它创建了一台虚拟机,可以在 Compute Engine 里看到,但不需要去连它——只向集群提交作业即可。

上传作业脚本到 GCS

Dataproc 需要把作业脚本放在它能读到的地方,课程做法是上传到 bucket 的code目录(生产上一般会用单独的 bucket 放代码,这里为了简单和共用一个):

gsutil -m cp -r 06_spark_sql.py gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql.py

这里有一个前置条件容易忽略:脚本不能写死 master。第 14 单元(14-creating-a-local-spark-cluster.md)中本地运行时脚本里有.master("local[*]"),而 06_spark_sql.py 里已经把它删掉了,只保留SparkSession.builder.appName('test').getOrCreate()。文档明确说明这一点很关键:Dataproc 运行作业时由它自己设置 master。

通过 Web UI 提交 PySpark 作业

打开刚创建的集群,点 Submit job,按下面内容填表单:

  • Job type:PySpark
  • Main Python filegs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql.py
  • 不需要依赖项,也不需要 jar 文件
  • Arguments:脚本的三个参数,输入和输出都指向 bucket:
    • --input_green=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2021/*/
    • --input_yellow=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2021/*/
    • --output=gs://dtc_data_lake_de-zoomcamp-nytaxi/report-2021

提交后等待。作业页面在运行期间会显示 driver 输出。成功条件是:作业结束后,bucket 中出现report-2021目录,里面是月度收入报告的 parquet 文件——这些数据就是刚创建的集群算出来的。

用 gcloud SDK 提交作业

Web UI 适合试验,但不适合生产:从 Airflow 之类的调度器里没法靠点按钮提交作业。向 Dataproc 提交作业有三种方式:Web UI、Google Cloud SDK 和 REST API。作业的详情页面会展示刚才那次提交的等价 REST 调用,可以从里面读出关键要素:集群名、Python 文件、参数。

用 SDK 时,在终端执行:

gcloud dataproc jobs submit pyspark \ --cluster=de-zoomcamp-cluster \ --region=europe-west6 \ gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql.py \ -- \ --input_green=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2020/*/ \ --input_yellow=gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2020/*/ \ --output=gs://dtc_data_lake_de-zoomcamp-nytaxi/report-2020

命令结构值得注意:--双横线之前配置的是提交本身(集群、区域、脚本),之后的部分原样传给作业,和本地spark-submit时把参数传给脚本的方式完全一致。这条命令跑的是 2020 年的数据,和 Web UI 里跑的 2021 年互为参数化验证。

权限报错的处理

文档中记录的失败现象:第一次运行上面的gcloud命令会报 permission denied(not authorized to request the resource)。原因是课程全程共用一个服务账号(模块一 Terraform 创建的那个),它没有提交 Dataproc 作业的权限。

文档给出的修复方式:打开 IAM & admin,找到该服务账号,添加Dataproc Administrator角色。文档同时说明:真实项目里应该拆分角色——给 Terraform 一个权限较大的角色,给 worker 和调度器一个只允许提交 Dataproc 作业、访问 bucket 等必要操作的窄角色;课程为了简单才在同一个账号上加角色。

更新策略后重新运行同一条gcloud命令即可通过:作业提交成功,终端显示与 Web UI 中相同的 driver 输出并正常结束。

结果验证

文档给出的验证方式是检查 bucket 里的输出:再为不同年份运行一次命令,可以看到 2020 的报告和 2021 的报告并排落在 bucket 中,两份都是集群算出的 parquet 报告。

后续衔接

文档指出的两个下一步:

  • 把这条链路接入 Airflow 时,最简单的做法是用一个 BashOperator 原样执行上面的gcloud命令,当然也有专门的 Dataproc operator 可用。
  • 如果希望结果落到数据仓库而不是 bucket(例如为了做仪表盘),第 16 单元(16-connecting-spark-to-bigquery.md)讲 Spark 直接写 BigQuery 的方式,比在 parquet 上建外部表再拷贝更直接。

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 16:57:47

ESP32蓝牙Beacon测距实战:RSSI转距离原理与校准滤波指南

搞 ESP32 开发的朋友应该都有这种感觉,从板子到手到真正让它按自己的想法跑起来,中间要翻过的山头真不少。之前几讲把环境搭建、GPIO、Wi-Fi 联网这些基础打完后,很多人在蓝牙这块又卡住了,尤其是想用蓝牙做点实际应用的时候&…

作者头像 李华
网站建设 2026/9/12 16:57:37

雷达Simulink仿真中的RF前端行为级建模与参数化实践

简介:这份资源面向雷达系统设计工程师与Simulink建模学习者,围绕射频前端行为在雷达系统级仿真中的整合展开,提供单站脉冲雷达目标探测与FMCW雷达距离/速度估计两套完整可运行模型。资源共9个文件,其中5个m脚本用于参数配置与仿真…

作者头像 李华
网站建设 2026/9/12 16:57:25

ESP32蓝牙Beacon测距实战:RSSI原理、代码实现与参数标定

1. 蓝牙beacon测距的总体设计思路开讲之前先说个事:这个系列的每一讲,我都在尽量控制篇幅,结果每讲写出来都能赶上小论文。倒不是因为废话多,而是ESP-IDF这套框架里,很多看起来“一行搞定”的功能,真要讲清…

作者头像 李华