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

如何在Databricks中基于3天时间窗为DataFrame分组生成序列ID?

解决方案:按分组生成3天时间窗口序列ID

以下分别提供Pandas和PySpark两种实现方案,适配Databricks环境,解决你遇到的连续3天被分配不同序列ID的问题。


Pandas 实现

适用于中小数据量场景,逻辑清晰易调试:

import pandas as pd

# 确保日期列格式为datetime
df['date'] = pd.to_datetime(df['date'])

# 按分组字段+日期排序,保证窗口计算的顺序性
df = df.sort_values(['service', 'phone_number', 'date'])

# 计算当前行与同组上一行的日期差(首次行差值设为0)
df['date_diff'] = df.groupby(['service', 'phone_number'])['date'].diff().dt.days.fillna(0)

# 标记新窗口:当日期差超过2天时,视为新窗口起点
df['new_window'] = (df['date_diff'] > 2).astype(int)

# 累加新窗口标记生成序列ID(每组从1开始计数)
df['seq'] = df.groupby(['service', 'phone_number'])['new_window'].cumsum() + 1

# 生成窗口说明:每个序列对应的起止日期
window_details = df.groupby(['service', 'phone_number', 'seq']).agg(
    window_start=('date', 'min'),
    window_end=('date', 'max')
).astype({'window_start': str, 'window_end': str}).reset_index()

# 合并窗口说明到原数据,并格式化描述文本
df = df.merge(window_details, on=['service', 'phone_number', 'seq'])
df['window_desc'] = df.apply(lambda x: f"窗口{x['seq']}: {x['window_start']} 至 {x['window_end']}", axis=1)

# 清理中间临时列
df = df.drop(['date_diff', 'new_window', 'window_start', 'window_end'], axis=1)

逻辑说明

  • 核心是通过日期差判断窗口连续性:同组内当前日期与前一行日期差≤2天时,属于同一3天窗口;超过2天则开启新窗口。
  • 对2023-02-16、17、18这类连续日期,日期差均为1,因此会被分配同一个seq值。

PySpark 实现

适用于大数据量场景,适配Databricks分布式计算环境:

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

# 确保日期列格式为date类型
df = df.withColumn("date", F.to_date("date"))

# 定义分组排序窗口:按service、phone_number分组,按date升序排序
group_sort_window = Window.partitionBy("service", "phone_number").orderBy("date")

# 计算当前行与同组上一行的日期差(首次行差值设为0)
df = df.withColumn(
    "date_diff",
    F.coalesce(F.datediff(F.col("date"), F.lag("date").over(group_sort_window)), F.lit(0))
)

# 标记新窗口:日期差>2天时标记为1,否则0
df = df.withColumn("new_window", F.when(F.col("date_diff") > 2, 1).otherwise(0))

# 累加新窗口标记生成序列ID(每组从1开始)
df = df.withColumn(
    "seq",
    F.sum("new_window").over(group_sort_window.rangeBetween(Window.unboundedPreceding, 0)) + 1
)

# 生成窗口说明:每个序列对应的起止日期及格式化描述
window_details = df.groupBy("service", "phone_number", "seq")\
    .agg(
        F.min("date").alias("window_start"),
        F.max("date").alias("window_end")
    )\
    .withColumn(
        "window_desc",
        F.concat(
            F.lit("窗口"), F.col("seq"), F.lit(": "),
            F.date_format("window_start", "yyyy-MM-dd"),
            F.lit(" 至 "),
            F.date_format("window_end", "yyyy-MM-dd")
        )
    )

# 合并窗口说明到原数据
df = df.join(window_details, on=["service", "phone_number", "seq"], how="left")

# 清理临时列
df = df.drop("date_diff", "new_window", "window_start", "window_end")

# 查看结果
df.display()  # Databricks环境推荐用display替代show

逻辑说明

  • 用PySpark窗口函数lag获取同组前一行日期,datediff计算天数差,通过sum over窗口累加新窗口标记生成序列ID。
  • 完全适配分布式计算,处理亿级数据无压力,且能正确识别连续3天为同一窗口。

问题根源分析

你之前的代码(@mozway)可能误用了rank()/row_number()这类基于排序的排名函数,这类函数仅按日期顺序生成序号,未考虑时间窗口的连续性,导致连续3天被分配不同ID。本文方案通过日期差判断窗口边界,从根本上解决了这个问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 09:07:04