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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 08:27:03