一. 前言:先決定答案掛在哪個時間點 #
時間序列最危險的 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:
- 是否先篩掉不需要的日期與欄位?
- 是否重複排序同一份資料?
period / every是否造成大量重疊?- 是否在正式流程一直輸出 boundary 欄位?
- 分組鍵是否基數極高?
- 是否把本可一次
.agg()的統計拆成多條視窗?
上線前,再逐項回答:
- 時間欄 dtype 與時區明確
- 每組時間索引遞增
- rolling 或 dynamic 的選擇有文件
-
every、period、offset有業務名稱 -
closed與label不靠預設猜測 - 空窗與補零政策明確
- 重複時間戳有處理規則
- 邊界、缺測、跨組測試已通過
- 升級 Polars 後重跑測試與效能基準
結語 #
Rolling 與 Dynamic Windows 的差別, 不在聚合函式, 而在「答案要掛在哪些時間點」。
每筆事件都要近期狀態, 用 rolling; 固定節奏要一列摘要, 用 dynamic。
然後把排序、邊界、時區、空窗與重複事件 都寫成可以測試的契約。 時間序列才不會只是在圖上看起來很合理。
拍拍君下次見。