如何在 Dataproc 上用 BigQuery connector 把 PySpark 结果写入 BigQuery
2026/9/12 5:16:55 网站建设 项目流程

如何在 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 Storage(GCS),在 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 角色,否则会报权限错误;
  • 本地终端可用gsutilgcloud

下面所有命令中的 bucket 名、集群名、临时桶名都是课程环境的示例值,换成你自己环境里对应的值即可;命令结构和参数保持不变。

把脚本的输出从 parquet 改成 BigQuery 表

原始脚本是 code/06_spark_sql.py,读取 green 和 yellow 两类 parquet、合并后按区域/月份/服务类型做收入聚合,结尾写 parquet:

df_result.coalesce(1) \ .write.parquet(output, mode='overwrite')

课程的做法是复制一份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 还需要一个临时 bucket:Spark 先把结果写到 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,作业类型选 PySpark,Main Python file 填上传后的脚本路径,参数给三个:两个 parquet 输入和现在指向 BigQuery 表的输出:

  • --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=trips_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 \ --cluster=de-zoomcamp-cluster \ --region=europe-west6 \ --jars=gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar \ gs://dtc_data_lake_de-zoomcamp-nytaxi/code/06_spark_sql_big_query.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=trips_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-2020,preview 里就是集群上算出来的每月收入行。

文档给出的该表 schema 前几列(作为示例结果展示,你跑出来的数据量不同,但列结构应一致):

FieldTypeMode
revenue_zoneINTEGERNULLABLE
revenue_monthTIMESTAMPNULLABLE
service_typeSTRINGREQUIRED
revenue_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询