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

修改Schema后PySpark DataFrame中C_0_0列值为NULL的问题排查

问题描述

将JSON字符串读取到PySpark DataFrame后,需要完成以下类型转换:

  • 数组C_0_0元素中的C_2_0从double类型改为decimal(6,6)
  • C_0_1下的字符串类型时间字段改为timestamp类型(保证无数据损失)
    但使用自定义transform_schema方法转换Schema后,C_0_0列全部变为NULL,需排查问题原因并解决。

读取JSON代码

rdd = sc.parallelize([json_str])
nested_df = hc.read.json(rdd)

当前Schema

# root
# |-- C_0_0: array (nullable = true)
# |    |-- element: struct (containsNull = true)
# |    |    |-- C_2_0: double (nullable = true)
# |    |    |-- C_2_1: double (nullable = true)
# |-- C_0_1: struct (nullable = true)
# |    |-- C_1_0: struct (nullable = true)
# |    |    |-- C_2_0: string (nullable = true)
# |    |    |-- C_2_1: string (nullable = true)
# |    |-- C_1_1: struct (nullable = true)
# |    |    |-- C_2_0: string (nullable = true)
# |    |    |-- C_2_1: double (nullable = true)

目标Schema

# root
# |-- C_0_0: array (nullable = true)
# |    |-- element: struct (containsNull = true)
# |    |    |-- C_2_0: decimal(6,6) (nullable = true)
# |    |    |-- C_2_1: double (nullable = true)
# |-- C_0_1: struct (nullable = true)
# |    |-- C_1_0: struct (nullable = true)
# |    |    |-- C_2_0: timestamp (nullable = true)
# |    |    |-- C_2_1: timestamp (nullable = true)
# |    |-- C_1_1: struct (nullable = true)
# |    |    |-- C_2_0: timestamp (nullable = true)
# |    |    |-- C_2_1: double (nullable = true)

建表语句

CREATE TABLE CompoundDataTypesSchema.TABLE_NAME ( C_0_0 ARRAY OF ROW ( C_2_0 DECIMAL(6,6) ,C_2_1 DOUBLE ) ,C_0_1 ROW ( C_1_0 ROW ( C_2_0 TIMESTAMP ,C_2_1 TIMESTAMP ),C_1_1 ROW ( C_2_0 TIMESTAMP ,C_2_1 DOUBLE ) ) ) in FILES_SERVICE

转换代码

def transform_schema(self, schema, parent=""):
        if schema == None:
            return StructType()
        new_schema = []    
        for f in schema.fields:
            if parent:
                field_name = parent + '.' + f.name
            else:
                field_name = f.name
            if isinstance(f.dataType, ArrayType):
                new_schema.append(StructField(f.name, ArrayType(self.transform_schema(f.dataType.elementType))))
            elif isinstance(f.dataType, StructType):
                new_schema.append(StructField(f.name, self.transform_schema(f.dataType)))
            else:
                new_datatype = self.changeDatatypeforNestedField()
                new_schema.append(StructField(f.name, new_datatype, f.nullable))
        return StructType(new_schema)
nested_df_schema = nested_df.schema
for f in nested_df_schema.fields:
    print("Name: ", f.name)
    col_name = f.name
    if isinstance(f.dataType, ArrayType):
        new_schema = ArrayType(self.transform_schema(f.dataType.elementType, parent = f.name))
        nested_df = nested_df.withColumn("col_name_json", to_json(col_name)).drop(col_name)
        nested_df = nested_df.withColumn(col_name, from_json("col_name_json", new_schema)).drop("col_name_json")
    elif isinstance(f.dataType, StructType):
        new_schema = self.transform_schema(f.dataType, parent = f.name)
        nested_df = nested_df.withColumn("col_name_json", to_json(col_name)).drop(col_name)
        nested_df = nested_df.withColumn(col_name, from_json("col_name_json", new_schema)).drop("col_name_json")
    else:
        new_datatype = self.changeDatatypeforNestedField()
        nested_df = nested_df.withColumn(col_name, nested_df[col_name].cast(new_datatype))

