如何在PySpark DataFrame中按车辆和EU分组对EU_variant排名拼接
问题描述
原始PySpark DataFrame:
id vehicle asIs EU EU_variant 1 A3345 PQ1298 FV1 FV1_variant 2 A3346 PQ1287 FV2 FV2_variant 3 A3346 PQ1207 FV2 FV2_variant 4 A3347 QP9 QP9_variant 5 A3347 QP9 QP9_variant 6 A3347 QP3 QP3_variant 7 A3348 MP6553 YR34 YR34_variant 8 A3348 MP6554 YR35 YR35_variant 9 A3348 MP6554 YR35 YR35_variant
需求:针对每个vehicle和EU的组合,对EU_variant按id顺序进行排名,将EU_variant与排名拼接成新列ECU_Variant_rank,期望结果如下:
id vehicle asIs EU EU_variant ECU_Variant_rank 1 A3345 PQ1298 FV1 FV1_variant FV1_variant(1) 2 A3346 PQ1287 FV2 FV2_variant FV2_variant(1) 3 A3346 PQ1207 FV2 FV2_variant FV2_variant(2) 4 A3347 QP9 QP9_variant QP9_variant(1) 5 A3347 QP9 QP9_variant QP9_variant(2) 6 A3347 QP3 QP3_variant QP3_variant(1) 7 A3348 MP6553 YR34 YR34_variant YR34_variant(1) 8 A3348 MP6554 YR35 YR35_variant YR35_variant(1) 9 A3348 MP6554 YR35 YR35_variant YR35_variant(2)
实现方案
通过PySpark的窗口函数即可实现,步骤如下:
1. 导入必要模块和函数
from pyspark.sql import Window from pyspark.sql.functions import row_number, concat, lit, col
2. 定义窗口规范
按vehicle和EU分组,以id升序排序(保证排名顺序和示例一致):
window_spec = Window.partitionBy("vehicle", "EU").orderBy("id")
3. 计算排名并生成目标列
用row_number()计算组内排名,再通过concat()拼接字符串生成新列:
result_df = df.withColumn("rank", row_number().over(window_spec)) \ .withColumn("ECU_Variant_rank", concat(col("EU_variant"), lit("("), col("rank"), lit(")"))) \ .drop("rank") # 不需要中间rank列可删除
完整代码示例
假设原始DataFrame名为df,完整代码:
from pyspark.sql import Window from pyspark.sql.functions import row_number, concat, lit, col # 定义窗口规则 window_spec = Window.partitionBy("vehicle", "EU").orderBy("id") # 生成结果DataFrame result_df = df.withColumn("rank", row_number().over(window_spec)) \ .withColumn("ECU_Variant_rank", concat(col("EU_variant"), lit("("), col("rank"), lit(")"))) \ .drop("rank") # 查看结果 result_df.show()
说明
partitionBy("vehicle", "EU"):确保排名在每个vehicle+EU的组合内独立计算orderBy("id"):保证排名顺序和示例一致,若需要其他排序逻辑,替换为对应字段即可(如asIs)concat()配合lit():将字符串类型的EU_variant、括号常量与数值类型的rank拼接,避免类型错误
内容的提问来源于stack exchange,提问作者karthik
相关产品推荐
相关产品推荐

