在PySpark中计算用户非重叠工作天数及休假天数
处理PySpark重叠任职记录的天数计算方案
核心思路
先对每个用户的任职记录按入职日期排序,合并重叠或连续的区间,再基于合并后的区间计算累计工作天数;未工作天数则用用户的总时间跨度(最早入职日到最晚离职日)减去累计工作天数。
具体实现步骤
1. 创建示例DataFrame
假设原始数据结构如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("WorkDaysCalculation").getOrCreate() data = [ ("user1", "2020-01-01", "2020-03-31"), ("user1", "2020-03-15", "2020-05-31"), ("user1", "2020-07-01", "2020-09-30"), ("user2", "2021-02-01", "2021-04-30"), ("user2", "2021-05-10", "2021-06-30"), ("user2", "2021-06-25", "2021-08-15") ] df = spark.createDataFrame(data, ["user", "hiring_date", "termination_date"]) df = df.withColumn("hiring_date", F.to_date("hiring_date")) df = df.withColumn("termination_date", F.to_date("termination_date"))
2. 合并重叠/连续的任职区间
通过窗口函数按用户分组,对入职日期排序,标记需要合并的区间:
# 定义窗口:按用户分组,按入职日期升序排列 window = Window.partitionBy("user").orderBy("hiring_date") # 计算上一个区间的结束日期,判断当前区间是否与上一个重叠/连续 df_with_prev = df.withColumn( "prev_termination", F.lag("termination_date").over(window) ).withColumn( "is_overlap", F.when( F.col("hiring_date") <= F.date_add(F.col("prev_termination"), 1), 1 ).otherwise(0) ) # 为每个连续的区间组分配ID df_with_group = df_with_prev.withColumn( "group_id", F.sum("is_overlap").over(window.rangeBetween(Window.unboundedPreceding, 0)) ) # 合并每个组的区间:取最早入职日和最晚离职日 merged_df = df_with_group.groupBy("user", "group_id").agg( F.min("hiring_date").alias("start_date"), F.max("termination_date").alias("end_date") ).drop("group_id")
3. 计算累计非重叠工作天数
# 计算每个合并区间的天数,再求和 work_days_df = merged_df.withColumn( "interval_days", F.datediff(F.col("end_date"), F.col("start_date")) + 1 # 包含首尾日期 ).groupBy("user").agg( F.sum("interval_days").alias("total_work_days") )
4. 计算未工作天数
需要先确定每个用户的总时间跨度(最早入职到最晚离职):
# 计算用户的总时间范围 user_time_range = df.groupBy("user").agg( F.min("hiring_date").alias("first_hire"), F.max("termination_date").alias("last_termination") ).withColumn( "total_days", F.datediff(F.col("last_termination"), F.col("first_hire")) + 1 ) # 关联工作天数,计算未工作天数 final_df = user_time_range.join(work_days_df, on="user").withColumn( "unworked_days", F.col("total_days") - F.col("total_work_days") ).select("user", "total_work_days", "unworked_days") final_df.show()
关键说明
- 合并区间时,连续的日期(如前一个区间结束于3月31日,下一个开始于4月1日)也视为连续区间,如果不需要合并连续日期,把
F.date_add(F.col("prev_termination"), 1)改成F.col("prev_termination")即可。 - 天数计算包含入职和离职当天,若不需要包含,去掉
+1。
内容的提问来源于stack exchange,提问作者JDDS
相关产品推荐
相关产品推荐

