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

基于PySpark统计员工非连续任职同一部门次数的技术问询

解决方案:PySpark统计员工非连续部门任职次数

核心思路

要统计非连续任职次数,本质是识别同一员工同一部门的连续任职段,每个连续段算1次。关键是先确保每个员工的记录在同一分区内按时间排序,再通过对比前后记录的部门ID来标记新的任职段,最后聚合计数。

步骤实现

  1. 数据预处理与分区排序
    先将日期字符串转为日期类型,再按emp_id分区(确保同一员工数据在同一Executor处理,避免跨Executor排序混乱),最后按begin日期排序。

  2. 标记新任职段
    使用窗口函数lag获取上一条记录的dept_id,如果当前dept_id与上一条不同,或是第一条记录,则标记为新的任职段开始。

  3. 聚合统计次数
    按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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 03:15:51