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

发布时间:2026/9/12 17:03:37
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-west6Zurich。已有一个可直接运行的 PySpark 脚本课程示例是 06_spark_sql.py它接收--input_green、--input_yellow、--output三个参数读取两份 parquet 数据按月、区域、车型统计收入再把结果写成 parquet 到--output指定的路径。终端上装有 Google Cloud SDKgcloud和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-clusterRegion and zone选你的 bucket 所在的区域让集群离数据近。文档中的 bucket 在europe-west6Zurich。Cluster type实际生产中通常用 standard一个 master 加若干 worker。课程因为只是做实验且数据量不大用 single node 就够了。Additional components勾选 Jupyter notebook 组件可以在集群上直接做实验再勾选 Docker文档中后文的课程单元会用到。master 和 worker 的机型保持默认点 Create。几分钟后集群进入 running 状态。Dataproc 背后为它创建了一台虚拟机可以在 Compute Engine 里看到但不需要去连它——只向集群提交作业即可。上传作业脚本到 GCSDataproc 需要把作业脚本放在它能读到的地方课程做法是上传到 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 typePySparkMain Python filegs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql.py不需要依赖项也不需要 jar 文件Arguments脚本的三个参数输入和输出都指向 bucket--input_greengs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2021/*/--input_yellowgs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2021/*/--outputgs://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 \ --clusterde-zoomcamp-cluster \ --regioneurope-west6 \ gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql.py \ -- \ --input_greengs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2020/*/ \ --input_yellowgs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2020/*/ \ --outputgs://dtc_data_lake_de-zoomcamp-nytaxi/report-2020命令结构值得注意--双横线之前配置的是提交本身集群、区域、脚本之后的部分原样传给作业和本地spark-submit时把参数传给脚本的方式完全一致。这条命令跑的是 2020 年的数据和 Web UI 里跑的 2021 年互为参数化验证。权限报错的处理文档中记录的失败现象第一次运行上面的gcloud命令会报 permission deniednot 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),仅供参考