问题原因分析
  1. transform_schema方法参数传递错误
    处理数组类型时,调用transform_schema未传递parent参数,导致内部无法识别字段层级,changeDatatypeforNestedField无法精准匹配C_0_0.C_2_0字段,可能给所有非结构/数组字段设置错误类型,最终from_json解析失败返回NULL。

  2. to_json参数使用错误
    代码中to_json(col_name)直接传入字符串列名,PySpark会将其当作字面量而非列对象,导致序列化的JSON内容错误,后续from_json无法解析,返回NULL。

  3. changeDatatypeforNestedField逻辑缺失
    该方法未根据字段路径区分转换规则,可能对C_0_0下的C_2_1错误应用decimal类型,破坏数组结构的解析规则。


解决方案

方案1:修复自定义Schema转换逻辑

from pyspark.sql.types import StructType, StructField, ArrayType, DecimalType, TimestampType, DoubleType
from pyspark.sql.functions import to_json, from_json, col

def change_datatype_for_field(field_path):
    # 根据字段路径匹配转换规则
    if field_path.endswith(".C_2_0"):
        if field_path.startswith("C_0_0"):
            return DecimalType(6,6)
        elif field_path.startswith("C_0_1"):
            return TimestampType()
    elif field_path.endswith(".C_2_1"):
        if field_path.startswith("C_0_1.C_1_0"):
            return TimestampType()
        else:
            return DoubleType()
    return None

def transform_schema(schema, parent=""):
    new_schema = []
    for f in schema.fields:
        current_path = f"{parent}.{f.name}" if parent else f.name
        if isinstance(f.dataType, ArrayType):
            # 传递parent参数,保留字段层级
            element_schema = transform_schema(f.dataType.elementType, parent=current_path)
            new_schema.append(StructField(f.name, ArrayType(element_schema), f.nullable))
        elif isinstance(f.dataType, StructType):
            struct_schema = transform_schema(f.dataType, parent=current_path)
            new_schema.append(StructField(f.name, struct_schema, f.nullable))
        else:
            new_datatype = change_datatype_for_field(current_path)
            # 未匹配到规则则保留原类型
            if new_datatype is None:
                new_datatype = f.dataType
            new_schema.append(StructField(f.name, new_datatype, f.nullable))
    return StructType(new_schema)

# 修复to_json参数,使用col()引用列
nested_df_schema = nested_df.schema
for f in nested_df_schema.fields:
    col_name = f.name
    if isinstance(f.dataType, ArrayType):
        new_schema = ArrayType(transform_schema(f.dataType.elementType, parent=col_name))
        nested_df = nested_df.withColumn(f"{col_name}_json", to_json(col(col_name))).drop(col_name)
        nested_df = nested_df.withColumn(col_name, from_json(col(f"{col_name}_json"), new_schema)).drop(f"{col_name}_json")
    elif isinstance(f.dataType, StructType):
        new_schema = transform_schema(f.dataType, parent=col_name)
        nested_df = nested_df.withColumn(f"{col_name}_json", to_json(col(col_name))).drop(col_name)
        nested_df = nested_df.withColumn(col_name, from_json(col(f"{col_name}_json"), new_schema)).drop(f"{col_name}_json")
    else:
        new_datatype = change_datatype_for_field(col_name)
        if new_datatype is not None:
            nested_df = nested_df.withColumn(col_name, nested_df[col_name].cast(new_datatype))

方案2:直接递归转换嵌套字段(更高效)

跳过JSON序列化/反序列化步骤,直接对嵌套字段做类型转换,避免解析失败风险:

from pyspark.sql.functions import col, array_transform, struct
from pyspark.sql.types import DecimalType, TimestampType

# 转换C_0_0数组中的C_2_0
nested_df = nested_df.withColumn(
    "C_0_0",
    array_transform(
        col("C_0_0"),
        lambda x: struct(
            x["C_2_0"].cast(DecimalType(6,6)).alias("C_2_0"),
            x["C_2_1"].alias("C_2_1")
        )
    )
)

# 转换C_0_1下的时间字段
nested_df = nested_df.withColumn(
    "C_0_1",
    struct(
        struct(
            col("C_0_1.C_1_0.C_2_0").cast(TimestampType()).alias("C_2_0"),
            col("C_0_1.C_1_0.C_2_1").cast(TimestampType()).alias("C_2_1")
        ).alias("C_1_0"),
        struct(
            col("C_0_1.C_1_1.C_2_0").cast(TimestampType()).alias("C_2_0"),
            col("C_0_1.C_1_1.C_2_1").alias("C_2_1")
        ).alias("C_1_1")
    )
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 01:55:35