↓快轉到主要內容
  1. 教學文章/

Polars Rolling / Dynamic Windows:時間窗聚合、邊界語意與 QA

·6 分鐘· loading · loading · ·
Python Polars Time-Series Rolling-Window Dynamic-Window Data-Engineering Testing
每日拍拍
作者
每日拍拍
科學家 X 科技宅宅
目錄
Python 學習 - 本文屬於一個選集。
§ 131: 本文

一. 前言:先決定答案掛在哪個時間點
#

時間序列最危險的 bug, 通常不是程式直接噴錯, 而是它很安靜地算出一張看起來合理的表。

嗨,我是拍拍君。 今天不做普通的 group_by, 而是處理兩種常被混在一起的時間窗:

  • 每一筆事件都要一個「往回看 30 分鐘」的統計;
  • 每 15 分鐘產生一列固定區間摘要。

前者適合 rolling window, 後者適合 dynamic window。 兩者都叫時間窗, 但輸出列數、錨點與邊界語意完全不同。

如果你還不熟 Polars Expression API, 先看Polars 基礎教學; 如果目標是大型 Parquet 的 pushdown 與 streaming, 則接著看 LazyFrame 最佳化。 這篇只把時間窗的契約講清楚。

二. 安裝與版本確認
#

建立一個小專案:

uv init polars-windows
cd polars-windows
uv add polars pytest
uv run python -c "import polars as pl; print(pl.__version__)"

本文使用現行 API:

  • DataFrame.rolling(...)
  • DataFrame.group_by_dynamic(...)
  • group_by= 作為額外分組參數

網路上的舊範例可能還會出現 group_by_rolling、groupby_dynamic、by= 或 check_sorted=。 不要直接複製; 先對照你安裝版本的官方文件。

三. 準備不規則事件資料
#

先做一份兩台裝置的溫度事件。 資料刻意不是固定頻率, 而且原始順序也被打亂:

from datetime import datetime

import polars as pl

events = pl.DataFrame(
    {
        "device": ["beta", "alpha", "alpha", "beta", "alpha", "beta"],
        "ts": [
            datetime(2026, 10, 2, 9, 38),
            datetime(2026, 10, 2, 9, 3),
            datetime(2026, 10, 2, 9, 17),
            datetime(2026, 10, 2, 9, 8),
            datetime(2026, 10, 2, 9, 44),
            datetime(2026, 10, 2, 9, 21),
        ],
        "temp_c": [23.8, 21.0, 22.0, 23.0, 24.0, 23.2],
        "quality": ["ok", "ok", "ok", "ok", "ok", "warn"],
    }
)

events = events.sort("device", "ts")
print(events)

時間窗的第一條契約是: 索引欄必須遞增排序。

有 group_by="device" 時, 要求變成每個 device 內的 ts 都要遞增。 因此不要只憑肉眼看到全表日期大致有序, 要把分組鍵和時間鍵一起寫進 sort。

排序是資料契約, 不是為了讓輸出看起來漂亮。

四. Rolling 與 Dynamic:先選對問題
#

兩種窗可以用一句話區分:

問題 API 視窗錨點 常見輸出列數
每筆事件發生時,回看一段時間 rolling 每筆索引值 約等於輸入列數
每隔固定時間產生摘要 group_by_dynamic 固定排程格線 等於非空視窗數

假設事件時間是 09:03、09:17、09:44。

rolling(period="30m") 會分別建立:

  • (08:33, 09:03]
  • (08:47, 09:17]
  • (09:14, 09:44]

但 group_by_dynamic(every="15m") 會建立像 [09:00, 09:15)、[09:15, 09:30) 的固定格子。

選擇關鍵不是「我要算平均」, 而是「我要在哪些時間點輸出答案」。

五. Dynamic Window:固定節奏的時間摘要
#

先做每 15 分鐘一列的裝置摘要:

dynamic_15m = (
    events.group_by_dynamic(
        "ts",
        every="15m",
        group_by="device",
        closed="left",
        label="left",
    )
    .agg(
        pl.len().alias("samples"),
        pl.col("temp_c").mean().round(2).alias("mean_temp"),
        pl.col("temp_c").min().alias("min_temp"),
        pl.col("temp_c").max().alias("max_temp"),
    )
    .sort("device", "ts")
)

print(dynamic_15m)

這裡的 every="15m" 表示: 每 15 分鐘啟動一個窗。

沒有指定 period 時, 視窗長度預設等於 every, 因此得到互不重疊的 15 分鐘格子。

label="left" 則表示結果的 ts 使用視窗左邊界作為標籤。 這個欄位不是「該組第一筆事件時間」。

六. every、period、offset 是三個旋鈕
#

group_by_dynamic 最重要的三個參數:

  • every:多久啟動一個新窗;
  • period:每個窗有多長;
  • offset:整套格線平移多少。

