PySpark DataFrame如何基于两列数据新增多列?
PySpark实现行转列:将daysAgo映射为dayN列
需求是将包含applianceName、customer、daysAgo、countAnomaliesByDay的DataFrame,转换为以applianceName和customer为维度,每个daysAgo值对应一个dayN列(如daysAgo=0对应day0),列值为对应的countAnomaliesByDay。
解决方案
使用PySpark的pivot函数实现行转列,核心步骤如下:
- 按维度字段分组,确定结果的行标识
- 对
daysAgo执行转置操作,将行值转为列名 - 聚合提取对应
countAnomaliesByDay的值 - 批量重命名列名以匹配
dayN格式
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col # 初始化SparkSession spark = SparkSession.builder.appName("pivot_demo").getOrCreate() # 创建示例DataFrame sample_data = [ ("app1", "cust1", 0, 100), ("app1", "cust1", 1, 200), ("app1", "cust1", 2, 300), ("app1", "cust1", 3, 400), ("app1", "cust1", 4, 500), ("app1", "cust1", 5, 600), ("app1", "cust1", 6, 700) ] schema = ["applianceName", "customer", "daysAgo", "countAnomaliesByDay"] df1 = spark.createDataFrame(sample_data, schema=schema) # 执行分组+转列操作 pivoted_df = df1.groupBy("applianceName", "customer") \ .pivot("daysAgo") \ .agg({"countAnomaliesByDay": "first"}) # 批量重命名列名 final_df = pivoted_df.select( col("applianceName"), col("customer"), *[col(str(i)).alias(f"day{i}") for i in range(7)] ) # 查看结果 final_df.show()
代码说明
groupBy("applianceName", "customer"):确保同一设备和客户的所有记录合并为一行pivot("daysAgo"):自动识别所有daysAgo的唯一值,将其转为列名agg({"countAnomaliesByDay": "first"}):提取每个分组下对应daysAgo的异常数,因每个分组内daysAgo值唯一,用first/max效果一致- 列重命名部分用列表推导式批量处理,若
daysAgo范围不确定,可通过pivoted_df.columns动态获取列名进行转换
执行结果
+-------------+--------+-----+-----+-----+-----+-----+-----+-----+ |applianceName|customer|day0 |day1 |day2 |day3 |day4 |day5 |day6 | +-------------+--------+-----+-----+-----+-----+-----+-----+-----+ |app1 |cust1 |100 |200 |300 |400 |500 |600 |700 | +-------------+--------+-----+-----+-----+-----+-----+-----+-----+
内容的提问来源于stack exchange,提问作者Karan Alang
相关产品推荐
相关产品推荐

