You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 21:40:16