Spark unionByName处理含#的结构体子列时解析异常问题排查
问题
在Azure Synapse中读取4个Azure存储账户的Parquet数据到DataFrame,使用unionByName合并以统一查询时触发org.apache.spark.sql.catalyst.parser.ParseException错误。数据包含名为ABC的结构体列,其子列#XYZ#(#为哈希值)单独查询正常,但合并操作报错,需确认是Spark bug还是操作疏漏。
合并代码示例
val data1 = spark.read.parquet("abfss://PATH1") val data2 = spark.read.parquet("abfss://PATH2") val data3 = spark.read.parquet("abfss://PATH3") val data4 = spark.read.parquet("abfss://PATH4") val data = data1 .unionByName(data2, allowMissingColumns=true) .unionByName(data3, allowMissingColumns=true) .unionByName(data4, allowMissingColumns=true) data.printSchema()
单独查询正常的代码
data1.select("ABC.#XYZ#").where(col("#XYZ#").isNotNull).show(5, truncate = false)
错误信息
org.apache.spark.sql.catalyst.parser.ParseException: 意外输入 '#',预期 {'ADD', 'AFTER', 'ALL', ...}(第1行,第0位) == SQL == #XYZ# ^^^
原因分析与解决方案
这不是Spark的bug,核心问题是列名包含Spark SQL语法的特殊字符#:
- 单独查询时使用Scala API的
col方法,直接通过字符串引用列名,不会触发SQL解析器,因此可以正常执行。 unionByName在处理列名匹配时,内部会调用SQL解析器处理列名,而未对特殊字符#做自动转义,导致解析失败。
解决方法
方法1:重命名特殊字符列
合并前将所有DataFrame中的ABC.#XYZ#列重命名为不含特殊字符的名称,合并后可按需改回:
// 定义重命名函数 def renameSpecialCols(df: DataFrame): DataFrame = { df.withColumnRenamed("ABC.#XYZ#", "ABC_XYZ") } // 对所有DataFrame重命名后执行合并 val data = renameSpecialCols(data1) .unionByName(renameSpecialCols(data2), allowMissingColumns=true) .unionByName(renameSpecialCols(data3), allowMissingColumns=true) .unionByName(renameSpecialCols(data4), allowMissingColumns=true)
方法2:显式指定带反引号的Schema
读取数据时指定包含反引号包裹特殊字符的Schema,确保SQL解析器能正确识别子列名:
// 构建正确的Schema,用反引号包裹特殊字符列名 val correctSchema = StructType(Array( StructField("ABC", StructType(Array( StructField("`#XYZ#`", StringType) // 根据实际数据类型调整 ))) )) // 读取数据时指定Schema val data1 = spark.read.schema(correctSchema).parquet("abfss://PATH1") val data2 = spark.read.schema(correctSchema).parquet("abfss://PATH2") val data3 = spark.read.schema(correctSchema).parquet("abfss://PATH3") val data4 = spark.read.schema(correctSchema).parquet("abfss://PATH4") // 执行合并操作 val data = data1 .unionByName(data2, allowMissingColumns=true) .unionByName(data3, allowMissingColumns=true) .unionByName(data4, allowMissingColumns=true)
内容的提问来源于stack exchange,提问作者ROODAY
相关产品推荐
相关产品推荐

