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

PySpark将点符号列转为嵌套JSON时出现列不存在错误

Spark根据带点列名生成嵌套JSON报错问题

原始代码与数据

df_renamed = df.withColumnRenamed("id","steps.id").withColumnRenamed("status_1","steps.status").withColumnRenamed("severity","steps.error.severity")

df_renamed.show(truncate=False)

输出结果:

+----------+-------+------+-----------------------+------------+--------------------+
|apiVersion|expired|status|steps.id               |steps.status|steps.error.severity|
+----------+-------+------+-----------------------+------------+--------------------+
|2         |false  |200   |mexican-curp-validation|200         |null                |
+----------+-------+------+-----------------------+------------+--------------------+  

期望输出

将带点列名转换为嵌套JSON结构,最终输出:

+----------+-------+------+-------------------------------------------------------------------------------------+
|apiVersion|expired|status|steps                                                                                |
+----------+-------+------+-------------------------------------------------------------------------------------+
|2         |false  |200   |{"id":"mexican-curp-validation","status":200,"error":{"severity":null}}               |
+----------+-------+------+-------------------------------------------------------------------------------------+  

尝试的代码

cols_list = [name for name in df_renamed.columns if "." in name]
df_new = df_renamed.withColumn("steps",F.to_json(F.struct(*cols_list)))
df_new.show()   

错误信息

df_new = df_renamed.withColumn("steps",F.to_json(F.struct(*cols_list)))
  File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/pyspark/sql/dataframe.py", line 3036, in withColumn
    return DataFrame(self._jdf.withColumn(colName, col._jc), self.sparkSession)
  File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/py4j/java_gateway.py", line 1321, in __call__
    return_value = get_return_value(
  File "/Users/../IdeaProjects/pocs/venvsd/lib/python3.9/site-packages/pyspark/sql/utils.py", line 196, in deco
    raise converted from None
pyspark.sql.utils.AnalysisException: Column 'steps.id' does not exist. Did you mean one of the following? [steps.id, expired, status, steps.status, apiVersion, steps.error.severity];
'Project [apiVersion#17, expired#18, status#19, steps.id#29, steps.status#37, steps.error.severity#44, to_json(struct(id, 'steps.id, status, 'steps.status, severity, 'steps.error.severity), Some(GMT+05:30)) AS steps#82]
+- Project [apiVersion#17, expired#18, status#19, steps.id#29, steps.status#37, severity#22 AS steps.error.severity#44]
   +- Project [apiVersion#17, expired#18, status#19, steps.id#29, status_1#21 AS steps.status#37, severity#22]
      +- Project [apiVersion#17, expired#18, status#19, id#20 AS steps.id#29, status_1#21, severity#22]
         +- Relation [apiVersion#17,expired#18,status#19,id#20,status_1#21,severity#22] csv

问题原因与解决方案

原因

Spark会将带点的列名字符串(如steps.id)解析为嵌套结构体字段引用(即从steps列中提取id字段),但你的数据中并没有steps结构体列,只有名为steps.id的普通列,因此报错提示找不到对应字段。

解决方案

使用F.col()明确引用完整列名,同时手动构建嵌套结构体来匹配JSON层级:

import pyspark.sql.functions as F

# 构建error子结构体
error_struct = F.struct(
    F.col("steps.error.severity").alias("severity")
).alias("error")

# 构建外层steps结构体
steps_struct = F.struct(
    F.col("steps.id").alias("id"),
    F.col("steps.status").alias("status"),
    error_struct
)

# 生成嵌套JSON列,并删除原带点列
df_new = df_renamed.withColumn("steps", F.to_json(steps_struct))
df_new = df_new.drop(*[name for name in df_renamed.columns if "." in name])

df_new.show(truncate=False)

通用处理思路(多层级列名)

如果有更多层级的带点列名(如a.b.c.d),可以通过递归解析列名的点分隔结构,自动构建嵌套结构体,示例逻辑:

def build_nested_struct(cols):
    struct_map = {}
    for col_name in cols:
        parts = col_name.split(".")
        current = struct_map
        for part in parts[:-1]:
            if part not in current:
                current[part] = {}
            current = current[part]
        current[parts[-1]] = F.col(col_name).alias(parts[-1])
    
    def build_struct_from_map(m):
        struct_fields = []
        for k, v in m.items():
            if isinstance(v, dict):
                struct_fields.append(build_struct_from_map(v).alias(k))
            else:
                struct_fields.append(v)
        return F.struct(*struct_fields)
    
    return build_struct_from_map(struct_map)

# 使用示例
cols_with_dots = [name for name in df_renamed.columns if "." in name]
nested_struct = build_nested_struct(cols_with_dots)
df_new = df_renamed.withColumn("steps", F.to_json(nested_struct)).drop(*cols_with_dots)
df_new.show(truncate=False)

内容的提问来源于stack exchange,提问作者RushHour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 02:36:03