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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 05:24:52