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,按课程文档的配置填写:
- Name:
de-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 file:
gs://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),仅供参考