Pyspark如何将相同时间戳的多行按name转新列合并?
PySpark行转列实现方案
这个需求属于典型的行转列(透视)场景,PySpark内置的groupBy+pivot组合可以直接实现,方案通用无需自定义复杂逻辑。
核心实现逻辑
按Timestamp字段分组,对name列做透视,用聚合函数取对应value值即可,无匹配数据默认填充null,可按需调整为nan。
代码示例
1. 样例数据构造(可跳过,直接套用核心逻辑到你的数据)
from pyspark.sql import SparkSession from pyspark.sql.functions import first import numpy as np # 初始化SparkSession spark = SparkSession.builder.appName("pivot_demo").getOrCreate() # 模拟输入数据 data = [ ("2024-01-01 00:00:00", "temperature", 23.5), ("2024-01-01 00:00:00", "humidity", 65), ("2024-01-01 01:00:00", "temperature", 22.1), ("2024-01-01 01:00:00", "pressure", 1013), ] df = spark.createDataFrame(data, schema=["Timestamp", "name", "value"])
2. 核心透视代码
pivot_df = df.groupBy("Timestamp") \ .pivot("name") \ .agg(first("value")) # 若需要将默认null替换为nan,增加以下代码 pivot_df = pivot_df.fillna(np.nan, subset=pivot_df.columns[1:]) # 查看输出结果 pivot_df.show()
方案特性
- 通用适配:无需提前知晓
name列的所有取值,pivot会自动识别所有不同的name生成对应列 - 灵活调整:如果同一个Timestamp下同一个name存在多条重复记录,可将
first聚合函数替换为sum(求和)、avg(求平均)、collect_list(收集为列表)等符合业务需求的聚合规则 - 性能优化:如果
name列的取值是已知固定的,可手动传入取值列表减少扫描开销,示例:.pivot("name", ["temperature", "humidity", "pressure"])
内容的提问来源于stack exchange,提问作者Afrobeta
相关产品推荐
相关产品推荐

