ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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),仅供参考
返回列表