You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Apache Spark读取BigQuery表中BIGNUMERIC类型时遇报错

问题

在Dataproc上运行Spark任务,通过Spark-BigQuery连接器读取包含BIGNUMERIC类型列的BigQuery表时,执行df.columns获取列名抛出ModuleNotFoundError: No module named 'google.cloud.spark',但跳过df.columns直接执行df.show()和df.printSchema()可正常运行,schema显示bignumeric类型列存在。

复现代码:

df = spark.read.format('bigquery').load('project_id.dataset_id.table_id')
columns = df.columns
print(f'*********Columns - {columns}**********')
df.show()
df.printSchema()

报错信息:

columns = df.columns() File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/dataframe.py", line
939, in columns File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/dataframe.py", line
256, in schema File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
871, in _parse_datatype_json_string File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
888, in _parse_datatype_json_value File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
577, in fromJson File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
577, in File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
434, in fromJson File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
890, in _parse_datatype_json_value File
"/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/types.py", line
736, in fromJson ModuleNotFoundError: No module named
'google.cloud.spark'

printSchema输出:

root
|-- col1: string (nullable = true)
|-- col2: bignumeric (nullable = true)

原因分析

Spark-BigQuery连接器将BigQuery的BIGNUMERIC类型映射为自定义类型google.cloud.spark.bigquery.types.BigNumeric,但PySpark在解析DataFrame schema(df.columns会触发schema的完整解析流程)时,需要加载该自定义类型对应的Python模块,而Dataproc默认环境中缺少这个依赖。

解决方案

1. 安装缺失的Python依赖

在Dataproc集群所有节点上安装google-cloud-spark包:

pip install google-cloud-spark

若使用Dataproc初始化动作,可将上述命令加入初始化脚本中,确保集群启动时自动安装。

2. 提交作业时指定连接器依赖包

提交Spark作业时,通过--packages参数指定包含BigNumeric类型定义的连接器完整包(注意版本与Spark版本兼容):

spark-submit \
  --packages com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.36.0 \
  your_spark_job.py

替换0.36.0为你实际使用的Spark-BigQuery连接器版本,Spark 3.x对应_2.12后缀,Spark 2.x对应_2.11后缀。

3. 绕过schema解析,直接通过BigQuery API获取列名

如果不想安装额外依赖,可通过BigQuery Python客户端直接读取表结构获取列名,再用于后续操作:

from google.cloud import bigquery

client = bigquery.Client()
table = client.get_table('project_id.dataset_id.table_id')
columns = [field.name for field in table.schema]
print(f'*********Columns - {columns}**********')

# 正常读取DataFrame
df = spark.read.format('bigquery').load('project_id.dataset_id.table_id')
df.show()
df.printSchema()

4. 将BIGNUMERIC映射为Spark原生Decimal类型

通过配置关闭BIGNUMERIC自定义类型支持,将其映射为Spark原生Decimal类型(注意:BigQuery BIGNUMERIC是78位精度,Spark Decimal最大支持38位,会丢失部分精度,仅适用于精度要求不高的场景):

spark.conf.set("spark.sql.bigquery.bignumeric.enabled", "false")
df = spark.read.format('bigquery').load('project_id.dataset_id.table_id')
columns = df.columns
print(f'*********Columns - {columns}**********')

内容的提问来源于stack exchange,提问作者Rajnil Guha

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 05:55:16