如何解决读取含JSON类型列的BigQuery表到Spark DataFrame的报错?
问题背景
我们可以成功将普通结构的BigQuery表读取为Spark DataFrame,普通表结构如下:
使用的PySpark读取代码如下:
from google.oauth2 import service_account from google.cloud import bigquery import json import base64 as bs from pyspark.sql.types import StructField, StructType, StringType, IntegerType, DoubleType, DecimalType schema = "schema_name" project_id = "project_id" table_name = "simple" # table_name = "jsonres" schema_table_name = str(project_id) + "." + str(schema) + "." + str(table_name) credentials_dict = {"Insert_actual_credentials": "here"} credentials = service_account.Credentials.from_service_account_info(credentials_dict) client = bigquery.Client(credentials=credentials, project=project_id) query = "SELECT * FROM `{}`;".format(schema_table_name) query_job = client.query(query) query_job.result() s = json.dumps(credentials_dict) res = bs.b64encode(s.encode('utf-8')) ans = res.decode("utf-8") try: df = spark.read.format('bigquery') \ .option("credentials", ans) \ .option("parentProject", project_id) \ .option("project", project_id) \ .option("mode", "DROPMALFORMED") \ .option('dataset', query_job.destination.dataset_id) \ .load(query_job.destination.table_id) df.printSchema() print(df) df.show() except Exception as exp: print(exp)
但当BigQuery表包含JSON类型列时(表结构如下),读取会报错:
报错信息
调用o1138.load时发生错误:
java.lang.IllegalStateException: Unexpected type: JSON at
com.google.cloud.spark.bigquery.SchemaConverters.getStandardDataType(SchemaConverters.java:355)
at
com.google.cloud.spark.bigquery.SchemaConverters.lambda$getDataType$3(SchemaConverters.java:303)
已尝试的方法
我们尝试在读取时手动指定Schema:
structureSchema = StructType([ \ StructField('x', StructType([ StructField('name', StringType(), True) ])), StructField("y", DecimalType(), True) \ ]) print(structureSchema) try: df = spark.read.format('bigquery') \ .option("credentials", ans) \ .option("parentProject", project_id) \ .option("project", project_id) \ .option("mode", "DROPMALFORMED") \ .option('dataset', query_job.destination.dataset_id) \ .schema(structureSchema) \ .load(query_job.destination.table_id) df.printSchema() print(df) df.show() except Exception as exp: print(exp)
但仍出现相同的java.lang.IllegalStateException: Unexpected type: JSON报错。
已知问题
GitHub上已有相关开源问题:Spark-BigQuery连接器读取含JSON类型字段的BigQuery表时会抛出异常。
可行替代方案
方案1:查询时将JSON列转为字符串后解析
修改BigQuery查询语句,用TO_JSON_STRING函数将JSON列转为字符串,读取后再用Spark的from_json函数解析为对应结构:
# 修改查询语句,将JSON列转为字符串 query = "SELECT TO_JSON_STRING(x) AS x_str, y FROM `{}`;".format(schema_table_name) query_job = client.query(query) query_job.result() # 读取数据后解析字符串为JSON结构 from pyspark.sql.functions import from_json # 定义JSON对应的Schema json_schema = StructType([StructField('name', StringType(), True)]) df = spark.read.format('bigquery') \ .option("credentials", ans) \ .option("parentProject", project_id) \ .option("project", project_id) \ .option('dataset', query_job.destination.dataset_id) \ .load(query_job.destination.table_id) # 解析字符串为JSON结构并删除临时列 df = df.withColumn('x', from_json(df.x_str, json_schema)).drop('x_str') df.show()
方案2:升级Spark-BigQuery连接器版本
部分新版本的连接器已修复JSON类型支持问题,升级连接器版本后可直接读取。例如在提交Spark任务时指定最新版本的包:
spark-submit --packages com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.36.0 your_script.py
升级后重新运行原读取代码即可。
方案3:导出为Parquet文件后读取
先将BigQuery表导出到GCS的Parquet文件,再用Spark读取Parquet文件,JSON列会自动转为Spark支持的结构类型:
# 将BigQuery表导出到GCS的Parquet文件 destination_uri = "gs://your-bucket/path/to/output.parquet" extract_job = client.extract_table( query_job.destination, destination_uri, format="PARQUET" ) extract_job.result() # Spark读取Parquet文件 df = spark.read.parquet(destination_uri) df.show()
内容的提问来源于stack exchange,提问作者Nandha

