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

如何在Polars中实现值中心滑动窗口?优先采用流式实现

在Polars中实现值中心滑动窗口(流式版本)

要实现你描述的“值中心”滑动窗口,我们可以分两种场景处理:小数据用表达式快速实现,大数据采用流式stateful批处理,无需加载全部数据到内存。

1. 小数据Eager模式实现

如果数据量不大,直接用条件表达式构造窗口即可,逻辑清晰高效:

import polars as pl

def sliding_window_eager(df: pl.DataFrame, col_name: str, window_size: int) -> pl.DataFrame:
    half = window_size // 2
    return df.with_columns(
        window=pl.when(pl.int_range(pl.count()).lt(half))
                # 前half个元素的窗口取前window_size个值
                .then(pl.col(col_name).head(window_size).list())
                .when(pl.int_range(pl.count()).ge(pl.count() - half))
                # 后half个元素的窗口取最后window_size个值
                .then(pl.col(col_name).tail(window_size).list())
                # 中间元素的窗口取当前位置前后half个元素
                .otherwise(pl.col(col_name).slice(pl.int_range(pl.count()) - half, window_size).list())
    ).select(col_name, "window")

# 示例使用
df = pl.DataFrame({"value": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]})
result = sliding_window_eager(df, "value", 3)
print(result)

输出结果:

shape: (10, 2)
┌───────┬─────────────┐
│ value ┆ window      │
│ ---   ┆ ---         │
│ i64   ┆ list[i64]   │
╞═══════╪═════════════╡
│ 1     ┆ [1, 2, 3]   │
│ 2     ┆ [1, 2, 3]   │
│ 3     ┆ [2, 3, 4]   │
│ 4     ┆ [3, 4, 5]   │
│ 5     ┆ [4, 5, 6]   │
│ 6     ┆ [5, 6, 7]   │
│ 7     ┆ [6, 7, 8]   │
│ 8     ┆ [7, 8, 9]   │
│ 9     ┆ [8, 9, 10]  │
│ 10    ┆ [8, 9, 10]  │
└───────┴─────────────┘

2. 大数据流式实现(Stateful Map Batches)

对于无法一次性加载到内存的大数据,用Polars的map_batches结合状态管理实现流式处理,核心是保留每个批次的尾部元素,确保跨批次窗口能正确构造:

import polars as pl
from typing import Optional, Tuple

def sliding_window_stream_processor(window_size: int):
    half = window_size // 2
    # 状态保留最后window_size个元素,用于构造尾部窗口
    state: Optional[pl.DataFrame] = None

    def processor(batch: pl.DataFrame) -> Tuple[pl.DataFrame, Optional[pl.DataFrame]]:
        nonlocal state
        # 合并上一批保留的状态与当前批次
        combined = pl.concat([state, batch]) if state is not None else batch

        # 可安全输出的行数:总行数减去half(避免窗口跨批次缺失元素)
        output_rows = combined.height - half
        if output_rows <= 0:
            # 元素不足,全部保留为状态(最多保留window_size个)
            state = combined.tail(window_size)
            return pl.DataFrame(), state

        # 为可输出的行生成对应窗口
        output = combined.slice(0, output_rows).with_columns(
            window=pl.col("value").slice(pl.int_range(output_rows) - half, window_size).list()
        )

        # 更新状态:保留最后window_size个元素
        state = combined.tail(window_size)
        return output, state

    return processor

# 示例使用(实际可替换为scan_csv/scan_parquet等流式数据源)
lf = pl.LazyFrame({"value": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]})
window_size = 3
processor = sliding_window_stream_processor(window_size)

# 应用流式处理
result_lf = lf.map_batches(processor, is_stateful=True)
final_result = result_lf.collect()

# 处理最后剩余的尾部元素
state = processor.__closure__[0].cell_contents
if state is not None and state.height > 0:
    tail_window = state.tail(window_size).select("value").to_series().to_list()
    tail_df = state.with_columns(window=pl.lit(tail_window))
    final_result = pl.concat([final_result, tail_df])

print(final_result)

流式实现说明

  • 状态管理:每个批次处理后保留最后window_size个元素,确保下一批次到来时能构造跨批次的窗口。
  • 边界处理:最后剩余的half个元素,其窗口直接取整个数据的最后window_size个值,与原生成器逻辑一致。
  • 内存效率:仅保留少量状态元素,无需加载全量数据,适合TB级别的流式数据处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 14:24:56