Spark用BigQueryOperator写BigQuery REPEATED RECORD字段Schema不匹配如何解决
报错原因
你遇到的Cannot add fields (field: actions.list)报错由两个问题共同导致:
- 你定义的Spark顶级Schema字段名和输入JSON的实际键名不匹配,会导致读出的
record、d字段为空值 - Spark的BigQuery连接器默认使用JSON作为中间写入格式时,会自动将Array类型字段包装为
{"list": [数组元素]}的结构,和BigQuery侧已经定义的REPEATED RECORD类型结构不匹配,被识别为要新增list子字段。
解决方法
步骤1:修正Schema定义,对齐JSON字段名
调整顶级Schema的字段名,和你输入JSON的recordName、date键名完全对应:
import org.apache.spark.sql.types._ val action: StructType = new StructType() .add("id", StringType) .add("name", StringType) .add("last", StringType) val actionsList = ArrayType(action, containsNull = true) // 顶级字段对齐JSON实际键名 val finalStruct: StructType = new StructType() .add("recordName", StringType) .add("date", StringType) .add("actions", actionsList)
步骤2:写入时指定Parquet作为中间格式
Parquet的数组结构和BigQuery的REPEATED类型天然对齐,不会额外包装list层,写入时新增对应配置即可:
val df = spark.read.schema(finalStruct).json(rdd) df.createOrReplaceTempView("myData") val finalDf = spark.sql("SELECT recordName as my_rec, date as inc_date, actions from myData") finalDf.write .mode("append") .format("bigquery") .option("intermediateFormat", "parquet") // 新增这一行即可解决list字段报错 // 保留你原本的其他BigQuery配置(表名、临时GCS桶等) .save()
注意事项
需要确认BigQuery侧actions字段的子字段名称、类型、顺序,和你Spark侧定义的action StructType完全一致,避免出现其他结构不匹配报错。
内容的提问来源于stack exchange,提问作者asdasd31
相关产品推荐
相关产品推荐

