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

如何将Pandas分组3天间隔序列函数转为无UDF的PySpark实现

解决方案:纯PySpark实现动态3天间隔分组序列

核心思路

原Pandas函数的逻辑是按service+phone_number分组,遍历排序后的日期:若当前日期与上一次触发新序列的基准日期间隔超过3天,则生成新序列号并更新基准日期;否则沿用当前序列号。这种状态依赖的逻辑可以通过Spark内置的aggregate函数处理分组后的日期数组实现,无需UDF。

完整代码实现

from pyspark.sql import functions as F
from pyspark.sql.types import DateType

# 1. 转换日期格式(假设输入为dd/MM/yyyy字符串,可根据实际格式调整)
df = df.withColumn("date", F.to_date("date", "dd/MM/yyyy"))

# 2. 分组并收集排序后的日期数组
grouped_df = df.groupBy("service", "phone_number").agg(
    F.sort_array(F.collect_list("date")).alias("sorted_dates")
)

# 3. 定义聚合逻辑:模拟原Pandas函数的循环判断
def accumulate(acc, date):
    # 判断是否需要生成新序列
    need_new_seq = date > F.date_add(acc.last_ref, 3)
    # 更新序列号和基准日期
    new_seq_num = F.when(need_new_seq, acc.seq_num + 1).otherwise(acc.seq_num)
    new_last_ref = F.when(need_new_seq, date).otherwise(acc.last_ref)
    # 记录当前日期对应的序列号
    new_seq_list = F.array_union(
        acc.seq_list,
        F.array(F.struct(date.alias("date"), new_seq_num.alias("seq")))
    )
    return F.struct(new_last_ref.alias("last_ref"), new_seq_num.alias("seq_num"), new_seq_list.alias("seq_list"))

# 初始状态:基准日期设为1970-01-01,序列号初始为0,结果列表为空
initial_state = F.struct(
    F.lit("1970-01-01").cast(DateType()).alias("last_ref"),
    F.lit(0).alias("seq_num"),
    F.array().alias("seq_list")
)

# 4. 对每个分组的日期数组应用聚合,生成日期-序列号映射
grouped_seq_map = grouped_df.withColumn(
    "agg_result",
    F.aggregate("sorted_dates", initial_state, accumulate)
).select(
    "service", "phone_number", F.explode("agg_result.seq_list").alias("seq_info")
).select(
    "service", "phone_number", "seq_info.date", "seq_info.seq"
)

# 5. 关联原始数据,得到最终结果
final_df = df.join(
    grouped_seq_map,
    on=["service", "phone_number", "date"],
    how="left"
).withColumnRenamed("seq", "seq_pandas")

逻辑验证

针对示例数据,该代码会生成与seq_pandas列完全一致的结果:

  • 同一分组内连续日期间隔≤3天时,序列号保持不变(如BBBB组的11/12/2021与09/12/2021)
  • 日期间隔>3天时,序列号递增(如BBBB组的14/01/2022与09/12/2021)
  • 相同日期的行共享同一序列号(如BBBB组的两行13/04/2022)
  • 不同分组独立生成序列(如AAAA组的23/05/2022重新从1开始)

注意事项

  • 确保日期格式正确,若输入日期格式不是dd/MM/yyyy,请修改to_date函数的第二个参数
  • 该方法基于Spark内置函数,无UDF依赖,完全兼容Unity Catalog Databricks Runtime 13.1

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 11:22:27