修改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))
问题原因分析
transform_schema方法参数传递错误
处理数组类型时,调用transform_schema未传递parent参数,导致内部无法识别字段层级,changeDatatypeforNestedField无法精准匹配C_0_0.C_2_0字段,可能给所有非结构/数组字段设置错误类型,最终from_json解析失败返回NULL。to_json参数使用错误
代码中to_json(col_name)直接传入字符串列名,PySpark会将其当作字面量而非列对象,导致序列化的JSON内容错误,后续from_json无法解析,返回NULL。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
相关产品推荐
相关产品推荐

