PySpark实现按年份将单行数据拆分为多行(按列头日期拆分)
高效实现PySpark DataFrame按年份拆分日期列的方案
刚接触Spark的话,真的不用手动迭代处理——PySpark提供了一套高效的内置函数组合,完全能搞定这个宽表转长表再重组的需求,而且性能比循环好太多(毕竟Spark是分布式引擎,循环会绕开它的优化机制)。我给你拆解一下具体步骤,代码直接就能用:
步骤1:构造示例DataFrame(模拟你的原始数据)
先把你的示例数据转换成PySpark可处理的DataFrame,方便后续演示:
from pyspark.sql import SparkSession from pyspark.sql.functions import split, concat, lit, col, expr # 初始化Spark会话 spark = SparkSession.builder.appName("date_column_split").getOrCreate() # 模拟你的原始数据 data = [ (1, "a001", 0, 32, 14, 108), (1, "a002", 80, 0, 0, 92) ] columns = ["id", "name", "01-Jan-10", "01-Feb-10", "01-Jan-11", "01-Feb-11"] df = spark.createDataFrame(data, columns) df.show()
步骤2:将宽表转为长表(Unpivot)
首先把所有日期列(01-Jan-10这类)转换成日期字符串-数值的键值对行,用stack函数一次性完成,这是PySpark处理宽表转长表的标准操作:
# 提取所有日期列(排除id和name) date_cols = [col_name for col_name in df.columns if col_name not in ["id", "name"]] # 构造stack表达式:stack(列数, 列名1, 值1, 列名2, 值2...) stack_expr = f"stack({len(date_cols)}, " + ", ".join([f"'{c}', {c}" for c in date_cols]) + ") as (date_str, value)" # 执行Unpivot操作 unpivoted_df = df.selectExpr("id", "name", stack_expr) unpivoted_df.show()
这一步之后,原本的每一行会拆成N行(N是日期列的数量),每行对应一个日期的数值。
步骤3:提取月份和年份
从日期字符串里解析出月份名称(Jan/Feb)和完整年份(2010/2011):
processed_df = unpivoted_df.withColumn("month", split(col("date_str"), "-")[1]) \ .withColumn("year", concat(lit("20"), split(col("date_str"), "-")[2]).cast("int")) \ .drop("date_str") # 丢弃不需要的原始日期字符串 processed_df.show()
这里用split按-分割日期字符串,索引1是月份,索引2是年份后两位,拼接20前缀得到完整年份(如果你的原始年份是四位,直接取split后的第三部分即可)。
步骤4:将长表转回目标宽表(Pivot)
最后按id、name、year分组,把月份转为列,用first聚合取对应数值(因为每个id+name+year+month只有一个值):
final_df = processed_df.groupBy("id", "name", "year") \ .pivot("month") \ .agg(expr("first(value)")) final_df.show()
这一步执行后,就得到了你想要的格式:每个id+name+year对应一行,月份列对应数值,年份单独成列。
一些注意事项
- 如果你的日期列年份是四位格式(比如
01-Jan-2010),提取年份时直接用split(col("date_str"), "-")[2]即可,不用拼接前缀。 - 如果存在同一个
id+name+year+month有多行数据的情况,可以把聚合函数换成sum(value)或者avg(value),根据你的业务需求调整。 stack函数从Spark 2.4版本开始支持,确保你的Spark版本符合要求。
内容的提问来源于stack exchange,提问作者Tibberzz
相关产品推荐
相关产品推荐

