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

基于非重叠时间戳区间与阈值的DataFrame分组问题(Dask/Pandas)

解决方案

Pandas 实现

因为需求是顺序依赖的非重叠分组(每行归属唯一组,分组起始位置由前序处理结果决定),需要通过循环迭代逐个确定有效分组:

代码示例

import pandas as pd

# 示例数据
data = {
    'timestamp': [2317614314, 2317614409, 2317614429, 2317614608, 2317622053, 2317622208],
    'col1': [1.551823, 1.206112, 1.022747, 2.082569, 2.260681, 2.355770],
    'col2': [1,3,2,3,1,2],
    'col3': [4,3,3,3,2,4],
    'col4': [44,25,48,59,15,46]
}
df = pd.DataFrame(data)

# 参数配置
interval_ps = 200  # 时间区间(皮秒)
threshold = 4      # col1求和阈值

# 确保数据按timestamp排序(实际场景可能需要)
df = df.sort_values('timestamp').reset_index(drop=True)

groups = []
current_idx = 0
n_rows = len(df)

while current_idx < n_rows:
    # 获取当前起始行的时间戳
    start_ts = df.loc[current_idx, 'timestamp']
    end_ts = start_ts + interval_ps
    
    # 筛选当前区间内的所有行(从current_idx开始)
    mask = (df['timestamp'] <= end_ts) & (df.index >= current_idx)
    candidate_group = df[mask]
    
    # 计算col1总和
    col1_sum = candidate_group['col1'].sum()
    
    if col1_sum > threshold:
        # 符合条件,加入分组
        groups.append(candidate_group)
        # 跳过当前组的所有行
        current_idx = candidate_group.index[-1] + 1
    else:
        # 不符合条件,跳过当前行
        current_idx += 1

# 查看结果
for i, group in enumerate(groups, 1):
    print(f"分组 {i}:")
    print(group)
    print("---")

输出结果

分组 1:
     timestamp      col1  col2  col3  col4
1   2317614409  1.206112     3     3    25
2   2317614429  1.022747     2     3    48
3   2317614608  2.082569     3     3    59
---
分组 2:
     timestamp      col1  col2  col3  col4
4   2317622053  2.260681     1     2    15
5   2317622208  2.355770     2     4    46
---

Dask 实现

Dask的并行特性处理顺序依赖的分组需要额外注意,常规分区并行会打破顺序逻辑,可通过以下两种方式实现:

方法1:全局排序后用 delayed 串行处理

如果数据量不是极大,可先全局排序,再用dask.delayed模拟Pandas的循环逻辑:

import dask.dataframe as dd
from dask.delayed import delayed

# 加载Dask DataFrame(示例用Pandas数据转换)
ddf = dd.from_pandas(df, npartitions=1)  # 全局排序需单分区,或重分区后排序
ddf = ddf.sort_values('timestamp').reset_index(drop=True)

# 转换为延迟对象
df_delayed = ddf.to_delayed()[0]  # 取唯一分区的延迟对象

@delayed
def find_groups(df, interval_ps, threshold):
    groups = []
    current_idx = 0
    n_rows = len(df)
    while current_idx < n_rows:
        start_ts = df.loc[current_idx, 'timestamp']
        end_ts = start_ts + interval_ps
        mask = (df['timestamp'] <= end_ts) & (df.index >= current_idx)
        candidate_group = df[mask]
        col1_sum = candidate_group['col1'].sum()
        if col1_sum > threshold:
            groups.append(candidate_group)
            current_idx = candidate_group.index[-1] + 1
        else:
            current_idx += 1
    return groups

# 执行计算
result_groups = find_groups(df_delayed, interval_ps, threshold).compute()

# 输出结果
for i, group in enumerate(result_groups, 1):
    print(f"分组 {i}:")
    print(group)
    print("---")

方法2:分区内处理+边界修正(适合大数据量)

如果数据量极大无法单分区处理,可按timestamp分区后处理:

  1. 按timestamp范围分区,确保每个分区的时间范围远大于设定的interval_ps
  2. 每个分区内执行Pandas的分组逻辑,记录每个分组的起始/结束timestamp
  3. 合并分区结果,处理跨分区的候选组(比如前一个分区的末尾组可能包含后一个分区的起始行)
  4. 最终去重,确保每行仅归属一个组

这种方法需要额外的边界处理逻辑,适合超大规模数据集。

关键注意事项

  • timestamp类型:皮秒级精度超出Pandasdatetime64的纳秒限制,直接用数值类型(int/float)存储和比较即可,无需转换为时间类型。
  • 顺序依赖:分组逻辑依赖前序行的处理结果,无法用常规的groupby实现,必须通过迭代或串行处理。
  • Dask的局限性:并行处理这类顺序依赖任务会有性能损耗,若数据量不大,优先用Pandas实现;大数据量场景需结合分区策略和边界修正。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 13:34:55