PySpark 2.x:提取Struct指定字段并移除原Struct对应字段
解决方案
1. 提取指定字段(Spark风格实现)
通过Spark内置列表达式批量生成提取逻辑,完全规避Python循环逐行处理:
先导入必要函数:
from pyspark.sql.functions import col, struct
定义要提取的字段列表:
extract_fields = ["name", "surname", "id"]
2. 重构properties列移除已提取字段(适配PySpark 2.x)
由于PySpark 2.x没有dropFields()方法,我们通过保留未提取字段并重新构造Struct的方式实现去重:
- 获取
properties列的所有Struct字段名:
all_struct_fields = df.schema["properties"].dataType.names
- 过滤出未被提取的字段:
remaining_fields = [field for field in all_struct_fields if field not in extract_fields]
- 构造新的
properties列,仅保留剩余字段:
new_properties = struct(*[col(f"properties.{field}").alias(field) for field in remaining_fields])
合并两步得到最终结果
将提取字段与重构后的properties列整合:
result_df = df.select( new_properties.alias("properties"), *[col(f"properties.{field}").alias(field) for field in extract_fields] )
执行后即可得到目标格式的DataFrame:
|properties|name|surname|id | |[foo, bar]|john|doe |123|
内容的提问来源于stack exchange,提问作者Jacek1
相关产品推荐
相关产品推荐

