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

任务变更时为人员分配RunID的Spark实现方案咨询

Spark中按人员任务连续批次分配全局RunID的解决方案

要实现的效果:给执行任务的人员分配RunID,规则是:

  • 同一个人连续做同一个任务,RunID不变
  • 任务变了(包括之后再做之前做过的任务),要分配新的RunID
  • 不同人的RunID要区分开,并且全局是递增的

示例数据里已经给出了期望的expectedRunID:

import pyspark.sql.functions as F
from pyspark.sql import Window

df = spark.createDataFrame(sc.parallelize([
            ['A', 'T1', '2023-01-01', 1],
            ['A', 'T1', '2023-01-02', 1],
            ['A', 'T2', '2023-01-03', 2],
            ['A', 'T2', '2023-01-04', 2],
            ['A', 'T2', '2023-01-05', 2],
            ['A', 'T1', '2023-01-06', 3],
            ['A', 'T1', '2023-01-07', 3],
            ['A', 'T1', '2023-01-08', 3],
            ['A', 'T1', '2023-01-09', 3],
            ['A', 'T1', '2023-01-10', 3],
            ['B', 'T1', '2023-01-01', 4],
            ['B', 'T1', '2023-01-02', 4],
            ['B', 'T1', '2023-01-03', 4],
            ['B', 'T1', '2023-01-04', 4],
            ['B', 'T1', '2023-01-05', 4],
            ]),
            ['Person', 'Task', 'Time', 'expectedRunID'])

原方案的问题

之前用dense_rank()按Person分区、Task排序的写法:

window = Window.partitionBy("Person").orderBy("Task")
df.withColumn("runID", F.dense_rank().over(window))

这个方法行不通,因为dense_rank只看Task的唯一值排序,不管时间顺序和任务的连续执行批次——比如人员A第二次做T1时,会被判定和第一次T1是同一个等级,RunID还是1,根本出不来期望的3。

正确解决方案

得分成两步来做:

  1. 先给每个人生成内部的连续批次ID:识别同一个人连续做的任务批次,任务一变,批次ID就加1
  2. 再基于每个人的内部批次ID,生成全局唯一递增的RunID

具体代码如下:

# 第一步:先按人员分组,按时间排序,定义窗口
person_window = Window.partitionBy("Person").orderBy("Time")

# 生成任务变更标记:如果当前任务和上一个不一样,或者是这个人的第一条记录,就标记为1,否则0
df = df.withColumn(
    "change_flag",
    F.when(
        F.lag("Task").over(person_window).isNull() | (F.col("Task") != F.lag("Task").over(person_window)),
        1
    ).otherwise(0)
)

# 对标记做累计求和,得到每个人内部的批次ID
df = df.withColumn(
    "person_run_id",
    F.sum("change_flag").over(person_window.rowsBetween(Window.unboundedPreceding, Window.currentRow))
)

# 第二步:生成全局RunID,按人员和内部批次排序,分配递增的ID
global_window = Window.orderBy("Person", "person_run_id")
df = df.withColumn(
    "runID",
    F.dense_rank().over(global_window)
)

# 查看结果,runID会和expectedRunID完全匹配
df.select("Person", "Task", "Time", "expectedRunID", "runID").show()

结果验证

执行后runID列和示例里的expectedRunID完全一致:

  • 人员A的T1(1-2号)→ runID=1,T2(3-5号)→ runID=2,再次执行的T1(6-10号)→ runID=3
  • 人员B的T1(1-5号)→ runID=4

内容的提问来源于stack exchange,提问作者WilliamEllisWebb

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 07:44:59