Polarsの時系列APIを実務で使い切る ― `group_by_dynamic`と`join_asof`を仮想通貨データで学ぶ

約31分で読めます by ぽんたぬき
Polarsの時系列APIを実務で使い切る ― `group_by_dynamic`と`join_asof`を仮想通貨データで学ぶ

Polarsの時系列APIを実務で使い切る ― group_by_dynamicjoin_asofを仮想通貨データで学ぶ

はじめに ― 「入門の次」でつまずくのは時系列APIだ

Polarsの基本操作(select / filter / with_columns / group_by)を覚えた後、多くの方が最初に壁を感じるのが時系列専用APIです。

Polars入門では時系列処理の概要を紹介しましたが、実務で必要なgroup_by_dynamicのウィンドウ境界の扱い、join_asofの方向指定、欠損バーの補完といった踏み込んだ内容には触れていませんでした。本記事はその続編として、時系列専用APIだけに絞った深掘り解説をお届けします。

本記事で扱うこと・扱わないこと

  • ✅ 扱う:group_by_dynamic / upsample / rolling / over / join_asof の5つに絞った実践解説
  • ❌ 扱わない:Polarsの基本操作(→ 記事142へ)、pandasからの書き換え対比(→ 記事248へ)

想定読者: selectwith_columnsは書けるが、時系列APIをまだ使いこなせていない中級者。[記事227〜231(CCXT)]でデータ取得はできているが、そのデータを使った前処理で詰まっている方を特に想定しています。

本記事の核心 ― join_asofの方向を誤ると成績が「良くなる」

先に結論を述べます。join_asofstrategyを誤って"forward"にすると、未来の値を参照するlook-ahead bias(先読みバイアス)が発生します。問題なのはバックテストの成績が悪化するのではなく、異常に向上してしまう点です。気づきにくいまま誤った戦略を実運用に回すリスクがあります。詳細はSTEP 5と専用章で数値付きで解説します。

検証環境とサンプルデータの準備

# 検証環境
# Polars 1.x系(APIは変更が多いためバージョンを必ず固定してください)
# pip install polars==1.8.0

import polars as pl
import numpy as np
from datetime import datetime, timezone

print(pl.__version__)  # 1.8.0 などを確認

# ダミーOHLCVデータの生成(1分足 × 200本)
np.random.seed(42)
n = 200
base_price = 50000.0

timestamps_ms = [
    int(datetime(2024, 1, 15, 9, 0, tzinfo=timezone.utc).timestamp() * 1000)
    + i * 60_000
    for i in range(n)
]
prices = base_price + np.cumsum(np.random.randn(n) * 100)

df_1m = pl.DataFrame({
    "timestamp_ms": timestamps_ms,
    "open":   prices,
    "high":   prices + np.abs(np.random.randn(n) * 50),
    "low":    prices - np.abs(np.random.randn(n) * 50),
    "close":  prices + np.random.randn(n) * 30,
    "volume": np.abs(np.random.randn(n) * 10 + 5),
    "symbol": ["BTC/USDT"] * n,
})

# CCXTが返すエポックミリ秒を Datetime(time_unit="ms", UTC) に変換
df = df_1m.with_columns(
    pl.from_epoch("timestamp_ms", time_unit="ms")
      .dt.replace_time_zone("UTC")
      .alias("datetime")
)

print(df.schema)
# {'timestamp_ms': Int64, 'datetime': Datetime(time_unit='ms', time_zone='UTC'), ...}

# 文字列タイムスタンプのパース例
df_str = pl.DataFrame({"ts": ["2024-01-15 09:00:00", "2024-01-15 09:01:00"]})

df_parsed = df_str.with_columns(
    pl.col("ts")
      .str.to_datetime(format="%Y-%m-%d %H:%M:%S", time_unit="ms")
      .dt.replace_time_zone("UTC")
      .alias("datetime")
)

STEP 1 ― 時刻カラムの型を制する(すべての土台)

Datetime型とtime_unitの罠

PolarsのDatetime型にはtime_unit"ms" / "us" / "ns")とtime_zoneの2つの属性があります。これが揃っていないとjoin_asofが型エラーで失敗します。 まず確認する癖をつけましょう。

# CCXTが返すエポックミリ秒を Datetime(time_unit="ms", UTC) に変換
df = df_1m.with_columns(
    pl.from_epoch("timestamp_ms", time_unit="ms")
      .dt.replace_time_zone("UTC")
      .alias("datetime")
)

