使用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