例如 every="15m", period="30m" 代表每 15 分鐘啟動一個 30 分鐘窗。因為 period > every, 相鄰視窗會重疊, 同一筆事件可能屬於多個窗。 這不是重複資料, 而是你明確要求的視窗模型。

反過來,若 period < every, 格線之間會出現空隙, 落在空隙中的事件不屬於任何窗。

offset 適合處理不是整點起算的營運週期。 例如每小時的班次在 10 分開始, 可以用 every="1h", offset="10m"。

七. closed 與 label:邊界不要靠猜
#

時間恰好落在 09:15 時, 它應該算前一窗還是後一窗? 答案由 closed 決定:

closed 左邊界 右邊界
"left" 包含 不包含
"right" 不包含 包含
"both" 包含 包含
"none" 不包含 不包含

資料管線常用左閉右開 [start, end), 因為相鄰固定窗不會同時吃到同一個邊界事件。

除錯時可以暫時加上 include_boundaries=True,再把 _lower_boundary、_upper_boundary 與成員時間一起印出來。這很適合 QA, 但官方文件提醒它會增加成本, 因為邊界輸出較難平行化。 確認語意後,不一定要留在正式結果裡。

八. 空窗不會自動出現
#

group_by_dynamic 只回傳有資料的視窗。 如果 10:00 到 10:15 沒有事件, 結果不會自動補一列 samples = 0。

這裡要先決定產品語意:

  • 沒有列:代表沒有事件;
  • 有列但值是 null:代表格線存在、沒有觀測;
  • 有列且值是 0:代表業務定義允許把缺測視為零。

需要完整格線時, 用 pl.datetime_range(..., interval="15m", eager=True) 先建立時間骨架再 join, 或評估 upsample()。 不要把缺列直接 fill_null(0), 因為缺列根本還不存在。

是否補零, 應由下游報表契約決定。

九. Rolling Window:每筆事件都回看
#

如果需求是「每次讀值到達時,計算過去 30 分鐘狀態」, 就用 rolling:

rolling_30m = (
    events.rolling(
        "ts",
        period="30m",
        group_by="device",
        closed="right",
    )
    .agg(
        pl.len().alias("samples_30m"),
        pl.col("temp_c").mean().round(2).alias("mean_30m"),
        pl.col("temp_c").max().alias("max_30m"),
        (pl.col("quality") == "warn").sum().alias("warnings_30m"),
    )
    .sort("device", "ts")
)

print(rolling_30m)

預設 offset=-period, 搭配 closed="right" 時, 每筆 t 的窗近似 (t - period, t]。

這很適合:

  • 每次交易發生時的近一小時成交量;
  • 每筆感測值的近 30 分鐘平均;
  • 每次請求到達時的近期錯誤率。

注意輸出 ts 是原本每筆事件的索引值, 不是固定排程標籤。

十. 分組與重複時間戳
#

多裝置資料一定要指定 group_by, 否則 alpha 的歷史會混進 beta 的視窗。

同一裝置也可能在同一微秒有兩筆事件。 Polars 可以聚合它們, 但你仍要先定義語意:

  • 兩筆都是有效樣本;
  • 後到資料覆蓋前一筆;
  • 依 event_id 去重;
  • 依 ingestion sequence 決定正式版本。

例如可以先保留 ingestion sequence,再用 unique(subset=["device", "ts"], keep="last") 選正式版本。不要把 unique() 當成無害清潔;它會改變 sample count、平均值與警報比例, 所以去重規則必須寫進資料契約。

十一. 時區與 DST:先分清楚「轉換」和「指定」
#

跨時區事件最好以 UTC 保存, 需要依在地日曆分窗時才轉到業務時區:

utc_events = pl.DataFrame(
    {
        "ts": pl.datetime_range(
            datetime(2026, 10, 2, 0, 0),
            datetime(2026, 10, 2, 3, 0),
            interval="30m",
            time_zone="UTC",
            eager=True,
        ),
        "value": [10, 11, 12, 13, 14, 15, 16],
    }
)

taipei_events = utc_events.with_columns(
    pl.col("ts").dt.convert_time_zone("Asia/Taipei")
)

convert_time_zone() 保留同一個瞬間, 只改變顯示與日曆時區。

replace_time_zone() 則會重解釋牆上時間, 並改變底層 timestamp。 它適合原始字串「本來就是某地時間」的情境, 不是一般的 UTC 轉在地時間。

日曆單位也不是固定秒數:

  • "1d" 是下一個日曆日的同一時間;
  • "24h" 是固定 24 小時;
  • 遇到 DST 的地區,兩者可能不同。

因此「每日報表」與「最近 24 小時」 不要共用一個模糊需求名稱。

十二. QA:用邊界資料證明結果
#

