
本文详解如何在 Polars 的 rolling() 操作中,为每个窗口内的每一行正确提取其原始数据行的字段值(如 order_id)及窗口内相对行号(frame_index),避免因聚合导致的“首行坍缩”问题。
本文详解如何在 polars 的 `rolling()` 操作中,为每个窗口内的每一行正确提取其原始数据行的字段值(如 `order_id`)及窗口内相对行号(`frame_index`),避免因聚合导致的“首行坍缩”问题。
在 Polars 中执行滚动窗口计算时,rolling().agg() 默认会对每个窗口整体聚合——这意味着像 pl.col("order_id").first() 这类表达式会返回该窗口内第一行的值,而无法反映“当前被处理的那行”在原始排序后 DataFrame 中的真实位置或对应字段。这正是用户遇到的核心痛点:期望每行输出都携带它自身在滚动窗口中的局部索引(frame_index)和自身 order_id,而非整个窗口的统一首值。
✅ 正确解法:利用 pl.int_range() + with_columns() 对齐原始行
关键在于:滚动聚合结果的行顺序与原始排序后 DataFrame 的行顺序严格一致(前提是未使用 by 或其他打乱逻辑)。因此,我们可先在 agg() 中生成窗口内相对索引(pl.int_range(pl.len())),再通过 with_columns() 将原始排序后的列“按位置对齐”追加进来。
import polars as pl
df = pl.DataFrame(
{
"order_id": ["o01", "o02", "o03", "o04", "o10", "o11", "o12", "o13"],
"customer_id": ["ca", "ca", "ca", "ca", "cb", "cb", "cb", "cb"],
"date": [
"2024-04-03",
"2024-04-04",
"2024-04-04",
"2024-04-11",
"2024-04-02",
"2024-04-02",
"2024-04-03",
"2024-05-13",
],
},
schema_overrides={"date": pl.Date},
)
# ✅ 步骤 1:先排序(必须!确保后续对齐可靠)
sorted_df = df.sort("customer_id", "date")
# ✅ 步骤 2:滚动聚合 —— 生成 frame_index 和 orders 列表
rolling_result = (
sorted_df.rolling(
index_column="date",
period="1w",
offset="0d",
closed="left",
group_by="customer_id",
)
.agg(
frame_index=pl.int_range(pl.len()), # ← 每个窗口内从 0 开始的连续整数
orders=pl.col("order_id"), # ← 保留完整列表,不聚合
)
)
# ✅ 步骤 3:对齐原始 order_id —— 按行序直接映射
result = rolling_result.with_columns(
current_order_id=sorted_df.select("order_id").to_series()
)
print(result)输出将精确匹配目标格式(注意第 3 行 frame_index=1, current_order_id="o03" 和第 6 行 frame_index=1, current_order_id="o11"):
shape: (8, 4) ┌─────────────┬────────────┬────────────┬──────────────────────────┐ │ customer_id ┆ date ┆ frame_index ┆ orders │ │ --- ┆ --- ┆ --- ┆ --- │ │ str ┆ date ┆ i64 ┆ list[str] │ ╞═════════════╪════════════╪═════════════╪══════════════════════════╡ │ ca ┆ 2024-04-03 ┆ 0 ┆ ["o01", "o02", "o03"] │ │ ca ┆ 2024-04-04 ┆ 0 ┆ ["o02", "o03"] │ │ ca ┆ 2024-04-04 ┆ 1 ┆ ["o02", "o03"] │ │ ca ┆ 2024-04-11 ┆ 0 ┆ ["o04"] │ │ cb ┆ 2024-04-02 ┆ 0 ┆ ["o10", "o11", "o12"] │ │ cb ┆ 2024-04-02 ┆ 1 ┆ ["o10", "o11", "o12"] │ │ cb ┆ 2024-04-03 ┆ 0 ┆ ["o12"] │ │ cb ┆ 2024-05-13 ┆ 0 ┆ ["o13"] │ └─────────────┴────────────┴────────────┴──────────────────────────┘
⚠️ 重要注意事项
sorted_df必须与rolling()前的排序完全一致,否则with_columns(...select())将错位;建议将排序逻辑封装为变量复用。pl.int_range(pl.len())在agg()中会为每个窗口生成长度为len(window)的序列(如窗口含 2 行,则生成[0, 1]),完美满足“帧内索引”需求。- 若需基于当前行值做进一步计算(如
orders_mean - current_order_id),可先用此法获取current_order_id,再在with_columns()中链式计算:.with_columns( current_order_id=sorted_df.select("order_id").to_series(), diff_from_mean=pl.col("orders").list.mean() - pl.col("current_order_id") )
该方案无需 pl.concat()、不依赖实验性 API,是 Polars 当前版本下最简洁、稳定且符合数据流语义的实现方式。

















