PySpark如何将Array Struct数组字段转为分号拼接的独立列
错误原因
你代码中使用的explode函数作用是将数组元素拆分为独立行,这是导致数据被拆成多行的核心原因,你的需求不需要拆分原始行,只需对每行内部的Struct数组做字段提取和拼接即可。
解决方案
Spark 2.4及以上版本可以直接使用内置的transform函数实现,无需自定义UDF,性能更高:
PySpark 代码示例
from pyspark.sql import functions as F result_df = sparkdf2.select( F.concat_ws(";", F.expr("transform(contents_json, item -> item.id_seller)")).alias("id_seller"), F.concat_ws(";", F.expr("transform(contents_json, item -> item.tot_product)")).alias("tot_product"), # 数值类型字段需要先转为字符串,避免concat_ws执行报错 F.concat_ws(";", F.expr("transform(contents_json, item -> string(item.tot_unit))")).alias("tot_unit"), F.concat_ws(";", F.expr("transform(contents_json, item -> item.prizes)")).alias("prizes") ) result_df.show()
实现逻辑说明
transform(contents_json, item -> item.字段名):遍历每行的contents_json数组,提取每个Struct元素的对应字段值,生成仅包含该字段值的新数组concat_ws(";", 字段值数组):将数组内所有元素用分号拼接为单个字符串,直接得到你需要的单元格格式
如果你的Spark版本低于2.4,可以使用自定义UDF实现:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def concat_struct_field(arr, field): return ";".join(str(i[field]) for i in arr) concat_udf = F.udf(concat_struct_field, StringType()) result_df = sparkdf2.select( concat_udf("contents_json", F.lit("id_seller")).alias("id_seller"), concat_udf("contents_json", F.lit("tot_product")).alias("tot_product"), concat_udf("contents_json", F.lit("tot_unit")).alias("tot_unit"), concat_udf("contents_json", F.lit("prizes")).alias("prizes") ) result_df.show()
Spark SQL 写法
如果习惯用SQL语法,也可以直接执行:
SELECT concat_ws(';', transform(contents_json, item -> item.id_seller)) AS id_seller, concat_ws(';', transform(contents_json, item -> item.tot_product)) AS tot_product, concat_ws(';', transform(contents_json, item -> string(item.tot_unit))) AS tot_unit, concat_ws(';', transform(contents_json, item -> item.prizes)) AS prizes FROM sparkdf2
内容的提问来源于stack exchange,提问作者user3661384
相关产品推荐
相关产品推荐

