市場データの逐次更新と完全性の監視
ノートブック Machine Learning for Trading
サマリー
日次の市場データを最新に保つ運用ワークフローを紹介します。初期履歴を読み込み、期間が少し重複する差分を取得し、欠損日を確認して、鮮度と検証上の問題を監視します。重複期間を設けることで、データ提供元による直近の足の修正も取り込めます。OHLCVデータの更新期間は現在の日付ではなく、最後に完了した足までとします。提供元から部分的な取引セッションや価格が未確定のセッションが返される場合があるためです。このガイドでは、取引所カレンダーなしで欠損を補完すると、週末や祝日の人工的な足が作られる可能性も指摘します。
差分更新、追記のみ、全件更新、バックフィルの各方法を比較し、週末を除外する欠損検出と、保存行数、鮮度、欠損、検証上の問題を示すダッシュボードを紹介します。確認方法にも限界があります。週末を除外しても取引所の祝日は考慮できず、取引終了後でも提供元の最新データが未完成の場合があります。バックテストでデータの完全性を信頼する前に、取引所カレンダーを考慮した確認が必要です。例が示すのはデータ運用と健全性の監視であり、トレード成績ではありません。
主なアイデア
- 差分更新では直近データを取得して保存済みの履歴と統合し、重複期間によって提供元の修正を取り込めます。
- 提供元が進行中の部分的な取引セッションを公開する場合があるため、更新は最後に完了した足までとします。
- 無条件に欠損を補完すると、週末や祝日に架空のOHLCV行が作られる可能性があります。
- 週末を除外する欠損検出でも、祝日を正しく特定するには取引所カレンダーが必要です。
- バックテスト前に、データ量、鮮度、欠損、検証上の問題を健全性レポートで追跡できます。
タグ
全文
# Incremental Updates: Keeping Market Data Fresh
# Incremental Updates: Keeping Market Data Fresh
**Docker image**: `ml4t`
**Purpose**: Walk through the update lifecycle — initial load, daily delta,
gap detection, and a small health dashboard — so the same flow that runs
in `data/download_all.py --update` is visible end-to-end.
**Learning objectives**:
1. Compare full-refresh and incremental update costs.
2. Drive a daily delta with `DataManager.update(..., fill_gaps=False)` and
understand why the `fill_gaps=True` default is unsafe for OHLCV data.
3. Verify completeness with `GapDetector(exclude_weekends=True)`.
4. Render a 3-panel health dashboard (volume, freshness, issues).
**Book reference**: §2.4 (storing data) — operational counterpart to the
storage benchmarks in notebooks 20 and 21.
**Prerequisites**: ml4t-data installed; live network access for Yahoo Finance.
**Why incremental updates**: a full refresh re-fetches every row a symbol has ever had,
and an incremental run fetches the sessions since the last one. On a universe of a few
hundred symbols with a decade of history, the first is millions of rows and the second is
hundreds, every day. Section 1 measures the ratio on the demo universe.
## Setup
```python
"""Incremental Updates — Keeping market data fresh with daily delta fetches."""
import logging
import shutil
from datetime import datetime
import matplotlib.pyplot as plt
import polars as pl
import structlog
# Quiet ml4t-data's structured logger so the notebook focuses on demo output.
structlog.configure(
wrapper_class=structlog.make_filtering_bound_logger(logging.WARNING),
)
from ml4t.data import DataManager
from ml4t.data.storage import HiveStorage
from ml4t.data.storage.backend import StorageConfig
from ml4t.data.update_manager import GapDetector
from ml4t.data.validation import OHLCVValidator
from utils.downloading import update_through_last_complete_bar
from utils.paths import display_path, get_output_dir
from utils.style import COLORS, show_with_alt
# Storage for this notebook's demos. Wipe any prior-run artifacts so the
# initial-load + update sequence is reproducible.
DEMO_DIR = get_output_dir(2, "incremental_updates")
if DEMO_DIR.exists():
shutil.rmtree(DEMO_DIR)
DEMO_DIR.mkdir(parents=True, exist_ok=True)
print(f"Demo storage: {display_path(DEMO_DIR)}")
```
### Declared parameters
The freshness thresholds are the two that carry a judgment. A daily equity feed that has
not moved in more than a weekend is behind, and one that has not moved in more than a week
is broken; those are the two lines the dashboard draws and colours against.
```python
DEMO_SYMBOLS = ["AAPL", "MSFT", "GOOGL"]
FULL_HISTORY_START = "2023-01-01"
FULL_HISTORY_END = "2024-12-31"
DEMO_PROVIDER = "yahoo"
UPDATE_LOOKBACK_DAYS = 7
FRESH_DAYS = 3 # up to this many days behind is fresh
STALE_DAYS = 7 # beyond this the feed is not merely behind
MAX_RETURN_THRESHOLD = 0.5 # the validator's extreme-return cutoff
STORAGE_COMPRESSION = "zstd"
```
---
## 1. Full Refresh vs Incremental Update
Let's demonstrate the difference with a concrete example.
### Step 1: Initial Load (Full History)
```python
config = StorageConfig(base_path=DEMO_DIR / "updates_demo", compression=STORAGE_COMPRESSION)
storage = HiveStorage(config=config)
dm = DataManager(storage=storage)
symbols = DEMO_SYMBOLS
print("=== Initial Load (Full History) ===")
for symbol in symbols:
dm.load(symbol, FULL_HISTORY_START, FULL_HISTORY_END, provider=DEMO_PROVIDER)
meta = dm.get_metadata(symbol)
print(f" {symbol}: stored ({meta['row_count']} rows)")
```
### Step 2: Incremental Update (Only New Data)
An update reads the last timestamp in storage, refetches a small overlap
(`UPDATE_LOOKBACK_DAYS`) so a bar the vendor revised replaces the stored one, and
merges everything since.
**Where it stops matters more than where it starts.** Yahoo returns the
current exchange date as a row with accumulating volume and no
open/high/low/close. A bar with no prices is not a bar, and the provider says
so:
```
DataValidationError: yahoo: Column 'open' contains 1 null values
```
`DataManager.update()` fetches to `datetime.now()`, so it asks for that row on
every trading day and raises. `update_through_last_complete_bar` performs the
same delta and **finds** the end of its window rather than computing one. It
starts a day before the current exchange date — in exchange time, not UTC,
since at 22:00 in New York the UTC date is already tomorrow — and steps back a
day at a time while the provider refuses the window.
The retreat is what a date cannot do. The placeholder row usually resolves a
few hours after the close, and sometimes it does not: Yahoo's 2026-09-03 daily
bar for AAPL was still `NaN, NaN, NaN, NaN, 37197362` eight and a half hours
later, while the same window ending 2026-09-02 was clean on every column. Both
look identical from the calendar, so the vendor has to be asked.
When a window is refused, ml4t-data logs it at error level. If you see
`Failed to fetch AAPL: yahoo: Column 'open' contains 1 null values` below and
then a row count on the next line, that is the retreat working: the first
window asked for a bar the vendor has not published, and the second one landed.
The other reason not to call `update()` blind is its `fill_gaps=True` default:
a calendar-unaware gap detector that treats every weekend and US holiday as a
missing trading day and forward-fills it, inflating each symbol's row count by
hundreds of phantom bars. For OHLCV data the right behavior is to leave
non-trading days absent and rely on a calendar-aware completeness check
downstream — which is what section 3 does with `GapDetector`.
```python
print("=== Incremental Update ===")
for symbol in symbols:
rows = update_through_last_complete_bar(
dm, storage, symbol, provider=DEMO_PROVIDER, lookback_days=UPDATE_LOOKBACK_DAYS
)
print(f" {symbol}: updated ({rows} rows)")
```
Each symbol started with 502 rows (2023-01 through 2024-12). The update
fetches the 7-day overlap plus every session through the last complete bar and
merges; the resulting row count equals the original 502 plus exactly the new
trading sessions — no synthetic weekend/holiday rows, and no half-formed bar
for the session that is still running.
---
## 2. Update Strategies
ml4t-data supports four strategies for different scenarios:
```python
pl.DataFrame(
{
"strategy": ["INCREMENTAL", "APPEND_ONLY", "FULL_REFRESH", "BACKFILL"],
"behavior": [
"Fetch data after last stored timestamp",
"Add new rows; never modify existing",
"Re-download and replace all data",
"Fetch missing periods inside existing range",
],
"use_case": [
"Daily updates (default, fastest)",
"Audit-safe archives",
"Recovery after corruption",
"Patch holes in historical data",
],
}
)
```
The `IncrementalUpdater` class provides low-level control over these strategies.
For most use cases, `DataManager.update()` (which uses `INCREMENTAL`) is sufficient.
---
## 3. Gap Detection
Before running a backtest, verify that your data is complete.
Gaps can occur from failed downloads, provider outages, or holiday handling.
```python
gap_detector = GapDetector(exclude_weekends=True)
print("=== Gap Detection ===")
for symbol in symbols:
key = f"equities/daily/{symbol}"
df = storage.read(key).collect()
gaps = gap_detector.detect_gaps(df, frequency="daily")
if gaps:
print(f"{symbol}: {len(gaps)} gap(s) found")
for gap in gaps[:5]:
print(f" {gap['start'].date()} -> {gap['end'].date()} ({gap['size_days']} days)")
if len(gaps) > 5:
print(f" ... and {len(gaps) - 5} more")
else:
print(f"{symbol}: complete (no gaps)")
```
`exclude_weekends=True` filters Saturdays and Sundays. US market holidays
(MLK Day, Good Friday, Thanksgiving, etc.) still register as one-day gaps
because the detector has no exchange calendar — fine for routine update
health checks, but pair with a calendar-aware completeness check before
trusting the result for backtest panels.
---
## 4. Data Health Dashboard
In production, monitor data health across your entire universe.
```python
def data_health_report(storage: HiveStorage, symbols: list[str]) -> pl.DataFrame:
"""Per-symbol freshness, gap count, and validation issue count."""
validator = OHLCVValidator(max_return_threshold=MAX_RETURN_THRESHOLD)
detector = GapDetector(exclude_weekends=True)
now = datetime.now()
rows = []
for symbol in symbols:
df = storage.read(f"equities/daily/{symbol}").collect()
last_date = df["timestamp"].max().replace(tzinfo=None)
days_stale = (now - last_date).days
gaps = detector.detect_gaps(df, frequency="daily")
result = validator.validate(df)
rows.append(
{
"symbol": symbol,
"status": "stale" if days_stale > STALE_DAYS else "fresh",
"rows": len(df),
"last_date": last_date.date(),
"days_stale": days_stale,
"gaps": len(gaps),
"issues": result.error_count if not result.passed else 0,
}
)
return pl.DataFrame(rows)
```
```python
report = data_health_report(storage, symbols)
report
```
```python
fig, axes = plt.subplots(1, 3, figsize=(14, 4), constrained_layout=True)
axes[0].barh(report["symbol"].to_list(), report["rows"].to_list(), color=COLORS["blue"])
axes[0].set_xlabel("Rows")
axes[0].set_title("Rows stored")
stale_days = report["days_stale"].to_list()
freshness_color = [
COLORS["positive"]
if d <= FRESH_DAYS
else COLORS["amber"]
if d <= STALE_DAYS
else COLORS["negative"]
for d in stale_days
]
# A bar of zero width draws nothing, so a fully fresh universe would render a blank panel.
# The markers below are what a reader sees instead of nothing.
axes[1].barh(report["symbol"].to_list(), stale_days, color=freshness_color)
zero_mask = [i for i, d in enumerate(stale_days) if d == 0]
if zero_mask:
axes[1].scatter(
[0] * len(zero_mask),
[report["symbol"].to_list()[i] for i in zero_mask],
color=[freshness_color[i] for i in zero_mask],
s=80,
marker="o",
zorder=3,
label="Fresh (0 days)",
)
axes[1].axvline(
FRESH_DAYS, color=COLORS["positive"], linestyle="--", alpha=0.6, label=f"{FRESH_DAYS} days"
)
axes[1].axvline(
STALE_DAYS, color=COLORS["amber"], linestyle="--", alpha=0.6, label=f"{STALE_DAYS} days"
)
axes[1].set_xlim(left=-0.5)
axes[1].set_xlabel("Days Since Update")
axes[1].set_title("Days since last update")
axes[1].legend(fontsize=8, loc="lower right")
gaps = report["gaps"].to_list()
issues = report["issues"].to_list()
axes[2].barh(report["symbol"].to_list(), gaps, color=COLORS["negative"], label="Gaps")
axes[2].barh(
report["symbol"].to_list(), issues, left=gaps, color=COLORS["amber"], label="Validation"
)
axes[2].set_xlabel("Count")
axes[2].set_title("Gaps and validation flags")
axes[2].legend(fontsize=8, loc="lower right")
fig.suptitle("Per-symbol storage health")
show_with_alt(
fig,
"Three horizontal-bar panels sharing a symbol axis. The left shows rows stored per "
"symbol, with bars of near-equal length. The middle shows days since the last update "
"against two dashed threshold lines, with a coloured marker at zero for each symbol "
"that is fully up to date. The right stacks gap counts and validation flags per "
"symbol.",
)
```
---
## 5. The Book's Update Workflow
The book's `data/download_all.py` script implements exactly this pattern
for all asset classes. Run it with `--update` to extend datasets to the present:
```bash
# Initial download (run once)
python data/download_all.py
# Update to present (run daily/weekly)
python data/download_all.py --update
```
Under the hood, `download_all.py` uses:
- `ETFDataManager.from_config("data/etfs/config.yaml")` → `manager.update()`
- `CryptoDataManager.from_config("data/crypto/config.yaml")` → `manager.update()`
- `MacroDataManager.from_config("data/macro/config.yaml")` → `manager.download_treasury_yields()`
- `FuturesDataManager.from_config("data/futures/config.yaml")` → `manager.download_all()`
Each manager reads its YAML config for symbols, date ranges, and provider settings,
then updates only what's new.
### Config-Driven Downloads
The ETF config at `data/etfs/config.yaml` defines:
```yaml
etfs:
provider: yahoo
start: '2006-01-01'
end: '2025-12-31'
frequency: daily
tickers:
us_equity_broad:
symbols: [SPY, QQQ, IWM, ...]
us_sectors:
symbols: [XLB, XLC, XLE, ...]
# ... 9 categories, 100 ETFs total
```
When you run `--update` in 2026+, it extends data beyond the configured end date
to the present — no config changes needed.
---
## Summary
### Key Patterns
| Pattern | Command | When |
|---------|---------|------|
| Initial load | `dm.load(symbol, start, end)` | First time |
| Daily update | `dm.update(symbol, lookback_days=7)` | Every trading day |
| Gap check | `updater.detect_gaps(df, "daily")` | Before backtesting |
| Full refresh | `UpdateStrategy.FULL_REFRESH` | After data corruption |
| Batch update | `download_all.py --update` | Cron job |
### Production Checklist
1. **Initial download**: `python data/download_all.py` (run once, ~10 min)
2. **Schedule updates**: Add `download_all.py --update` to cron (daily at 6 PM)
3. **Monitor health**: Check freshness, gaps, and validation before backtests
4. **Validate data**: Use `OHLCVValidator` on every load (see `13_data_quality_framework`)
### Key Takeaways
- **`lookback_days` is the operating lever**, not the symbol set. Set it once
to cover provider revision windows (Yahoo retro-adjusts ~5 trading days; CRSP
~30) and the strategy generalizes across instruments.
- **Gap detection runs against stored data, not stream data.** A daily cron
that finishes with `detect_gaps()` catches missed trading days before any
backtest reads stale partitions.
- **Full refresh is the recovery path.** Use `UpdateStrategy.FULL_REFRESH`
when validation flags a regression; never patch a corrupted parquet in place.
- **Configurable end dates eliminate config drift.** Symbol-list YAMLs use
today-as-default; `--update` extends data beyond the file's `end:` field
without needing edits each calendar year.
- **Cron + idempotent CLI is the production interface.** The same
`download_all.py --update` works in a notebook, a CI run, and a 6 PM cron.
Next: `20_storage_benchmark_file` compares file-format throughput for the
parquet write path implicit in every update; `21_storage_benchmark_database`
extends the same comparison to database engines.
### Cross-References
- **Data quality**: `13_data_quality_framework` — validation and anomaly detection.
- **DataManager basics**: `18_data_management` — fetch, batch, Universe, storage.
- **Storage benchmarks**: `20_storage_benchmark_file` (file formats) and
`21_storage_benchmark_database` (database engines).
- **Download scripts**: `data/download_all.py` — the book's orchestrator.
- **ml4t-data docs**: [ml4trading.io/docs/data/user-guide/incremental-updates/](https://ml4trading.io/docs/data/user-guide/incremental-updates/)
出典を明記したうえで、ライセンスに従って全文を掲載しています。 ライセンス: MIT
この要約は原文をもとにStratmillのリサーチエージェントが作成したもので、出典の複製ではありません。