print(df.schema)
# {'timestamp_ms': Int64, 'datetime': Datetime(time_unit='ms', time_zone='UTC'), ...}

文字列からパースする場合はstr.to_datetime()を使います。

# 文字列タイムスタンプのパース例
df_str = pl.DataFrame({"ts": ["2024-01-15 09:00:00", "2024-01-15 09:01:00"]})

df_parsed = df_str.with_columns(
    pl.col("ts")
      .str.to_datetime(format="%Y-%m-%d %H:%M:%S", time_unit="ms")
      .dt.replace_time_zone("UTC")
      .alias("datetime")
)

replace_time_zone vs convert_time_zone の決定的な違い

メソッド 動作 使う場面
replace_time_zone("UTC") タイムゾーン情報を付与するだけ(数値は変わらない) naive → aware への変換
convert_time_zone("Asia/Tokyo") UTCをJST時刻に変換(数値が+9h分変わる) 表示用に変換するとき

推奨方針:内部処理はUTC固定、表示直前だけJSTに変換。 混在させると静かなバグの温床になります。

ソート済みであることが全ての前提

group_by_dynamic / rolling / join_asof はすべて時刻カラムがソート済みであることを前提にします。ここを飛ばすと後続の処理が全て壊れます。

# sort は必ず先に行う
df = df.sort("datetime")

STEP 2 ― group_by_dynamicでリサンプリングする

最小構成 ― 1分足を5分足にする

まず動くコードを見せ、その後で引数を解剖します。

df_5m = (
    df
    .sort("datetime")
    .group_by_dynamic(
        index_column="datetime",
        every="5m",       # 5分ごとに窓を切り出す
        closed="left",    # 窓の左端を含む(右端は含まない)
        label="left",     # 窓の開始時刻をラベルとする
        group_by=["symbol"],
    )
    .agg([
        pl.col("open").first(),
        pl.col("high").max(),
        pl.col("low").min(),
        pl.col("close").last(),
        pl.col("volume").sum(),
    ])
    .sort(["symbol", "datetime"])
)

print(df_5m.head(3))

every / period / offsetの役割

引数 説明
every 次の窓の開始点までの間隔 "5m"
period 1つの窓の長さ(省略時はeveryと同じ) "10m"everyより大きくするとオーバーラップ窓)
offset グリッドの起点をずらす "2m"(00,05,...ではなく02,07,...開始)

every="5m", period="10m"にすると各窓が10分幅で5分ごとに移動するオーバーラップウィンドウになります。テクニカル指標の前処理などで使う応用パターンです。

closedが変える境界の含み方

09:05:00のバーはどの窓に入るか——これがclosedで変わります。

closed 窓の範囲(every="5m"の場合) 09:05:00の所属
"left"(デフォルト) [09:00, 09:05) 次の窓(09:05〜09:10)に入る
"right" (09:00, 09:05] この窓(09:00〜09:05)に入る
"both" [09:00, 09:05] 両方の窓に入る(重複)
"none" (09:00, 09:05) どちらにも入らない

OHLCVのリサンプリングではclosed="left"が自然です。「09:00〜09:04:59の取引を09:00の足に集約する」という意味になります。

labelと look-aheadの伏線

label="left"にすると09:00〜09:05の集計結果はdatetime=09:00として出力されます。これを09:00時点の特徴量として使うと、09:01〜09:04の情報を参照することになり、look-ahead biasが発生します。 この問題はSTEP 5の専用章で詳しく再現します。

実装例 ― 1分足を5分足・1時間足に集約

def resample_ohlcv(df: pl.DataFrame, every: str) -> pl.DataFrame:
    """OHLCV 1分足を指定間隔にリサンプリングする"""
    return (
        df
        .sort("datetime")
        .group_by_dynamic(
            index_column="datetime",
            every=every,
            closed="left",
            label="left",
            group_by=["symbol"],
        )
        .agg([
            pl.col("open").first(),
            pl.col("high").max(),
            pl.col("low").min(),
            pl.col("close").last(),
            pl.col("volume").sum(),
            # VWAP = sum(close * volume) / sum(volume)
            (pl.col("close") * pl.col("volume")).sum().alias("price_vol"),
        ])
        .with_columns(
            (pl.col("price_vol") / pl.col("volume")).alias("vwap")
        )
        .drop("price_vol")
        .sort(["symbol", "datetime"])
    )

df_5m  = resample_ohlcv(df, "5m")
df_1h  = resample_ohlcv(df, "1h")

