如何在PySpark中从日期数据集提取每年最后一个工作日?
解决Spark提取每年最后一个工作日的问题
问题分析
你需要从日期数据集中提取每年的最后一个工作日,规则是:若12月31日为周六或周日,则取12月30日。之前硬编码匹配12月31日的方式在2022年失效,因为该年12月31日是周六,而数据集中只有12月30日的记录。
原代码存在几个问题:
- 硬构造12-31日期匹配,无法覆盖该日期不在数据集的场景
- 使用
collect()遍历年份,不符合Spark分布式处理的最佳实践 - 存在语法错误:
concat_ws拼写错误、参数分隔符误用点号、filter语句未闭合、all_years变量未定义
解决方案
方案一:严格按规则计算目标日期后匹配
先根据年份计算出符合规则的年末目标日期,再从原数据集中匹配对应的记录:
from pyspark.sql import functions as F from pyspark.sql import types as T # 先确保Date列是日期类型(若原数据是字符串格式) df = df.withColumn("Date", F.to_date(F.col("Date"), "dd-MM-yyyy")) # 提取年份列 df = df.withColumn("Year", F.year(F.col("Date"))) # 生成每个年份的目标年末日期 target_dates = df.select("Year").distinct().orderBy("Year") \ .withColumn("Dec31", F.to_date(F.concat_ws("-", F.lit("31"), F.lit("12"), F.col("Year")), "dd-MM-yyyy")) \ .withColumn("DayOfWeek", F.dayofweek(F.col("Dec31"))) # Spark中dayofweek返回1=周日, 7=周六 .withColumn( "年末日期", F.when(F.col("DayOfWeek").isin(1, 7), F.date_sub(F.col("Dec31"), 1)) .otherwise(F.col("Dec31")) ) # 关联原数据集,筛选存在的目标日期 result = df.join( target_dates, (df["Year"] == target_dates["Year"]) & (df["Date"] == target_dates["年末日期"]), how="inner" ).select("年末日期").distinct().orderBy("年末日期") # 查看结果 result.show()
方案二:直接取每年12月的最大日期
如果你的数据集已经包含了符合规则的年末工作日(比如2022年只有12月30日),可以直接筛选12月的日期,按年份分组取最大值:
from pyspark.sql import functions as F # 转换Date列为日期类型 df = df.withColumn("Date", F.to_date(F.col("Date"), "dd-MM-yyyy")) # 筛选12月数据,按年份取最大日期 result = df.filter(F.month(F.col("Date")) == 12) \ .groupBy(F.year(F.col("Date")).alias("Year")) \ .agg(F.max(F.col("Date")).alias("年末日期")) \ .orderBy("年末日期") \ .select("年末日期") result.show()
两种方案都能得到你期望的输出:
| 年末日期 |
|---|
| 31-12-2020 |
| 31-12-2021 |
| 30-12-2022 |
内容的提问来源于stack exchange,提问作者SaurabhShelar
相关产品推荐
相关产品推荐

