Cập nhật dữ liệu thị trường từng phần, phát hiện khoảng trống và giám sát tình trạng
Tóm tắt
Notebook giải thích vòng đời dữ liệu thị trường OHLCV hằng ngày: tải lịch sử ban đầu, lấy phần dữ liệu bổ sung có một khoảng chồng lấn nhỏ, phát hiện phiên bị thiếu và tóm tắt độ mới cùng vấn đề xác thực của từng mã. Notebook trình bày lý do cửa sổ cập nhật nên dừng ở nến hoàn chỉnh gần nhất từ nhà cung cấp, vì dòng dữ liệu của ngày hiện tại có thể chứa giá chưa hoàn chỉnh, và lý do khoảng chồng lấn giúp ghi nhận các lần nhà cung cấp chỉnh sửa. Notebook cũng nêu các cách tiếp cận cập nhật từng phần, chỉ nối thêm, làm mới toàn bộ và điền dữ liệu thiếu cho các nhu cầu cập nhật khác nhau.
Các ví dụ dùng cổ phiếu US theo ngày và báo cáo bảng ban đầu có 502 dòng mỗi mã trước khi thêm các phiên mới có dữ liệu. Kiểm tra khoảng trống có xét cuối tuần loại thứ Bảy và Chủ nhật, nhưng không có lịch sở giao dịch nên ngày nghỉ thị trường vẫn có thể bị xem là khoảng trống. Notebook lưu ý rằng công cụ điền khoảng trống không xét lịch có thể tạo các dòng OHLCV giả. Bảng điều khiển theo dõi số dòng đã lưu, tuổi dữ liệu, khoảng trống và vấn đề xác thực; các kiểm tra này hỗ trợ giám sát thường xuyên nhưng không thay thế việc rà soát tính đầy đủ theo lịch trước khi kiểm thử lịch sử.
Ý chính
- Dùng thao tác lấy dữ liệu từng phần có chồng lấn để ghi nhận phiên mới và các lần nhà cung cấp chỉnh sửa.
- Dừng cập nhật ở nến hoàn chỉnh gần nhất từ nhà cung cấp để tránh nạp dữ liệu ngày chưa hoàn tất.
- Tránh điền giá trị gần nhất về phía trước vào các khoảng trống lịch trong dữ liệu OHLCV vì cuối tuần và ngày lễ không phải phiên giao dịch.
- Loại trừ cuối tuần giúp kiểm tra khoảng trống, nhưng cần lịch sở giao dịch để phân biệt ngày lễ với dữ liệu bị thiếu.
- Theo dõi độ mới, số lượng khoảng trống và vấn đề xác thực trước khi dùng dữ liệu đã lưu trong backtest.
Thẻ
Toàn văn
# 19_incremental_updates.py
```py
# ---
# jupyter:
# jupytext:
# cell_metadata_filter: tags,-all
# text_representation:
# extension: .py
# format_name: percent
# format_version: '1.3'
# jupytext_version: 1.19.3
# kernelspec:
# display_name: Python 3 (ipykernel)
# language: python
# name: python3
# ---
# %% [markdown]
# # 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.
# %% [markdown]
# ## Setup
# %%
"""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)}")
# %% [markdown]
# ### 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.
# %% tags=["parameters"]
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"
# %% [markdown]
# ---
#
# ## 1. Full Refresh vs Incremental Update
#
# Let's demonstrate the difference with a concrete example.
# %% [markdown]
# ### Step 1: Initial Load (Full History)
# %%
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)")
# %% [markdown]
# ### 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`.
# %%
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)")
# %% [markdown]
# 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.
# %% [markdown]
# ---
#
# ## 2. Update Strategies
#
# ml4t-data supports four strategies for different scenarios:
# %%
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",
],
}
)
# %% [markdown]
# The `IncrementalUpdater` class provides low-level control over these strategies.
# For most use cases, `DataManager.update()` (which uses `INCREMENTAL`) is sufficient.
# %% [markdown]
# ---
#
# ## 3. Gap Detection
#
# Before running a backtest, verify that your data is complete.
# Gaps can occur from failed downloads, provider outages, or holiday handling.
# %%
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)")
# %% [markdown]
# `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.
# %% [markdown]
# ---
#
# ## 4. Data Health Dashboard
#
# In production, monitor data health across your entire universe.
# %%
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)
# %%
report = data_health_report(storage, symbols)
report
# %%
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.",
)
# %% [markdown]
# ---
#
# ## 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.
# %% [markdown]
# ---
#
# ## 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/)
```Hiển thị toàn văn kèm ghi nguồn theo giấy phép của tài liệu gốc. Giấy phép: MIT
Bản tóm tắt này do tác nhân nghiên cứu của Stratmill biên soạn từ tài liệu gốc; đây không phải bản sao của tài liệu.