如何在PySpark DataFrame中根据sub和rank列扩展为多列?
PySpark DataFrame按rank映射sub值到对应列
现有如下PySpark DataFrame:
+--------------------------------+-------+------------------------------------+-----------+------------------+------------------+ |id |order |cart |item |sub |rank | +--------------------------------+-------+------------------------------------+-----------+------------------+------------------+ |694e45100f52475d8dac13b829e02394|null |baa4b664-8a84-4abb-8919-01836340a30a|826017462 |697030209 |6 | |532f20768ce64599b6f79bb17b9fc597|null |aa45d078-dfc6-4710-adac-4385b43d3eba|10402650 |13893731 |3 | |fae01765055f4df0a1041de0669393e6|null |f7b498d6-ea9c-47e7-9220-ec18070f9ee8|293214825 |1226495201 |4 | |f2135f7919713625e044001517f43a86|null |d4873f7a-2219-49d9-a0e8-cb41d24f8aec|814735847 |222288904 |2 | |6e74f194b9d2403cb6d28f0a33a00a9a|null |20753b35-36bd-4679-9d64-c480c9c19b97|10291833 |10313028 |1 | +--------------------------------+-------+------------------------------------+-----------+------------------+------------------+
需要根据rank列的值,将sub列的内容放到对应的sub{rank}列中,其余subN列填充0,得到如下结果:
+--------------------------------+-------+------------------------------------+-----------+------------------+------------------+-------------+-------------+-------------+-------------+-------------+-------------+-------------+ |id |order |cart |item |sub |rank | sub1 | sub2 | sub3 | sub4 | sub5 | sub6 | sub7 | +--------------------------------+-------+------------------------------------+-----------+------------------+------------------+-------------+-------------+-------------+-------------+-------------+-------------+-------------+ |694e45100f52475d8dac13b829e02394|null |baa4b664-8a84-4abb-8919-01836340a30a|826017462 |697030209 |6 | 0 | 0 | 0 | 0 | 0 |697030209 | 0 | |532f20768ce64599b6f79bb17b9fc597|null |aa45d078-dfc6-4710-adac-4385b43d3eba|10402650 |13893731 |3 | 0 | 0 |13893731 | 0 | 0 | 0 | 0 | |fae01765055f4df0a1041de0669393e6|null |f7b498d6-ea9c-47e7-9220-ec18070f9ee8|293214825 |1226495201 |4 | 0 | 0 | 0 | 1226495201 | 0 | 0 | 0 | |f2135f7919713625e044001517f43a86|null |d4873f7a-2219-49d9-a0e8-cb41d24f8aec|814735847 |222288904 |2 | 0 | 222288904 | 0 | 0 | 0 | 0 | 0 | |6e74f194b9d2403cb6d28f0a33a00a9a|null |20753b35-36bd-4679-9d64-c480c9c19b97|10291833 |10313028 |1 | 10313028 | 0 | 0 | 0 | 0 | 0 | 0 | +--------------------------------+-------+------------------------------------+-----------+------------------+------------------+--------------------------------------------------------------------------------------------------
解决方案
可以通过PySpark的条件函数when/otherwise,配合列表推导式批量生成目标列:
from pyspark.sql import functions as F # 假设原始DataFrame名为df # 定义需要生成的sub列范围(这里是1到7) max_rank = 7 # 批量生成sub1到sub7的列表达式 sub_cols = [ F.when(F.col("rank") == n, F.col("sub")).otherwise(F.lit(0)).alias(f"sub{n}") for n in range(1, max_rank + 1) ] # 合并原始列与新生成的列 result_df = df.select("*", *sub_cols) # 查看结果 result_df.show(truncate=False)
灵活适配最大rank值
如果不确定最大rank值,可以先从DataFrame中自动获取:
# 获取数据中的最大rank值 max_rank = df.agg(F.max("rank")).first()[0] # 后续代码同上 sub_cols = [ F.when(F.col("rank") == n, F.col("sub")).otherwise(F.lit(0)).alias(f"sub{n}") for n in range(1, max_rank + 1) ] result_df = df.select("*", *sub_cols)
这段代码的逻辑是:对每个目标列subN,判断当前行的rank是否等于N,匹配则填充sub的值,否则填充0,最后将所有新列与原始列组合,得到所需结果。
内容的提问来源于stack exchange,提问作者Shibu
相关产品推荐
相关产品推荐

