
如何在 Dataproc 上用 BigQuery connector 把 PySpark 结果写入 BigQuery【免费下载链接】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在 Data Engineering Zoomcamp 的批处理模块里前面几单元已经完成了这样一条链路把 NYC 出租车 parquet 数据上传到 Google Cloud StorageGCS在 Dataproc 集群上用 PySpark 做聚合把结果写回 GCS 的一个文件夹。这一篇解决最后一步不再写 bucket而是用 BigQuery connector 把 PySpark 的聚合结果直接写入 BigQuery 表让报表可以直接被仪表盘使用。适用前提均来自课程前序单元的设置GCS bucket 里已有按 green/yellow 分目录的 parquet 数据课程示例中位于gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2020/*/与gs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2020/*/Dataproc 集群已创建并运行示例集群名de-zoomcamp-cluster区域europe-west6创建过程见 15-setting-up-a-dataproc-cluster.md提交作业用的服务账号已有 Dataproc Administrator 角色否则会报权限错误本地终端可用gsutil和gcloud。下面所有命令中的 bucket 名、集群名、临时桶名都是课程环境的示例值换成你自己环境里对应的值即可命令结构和参数保持不变。把脚本的输出从 parquet 改成 BigQuery 表原始脚本是 code/06_spark_sql.py读取 green 和 yellow 两类 parquet、合并后按区域/月份/服务类型做收入聚合结尾写 parquetdf_result.coalesce(1) \ .write.parquet(output, modeoverwrite)课程的做法是复制一份06_spark_sql_big_query.py仓库中即 code/06_spark_sql_big_query.py数据读取和聚合 SQL 完全不动只改输出部分df_result.write.format(bigquery) \ .option(table, output) \ .save()两点变化--output不再是一个 GCS 文件夹而是schema.table形式的 BigQuery 表从命令行传入。课程示例中数据集是trips_data_all所以输出为trips_data_all.reports-2020不再需要coalesce(1)去合并输出文件多文件归并交给 BigQuery 处理。connector 还需要一个临时 bucketSpark 先把结果写到 GCS再加载进 BigQuery。脚本里这样配置temporaryGcsBucket是 Spark 配置项值必须是 GCS bucket 名spark.conf.set(temporaryGcsBucket, dataproc-temp-europe-west6-828225226997-fckhkym8)文档中这个值来自课程自己的 Dataproc 集群创建集群时 Dataproc 会自动建一个dataproc-temp-...形式的临时桶这里复用它。换成你自己的集群时用你自己集群对应的临时桶名替换。脚本改好后上传到 bucket供 Dataproc 读取gsutil -m cp -r 06_spark_sql_big_query.py gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql_big_query.py第一次提交作业快速失败报 Failed to find data source: bigquery在 Dataproc UI 里点 submit job作业类型选 PySparkMain Python file 填上传后的脚本路径参数给三个两个 parquet 输入和现在指向 BigQuery 表的输出--input_greengs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2020/*/--input_yellowgs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2020/*/--outputtrips_data_all.reports-2020文档记录的实际结果是作业很快失败驱动输出里是py4j.protocol.Py4JJavaError: An error occurred while calling o119.save. java.lang.ClassNotFoundException: Failed to find data source: bigquery.原因很直接与 GCS 不同BigQuery connector 并不随每个 Dataproc 集群自带集群也不会自动推断你要用它必须显式把 connector jar 提供给作业。用 --jars 指定 connector jar 重新提交修复方式和本地 Spark 挂 GCS connector jar 的思路一致给作业一个 jar。Google 把 BigQuery connector 放在公开 bucketgs://spark-lib/bigquery/下提交时直接用--jars指向它gcloud dataproc jobs submit pyspark \ --clusterde-zoomcamp-cluster \ --regioneurope-west6 \ --jarsgs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar \ gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql_big_query.py \ -- \ --input_greengs://dtc_data_lake_de-zoomcamp-nytaxi/pq/green/2020/*/ \ --input_yellowgs://dtc_data_lake_de-zoomcamp-nytaxi/pq/yellow/2020/*/ \ --outputtrips_data_all.reports-2020命令中--之前的部分配置提交本身集群、区域、jar、脚本--之后的参数原样传给脚本和本地spark-submit的用法一致。关于 jar 版本文档给了两条需要注意的信息最新 Spark 版本与 BigQuery connector 之间可能出现兼容问题。各 Spark 版本对应的 jar 下载链接可以在 Spark BigQuery connector 的仓库GoogleCloudDataproc/spark-bigquery-connector托管在 GitHub里找到Dataproc 2.2 版本起对应 GCE 2.1 及更新的 Dataproc 镜像connector 已预装在新集群上这类集群提交作业时完全不需要--jars参数。验证表出现在 BigQuery 里这次作业正常跑完。一个提交前悬着的问题也被解答如果输出表还不存在会怎样答案是 Spark 会直接创建它。打开 BigQuery 刷新后trips_data_all数据集下出现了reports-2020preview 里就是集群上算出来的每月收入行。文档给出的该表 schema 前几列作为示例结果展示你跑出来的数据量不同但列结构应一致FieldTypeModerevenue_zoneINTEGERNULLABLErevenue_monthTIMESTAMPNULLABLEservice_typeSTRINGREQUIREDrevenue_monthly_fareFLOATNULLABLE限制与收尾写入是两段式的结果先落到temporaryGcsBucket指定的 GCS bucket再加载进 BigQuery所以临时桶配置漏掉或桶不可用都会导致写入失败spark-bigquery-latest_2.12.jar中的_2.12对应 Scala 2.12 构建这是文档使用的 jar 名换 Spark 版本时按 connector 仓库提供的对应版本 jar 选择完成这一步后课程「在云上跑 Spark」的部分就闭环了本地 Spark 连 GCS、创建 Dataproc 集群、用 UI 和gcloud提交作业、结果直写 BigQuery。后续接 Airflow 时最简单的方式就是用一个 BashOperator 跑这条gcloud dataproc jobs submit命令这也是课程文档给出的下一步方向。主要参考文档16-connecting-spark-to-bigquery.md、15-setting-up-a-dataproc-cluster.md、13-connecting-to-google-cloud-storage.md。【免费下载链接】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),仅供参考