使用PySpark写入BigQuery JSON字段失败,求解决方法
问题背景
使用GCP Dataproc 2.1集群(Spark 3.4、Java 17、Scala 2.13),依赖gs://spark-lib/bigquery/spark-bigquery-latest.jar,尝试将JSON对象写入BigQuery表的JSON类型字段(表仅含一个message字段,类型为JSON),执行代码时触发Unsupported field type: JSON错误。
原代码
rdd1 = spark.read.json("gs://sample_data_x/sample_json/*.json").rdd rdd2 = rdd1.map( lambda x : {"message" : x} ) df = rdd2.toDF() df.printSchema() df.show() df.write.format('bigquery')\ .option('table', ( ''))\ .option("project", "") \ .option("temporaryGcsBucket","")\ .mode("overwrite")\ .save()
错误信息
Caused by: com.google.cloud.spark.bigquery.repackaged.com.google.cloud.bigquery.BigQueryException: Unsupported field type: JSON at com.google.cloud.spark.bigquery.repackaged.com.google.cloud.bigquery.Job.reload(Job.java:419) at com.google.cloud.spark.bigquery.repackaged.com.google.cloud.bigquery.Job.waitFor(Job.java:252) at com.google.cloud.bigquery.connector.common.BigQueryClient.createAndWaitFor(BigQueryClient.java:333) at com.google.cloud.bigquery.connector.common.BigQueryClient.createAndWaitFor(BigQueryClient.java:323) at com.google.cloud.bigquery.connector.common.BigQueryClient.loadDataIntoTable(BigQueryClient.java:564) at com.google.cloud.spark.bigquery.write.BigQueryWriteHelper.loadDataToBigQuery(BigQueryWriteHelper.java:134) at com.google.cloud.spark.bigquery.write.BigQueryWriteHelper.writeDataFrameToBigQuery(BigQueryWriteHelper.java:107) ... 44 more
可行解决方案
原代码报错的核心原因是:Spark将Row对象推断为StructType,而BigQuery Spark连接器默认不会将StructType映射到BigQuery的JSON类型,而是尝试映射为RECORD类型,导致类型不兼容。
解决思路是将Spark中的结构化数据转换为JSON格式的字符串,再写入BigQuery的JSON字段(BigQuery会自动将合法的JSON字符串解析为JSON类型),具体代码如下:
from pyspark.sql.functions import to_json, struct # 读取JSON文件,得到结构化DataFrame df_input = spark.read.json("gs://sample_data_x/sample_json/*.json") # 将整行数据转换为JSON字符串,对应BigQuery的message字段 df = df_input.select(to_json(struct("*")).alias("message")) df.printSchema() df.show() # 写入BigQuery df.write.format('bigquery')\ .option('table', '你的项目ID.数据集ID.表名')\ .option("project", "你的项目ID") \ .option("temporaryGcsBucket","你的临时GCS桶")\ .mode("overwrite")\ .save()
关键说明
to_json+struct("*"):将DataFrame的所有列打包成一个结构体,再转换为JSON格式的字符串,确保输出的message字段是符合JSON规范的字符串。- 类型兼容性:BigQuery的JSON类型支持接收合法的JSON字符串作为输入,无需额外配置类型映射,连接器会自动识别并写入正确的字段类型。
内容的提问来源于stack exchange,提问作者Bhargav Velisetti
相关产品推荐
相关产品推荐

