基于PySpark统计员工非连续任职同一部门次数的技术问询
解决方案:PySpark统计员工非连续部门任职次数
核心思路
要统计非连续任职次数,本质是识别同一员工同一部门的连续任职段,每个连续段算1次。关键是先确保每个员工的记录在同一分区内按时间排序,再通过对比前后记录的部门ID来标记新的任职段,最后聚合计数。
步骤实现
数据预处理与分区排序
先将日期字符串转为日期类型,再按emp_id分区(确保同一员工数据在同一Executor处理,避免跨Executor排序混乱),最后按begin日期排序。标记新任职段
使用窗口函数lag获取上一条记录的dept_id,如果当前dept_id与上一条不同,或是第一条记录,则标记为新的任职段开始。聚合统计次数
按emp_id和dept_id分组,统计每个组内的新任职段标记数量,即为非连续任职次数。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("DeptTenureCount").getOrCreate() # 模拟原始数据(替换为你的数据源) data = [ ("07a0fcf5", "30.06.2021", "30.06.2021", "1443"), ("07a0fcf5", "01.07.2021", "01.07.2021", "1443"), ("07a0fcf5", "02.07.2021", "02.07.2021", "1269"), ("07a0fcf5", "03.07.2021", "11.07.2021", "1269"), ("07a0fcf5", "12.07.2021", "14.07.2021", "1269"), ("07a0fcf5", "15.07.2021", "15.07.2021", "1273"), ("07a0fcf5", "16.07.2021", "30.08.2021", "1273"), ("07a0fcf5", "31.08.2021", "05.10.2021", "1273"), ("07a0fcf5", "06.10.2021", "21.02.2022", "1269"), ("07a0fcf5", "24.02.2022", "23.06.2022", "1269"), ("07a0fcf5", "24.06.2022", "01.01.9999", "1269"), ("07d06bee", "28.06.2021", "29.06.2021", "1273"), ("07d06bee", "30.06.2021", "30.06.2021", "1287"), ("07d06bee", "01.07.2021", "26.07.2021", "1443"), ("07d06bee", "27.07.2021", "27.07.2021", "1287"), ("07d06bee", "28.07.2021", "01.08.2021", "1443"), ("07d06bee", "02.08.2021", "01.01.9999", "1287"), ("07d1fdd3", "25.05.2021", "25.05.2021", "1256"), ("07d1fdd3", "26.05.2021", "26.05.2021", "1256"), ("07d1fdd3", "27.05.2021", "27.05.2021", "1256"), ("07d1fdd3", "28.05.2021", "06.06.2021", "1256"), ("07d1fdd3", "07.06.2021", "18.06.2021", "1256"), ("07d1fdd3", "19.06.2021", "20.06.2021", "1256"), ("07d1fdd3", "21.06.2021", "21.06.2021", "1256"), ("07d1fdd3", "22.06.2021", "06.07.2021", "1256"), ("07d1fdd3", "07.07.2021", "13.07.2021", "1098"), ("07d1fdd3", "14.07.2021", "16.08.2021", "1098"), ("07d1fdd3", "17.08.2021", "25.08.2021", "1098"), ("07d1fdd3", "26.08.2021", "26.08.2021", "1098"), ("07d1fdd3", "27.08.2021", "06.09.2021", "1098") ] df = spark.createDataFrame(data, ["emp_id", "begin", "end", "dept_id"]) # 1. 转换日期格式,按emp_id分区、begin排序 df = df.withColumn("begin_date", F.to_date(F.col("begin"), "dd.MM.yyyy")) window_spec = Window.partitionBy("emp_id").orderBy("begin_date") # 2. 标记新的任职段:当前dept_id与上一条不同,或为第一条记录则标记为1 df = df.withColumn("prev_dept", F.lag("dept_id").over(window_spec)) df = df.withColumn( "is_new_segment", F.when( F.col("prev_dept").isNull() | (F.col("dept_id") != F.col("prev_dept")), 1 ).otherwise(0) ) # 3. 按emp_id和dept_id分组,统计新段数量 result_df = df.groupBy("emp_id", "dept_id").agg(F.sum("is_new_segment").alias("times_worked")) # 展示结果 result_df.orderBy("emp_id", "dept_id").show()
关键说明
- 分区策略:通过
partitionBy("emp_id")确保同一员工的所有记录落在同一Executor,彻底解决跨Executor排序导致的窗口函数失效问题。 - 日期排序:必须将字符串日期转为日期类型再排序,否则字符串排序会出现逻辑错误(比如"01.10.2021"会错误排在"02.09.2021"前面)。
- 标记逻辑:
lag函数获取上一条部门ID,对比后标记新段,最后求和得到每个部门的任职次数,完全符合非连续任职的统计需求。
验证结果
运行代码后输出结果与期望一致:
+---------+-------+------------+ | emp_id|dept_id|times_worked| +---------+-------+------------+ |07a0fcf5| 1443| 1| |07a0fcf5| 1269| 2| |07a0fcf5| 1273| 1| |07d06bee| 1273| 1| |07d06bee| 1287| 3| |07d06bee| 1443| 2| |07d1fdd3| 1256| 1| |07d1fdd3| 1098| 1| +---------+-------+------------+
内容的提问来源于stack exchange,提问作者uri_leo
相关产品推荐
相关产品推荐