時間窗測試不應只用「每五分鐘一筆」的完美資料。 至少要包含:

  • 剛好落在左右邊界的事件;
  • 同一時間的重複事件;
  • 分組交錯但組內有序的事件;
  • 中間缺一整個窗;
  • 時區轉換或 DST 邊界;
  • 單一事件與空資料集。

先測最容易出錯的左閉右開語意:

from polars.testing import assert_frame_equal


def test_dynamic_window_is_left_closed() -> None:
    source = pl.DataFrame(
        {
            "ts": [
                datetime(2026, 10, 2, 9, 0),
                datetime(2026, 10, 2, 9, 15),
                datetime(2026, 10, 2, 9, 29),
            ],
            "value": [1, 10, 100],
        }
    )

    actual = source.group_by_dynamic(
        "ts",
        every="15m",
        closed="left",
        label="left",
    ).agg(
        pl.len().alias("samples"),
        pl.col("value").sum().alias("total"),
    )

    expected = pl.DataFrame(
        {
            "ts": [
                datetime(2026, 10, 2, 9, 0),
                datetime(2026, 10, 2, 9, 15),
            ],
            "samples": [1, 2],
            "total": [1, 110],
        },
        schema_overrides={"samples": pl.UInt32},
    )

    assert_frame_equal(actual, expected)

再測 rolling 的輸出基數和分組隔離:

def test_rolling_keeps_one_row_per_event() -> None:
    actual = events.rolling(
        "ts",
        period="30m",
        group_by="device",
        closed="right",
    ).agg(pl.len().alias("samples"))

    assert actual.height == events.height

    first_per_device = (
        actual.sort("device", "ts")
        .group_by("device", maintain_order=True)
        .first()
    )
    assert first_per_device["samples"].to_list() == [1, 1]

最後加上不依賴完整 snapshot 的 invariant:

assert dynamic_15m["samples"].sum() == events.height
assert dynamic_15m["samples"].min() >= 1
assert rolling_30m.height == events.height
assert rolling_30m["samples_30m"].min() >= 1

第一個總和 invariant 只適用於 不重疊、沒有 gap 的 dynamic windows。 若 period != every, 同一事件可能被算多次或完全沒被算到, 測試也要跟著契約改。

十三. 效能與上線檢查
#

時間窗慢時, 先檢查資料形狀, 不要急著改成 Python loop:

  1. 是否先篩掉不需要的日期與欄位?
  2. 是否重複排序同一份資料?
  3. period / every 是否造成大量重疊?
  4. 是否在正式流程一直輸出 boundary 欄位?
  5. 分組鍵是否基數極高?
  6. 是否把本可一次 .agg() 的統計拆成多條視窗?

上線前,再逐項回答:

  • 時間欄 dtype 與時區明確
  • 每組時間索引遞增
  • rolling 或 dynamic 的選擇有文件
  • every、period、offset 有業務名稱
  • closed 與 label 不靠預設猜測
  • 空窗與補零政策明確
  • 重複時間戳有處理規則
  • 邊界、缺測、跨組測試已通過
  • 升級 Polars 後重跑測試與效能基準

結語
#

Rolling 與 Dynamic Windows 的差別, 不在聚合函式, 而在「答案要掛在哪些時間點」。

每筆事件都要近期狀態, 用 rolling; 固定節奏要一列摘要, 用 dynamic。

然後把排序、邊界、時區、空窗與重複事件 都寫成可以測試的契約。 時間序列才不會只是在圖上看起來很合理。

拍拍君下次見。

延伸閱讀
#

Python 學習 - 本文屬於一個選集。
§ 131: 本文

相關文章

Polars LazyFrame 最佳化:Pushdown、Streaming 與 Explain
·6 分鐘· loading · loading
Python Polars LazyFrame Query Optimization Parquet Streaming Data-Engineering
DuckDB ASOF Join 實戰:時間序列對齊與最近事件查詢
·7 分鐘· loading · loading
Python DuckDB SQL ASOF Join Time-Series Temporal Data Data-Engineering
Python fsspec 實戰:統一讀寫本機、S3、HTTP 與資料管線路徑
·7 分鐘· loading · loading
Python Fsspec Filesystem S3 Data-Engineering ETL
Python pytest fixtures 進階:conftest、factory 與測試資料管理
·8 分鐘· loading · loading
Python Pytest Fixtures Testing Conftest Monkeypatch Developer-Tools
Textual Command Palette + Actions:快捷鍵、命令搜尋與可測試操作
·8 分鐘· loading · loading
Python Textual TUI Command-Palette Actions Key-Bindings Testing
Python JSON 實戰:解析、序列化、自訂型別與大型資料
·7 分鐘· loading · loading
Python Json Serialization Standard-Library JSON Lines Data-Engineering Developer-Tools