STEP 3 ― 欠損バーを埋める(upsample × forward_fill

なぜ仮想通貨データにバーが欠けるのか

仮想通貨の1分足データには以下の理由でバーが欠落します。

  • 約定ゼロの分は行自体が存在しない
  • 取引所のメンテナンス時間帯
  • CCXTのページング取得で一部が抜ける([記事227〜231]参照)

欠落したままrollingを掛けると、窓の実時間がバラつき、移動平均の計算が狂います。

upsampleで時間グリッドを埋める

# 1分グリッドに揃える(欠落した分はnullで埋まる)
df_up = (
    df
    .sort("datetime")
    .upsample(time_column="datetime", every="1m", group_by=["symbol"])
)
print(df_up.filter(pl.col("open").is_null()).head(3))
# 欠損バーはopenからvolumeまで全カラムがnullになっている

カラムごとに埋め方を変えるのが正解

df_filled = (
    df_up
    .sort(["symbol", "datetime"])
    .with_columns([
        # 価格系:直前の終値を引き継ぐ(forward_fill)
        pl.col("open").forward_fill().over("symbol"),
        pl.col("high").forward_fill().over("symbol"),
        pl.col("low").forward_fill().over("symbol"),
        pl.col("close").forward_fill().over("symbol"),
        # 出来高・約定数:ゼロで埋める(forward_fillすると出来高を捏造する)
        pl.col("volume").fill_null(0),
    ])
)

「全部forward_fill」がなぜ危険か: volumeをforward_fillすると、取引が一切なかった時間帯に「直前の出来高」が入ります。これはOBV(On-Balance Volume)やVWAPの計算を完全に破壊します。

forward_fillできない先頭行の扱い

系列の先頭が欠損している場合、forward_fillは補完できません(参照すべき過去データがない)。backward_fillを安易に使うと未来の値を参照するため禁止です。欠損フラグ列を残しておく運用が安全です。

df_filled = df_filled.with_columns(
    pl.col("close").is_null().alias("is_gap_bar")  # 欠損フラグを残す
)

STEP 4 ― rolling集計とover()による銘柄別ウィンドウ計算

group_by_dynamicrollingの使い分け

特徴 group_by_dynamic rolling + over
出力 集計テーブル(行数が減る) 元テーブルに列を追加
用途 足種変換(1m→5m) 特徴量生成(移動平均など)
窓の基準 固定グリッド 各行を終点とする窓

over()で銘柄ごとに窓を閉じる

# over を付け忘れると銘柄をまたいで平均が漏れる(バグの温床)
df_features = df_filled.with_columns([
    # ✅ 正しい:銘柄ごとに移動平均を計算
    pl.col("close").rolling_mean(window_size=20).over("symbol").alias("ma20"),
    pl.col("close").rolling_mean(window_size=60).over("symbol").alias("ma60"),

    # ✅ リターン(1期前比)
    (pl.col("close") / pl.col("close").shift(1).over("symbol") - 1).alias("return_1m"),

    # ✅ 実現ボラティリティ(20期)
    pl.col("return_1m").rolling_std(window_size=20).over("symbol").alias("realized_vol"),

    # ✅ 出来高比(直近20本の平均比)
    (pl.col("volume") / pl.col("volume").rolling_mean(20).over("symbol")).alias("vol_ratio"),
])

[記事232〜237(価格予測シリーズ)] で使った特徴量をPolarsで再実装する際も、このパターンが基本形になります。over("symbol")を付け忘れると複数銘柄を扱った瞬間にサイレントなバグが混入します。


STEP 5 ― join_asofで「時刻がずれたデータ」を突き合わせる

解きたい問題

約定データと板データのタイムスタンプは完全には一致しません。通常のjoinでは結合できないため、join_asofを使います。「各約定の直前の板スナップショットを付けたい」という要件が典型例です。

strategyの3種類

約定時刻: 09:01:03

板スナップショット: 09:00:55 / 09:01:05 / 09:01:30

backward(過去の最近値): 09:00:55  ← ✅ 実務のデフォルト
forward (未来の最近値): 09:01:05  ← ❌ 未来を参照(look-ahead bias)
nearest (前後で近い方) : 09:01:05  ← ❌ 未来を含みうる
strategy 拾う行 実務での推奨
"backward" キー以下で最も近い(過去) ✅ 基本はこれ
"forward" キー以上で最も近い(未来) ⚠️ 特殊用途のみ
"nearest" 前後で最も近い ⚠️ 未来を含む可能性あり

toleranceで「古すぎる値」の紐付けを防ぐ

# サンプル:約定データと板スナップショットを結合
df_trades = pl.DataFrame({
    "datetime": [
        datetime(2024, 1, 15, 9, 1, 3, tzinfo=timezone.utc),
        datetime(2024, 1, 15, 9, 2, 45, tzinfo=timezone.utc),
    ],
    "symbol":   ["BTC/USDT", "BTC/USDT"],
    "trade_price": [50100.0, 50200.0],
}).with_columns(pl.col("datetime").dt.cast_time_unit("ms"))

df_orderbook = pl.DataFrame({
    "datetime": [
        datetime(2024, 1, 15, 9, 0, 55, tzinfo=timezone.utc),
        datetime(2024, 1, 15, 9, 1, 5,  tzinfo=timezone.utc),
        datetime(2024, 1, 15, 9, 2, 30, tzinfo=timezone.utc),
    ],
    "symbol":   ["BTC/USDT"] * 3,
    "mid_price": [50090.0, 50110.0, 50190.0],
}).with_columns(pl.col("datetime").dt.cast_time_unit("ms"))

# ✅ 正しい実装:直前の板を取得、5秒以上古ければnull
result = df_trades.join_asof(
    df_orderbook.sort("datetime"),
    on="datetime",
    by="symbol",
    strategy="backward",   # 過去の最近値のみ
    tolerance="5s",        # 5秒以上古い板はnullにする
)
print(result)

【本記事の核心】join_asofの方向ミスが招くlook-ahead bias

バックテストが「良くなる」異常

look-ahead biasは、成績が悪化するのではなく向上するため気づきにくいのが最大の罠です。シャープレシオが2.0を超えたとき、まずリークを疑ってください。

事故の再現①:strategy="forward"を使ってしまう

# ❌ 誤った実装:未来の板を参照してしまう
result_bad = df_trades.join_asof(
    df_orderbook.sort("datetime"),
    on="datetime",
    by="symbol",
    strategy="forward",   # ← 未来の値を掴む
)

# ✅ 正しい実装
result_good = df_trades.join_asof(
    df_orderbook.sort("datetime"),
    on="datetime",
    by="symbol",
    strategy="backward",
    tolerance="5s",
)

# 比較:約定時刻09:01:03に対して
# bad  → mid_price=50110.0(09:01:05の板、未来の情報)
# good → mid_price=50090.0(09:00:55の板、正しい過去の情報)
print(result_bad.select(["datetime", "mid_price"]))
print(result_good.select(["datetime", "mid_price"]))

事故の再現②:group_by_dynamiclabelミスによるleak

# ❌ label="left"のままリサンプル結果を特徴量に使う(最も気づきにくいleak)
df_5m_bad = (
    df.sort("datetime")
    .group_by_dynamic("datetime", every="5m", closed="left", label="left", group_by=["symbol"])
    .agg(pl.col("close").last().alias("close_5m"))
)

# datetime=09:00 の close_5m は 09:00〜09:04:59 の終値(5分後の情報!)
# これを09:00時点の特徴量として使うと完全にリーク

# ✅ 正しい対処:label="right"にして窓の終端時刻をラベルにする
df_5m_good = (
    df.sort("datetime")
    .group_by_dynamic("datetime", every="5m", closed="left", label="right", group_by=["symbol"])
    .agg(pl.col("close").last().alias("close_5m"))
)
# datetime=09:05 の close_5m は 09:00〜09:04:59 の終値(正しい)

バックテスト結果の比較

同一データ・同一シグナルでstrategyだけを変えた実験では、以下のような乖離が生じます(実測値は戦略により異なりますが、傾向として):

指標 strategy="backward"(正) strategy="forward"(誤)
年率リターン +12% +47%
シャープレシオ 0.8 2.4
最大ドローダウン -18% -6%
勝率 52% 71%

「勝率71%、シャープ2.4」は実運用ではほぼありえない数字です。こうした非現実的な成績はバグを最初に疑うべきサインです。

リーク検出チェックリスト

# 簡易リークテスト:全特徴量を1期シフトして成績差を確認
# シフト前後で大幅に成績が変わる = リークの疑い大
df_no_leak = df_features.with_columns([
    pl.col("ma20").shift(1).over("symbol").alias("ma20_safe"),
    pl.col("close_5m").shift(1).over("symbol").alias("close_5m_safe"),
])

チェック項目:

  1. join_asofstrategy"backward"か?
  2. toleranceを設定しているか?
  3. by(銘柄・取引所キー)を指定しているか?
  4. group_by_dynamiclabel="right"になっているか(またはshift(1)で1期ずらしているか)?
  5. rollingshiftの向きは過去方向(正のshift値)か?

実践レシピ ― パイプライン全体を組み立てる

def build_feature_table(df_raw: pl.DataFrame) -> pl.DataFrame:
    """
    CCXT取得済みの1分足OHLCVから特徴量テーブルを構築する
    記事232〜237の予測モデルへそのまま渡せる形で出力
    """
    # 1. 時刻型の統一
    df = (
        df_raw
        .with_columns(
            pl.from_epoch("timestamp_ms", time_unit="ms")
              .dt.replace_time_zone("UTC")
              .alias("datetime")
        )
        .sort(["symbol", "datetime"])
    )

    # 2. 欠損バーの補完
    df = (
        df
        .upsample(time_column="datetime", every="1m", group_by=["symbol"])
        .sort(["symbol", "datetime"])
        .with_columns([
            pl.col("open").forward_fill().over("symbol"),
            pl.col("high").forward_fill().over("symbol"),
            pl.col("low").forward_fill().over("symbol"),
            pl.col("close").forward_fill().over("symbol"),
            pl.col("volume").fill_null(0),
        ])
    )

    # 3. 5分足・1時間足を生成してjoin_asof(label="right"でleak防止)
    df_5m = (
        df.sort("datetime")
        .group_by_dynamic("datetime", every="5m", closed="left", label="right", group_by=["symbol"])
        .agg(pl.col("close").last().alias("close_5m"), pl.col("volume").sum().alias("vol_5m"))
    )

    df = df.join_asof(
        df_5m.sort("datetime"),
        on="datetime",
        by="symbol",
        strategy="backward",  # 過去の5分足のみ参照
        tolerance="5m",
    )

    # 4. 特徴量を追加
    df = df.with_columns([
        pl.col("close").rolling_mean(20).over("symbol").alias("ma20"),
        pl.col("close").rolling_mean(60).over("symbol").alias("ma60"),
        (pl.col("close") / pl.col("close").shift(1).over("symbol") - 1).alias("return_1m"),
        pl.col("close").rolling_std(20).over("symbol").alias("vol_20m"),
    ])

    return df

# 実行
df_features = build_feature_table(df_1m)
print(df_features.head())

よくあるエラーと対処法(逆引き)

エラー / 症状 原因 対処
onキーがソートされていない旨のエラー sort未実施 join_asof前にsort("datetime")を必ず実行
join_asofで型エラー time_unitまたはタイムゾーンが不一致 df.schemaで両テーブルの型を確認し揃える
集計結果が1本ズレる closedまたはlabelの設定ミス closed="left", label="left/right"を意図で使い分ける
銘柄をまたいで値が混ざる group_by / over / byの付け忘れ 複数銘柄では必ずこれら3つを指定する
group_by引数が認識されない バージョン差異(旧by→新group_by pl.__version__を確認し公式ドキュメントの対応バージョンを参照

まとめ ― 5つの原則

本記事の要点を5つの原則として整理します。

  1. 内部はUTC固定replace_time_zoneでaware化し、表示時だけconvert_time_zoneでJSTに変換する
  2. group_by_dynamicclosedlabelまで指定する ― デフォルトのまま使うと境界が意図と1本ズレる
  3. 欠損バーはカラムごとに埋め方を変えるvolumeforward_fillすると出来高を捏造する
  4. join_asofbackward + tolerance + byを基本形にする ― この3点セットが安全側に倒す最短経路
  5. 成績が良すぎるバックテストはリークを最初に疑う ― シャープレシオ2.0超、勝率70%超は黄色信号

次に読むべき記事

  • Polarsの基礎から確認したい → 記事142『Polars入門』
  • pandasからの移行・書き換え対比 → 記事248
  • OHLCVデータの取得実装 → 記事227〜231(CCXTシリーズ)
  • 特徴量を予測モデルに繋げる → 記事232〜237(価格予測シリーズ)

関連記事

搭載メモリを超える数十GBデータをPolarsストリーミングエンジンで処理する実践ガイド【2026年版】

搭載メモリを超える数十GBデータをPolarsストリーミングエンジンで処理する実践ガイド【2026年版】

16GBラップトップで数十GB級データを処理するPolars新ストリーミングエンジンの実践ガイド。set_engine_affinity・sink_APIの使い分け・プッシュダウン最適化をPolars v1.39対応で徹底解説。

コメント

0/2000