주문 흐름 불균형과 수익률·실행 가능 마크아웃 비교
코드 Machine Learning for Trading
요약
이 노트북은 주문장 활동이 NVDA의 단기 수익률을 예측하는지 연구합니다. 주문별 시장 데이터에서 추가, 취소, 체결 정보를 이용해 거리에 따라 가중한 주문장 압력 지표를 만들고, 부호가 있는 활동을 평활화해 중간 가격 부근의 유동성 변화를 나타냅니다. 여러 미래 기간에 걸쳐 주문 흐름 불균형을 평가하고, 중간가, 지연 조정 가격, 실행 가능 가격을 포함한 여러 마크아웃 정의로 수익률 상관관계와 십분위 패턴을 비교합니다. 이 비교는 지연과 매수·매도 호가 스프레드를 고려하면 통계적 연관성이 거래 가능한 성과와 다를 수 있음을 보여줍니다.
근거는 단일 주식과 짧은 거래일 표본으로 제한됩니다. 미래 구간이 서로 겹치면 단순 표준오차가 지나치게 낙관적일 수 있으며, 추정치가 약하다는 사실만으로 작은 효과의 존재나 부재를 입증할 수 없다고 문서는 강조합니다. 비용 크기는 측정된 실행 비용이 아니라 일반적인 참고값이며, 무조건부 상관관계는 조건부 패턴을 가릴 수 있습니다. 이 문서는 주문 흐름을 독립적인 방향 예측보다 불리한 조건을 거르는 필터로 활용하는 편이 더 타당하다고 봅니다.
핵심 아이디어
- 주문장 압력은 메시지 방향, 동작, 규모, 중간가와의 거리를 결합해 유동성 변화를 요약합니다.
- 작은 상관관계는 불확실성과 유효 표본 크기를 함께 고려해 해석해야 합니다.
- 겹치는 미래 수익률 구간은 단순 표준오차의 독립성 가정을 위반합니다.
- 실행 가능한 결과에는 지연과 스프레드 비용이 포함되므로 중간가 기준 마크아웃은 잠재 수익률을 과대평가할 수 있습니다.
- 제한된 근거는 주문 흐름을 방향 예측보다 트레이딩 필터로 연구하는 편을 뒷받침합니다.
태그
전문
# 09_databento_mbo_analysis.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]
# # DataBento MBO: Order Flow Predictability
#
# **Chapter 3: Market Microstructure**
#
# **Docker image**: `ml4t`
#
# ## Purpose
#
# Quantify whether one-minute order-flow imbalance carries predictive
# information for next-minute returns on NVDA, and contrast midprice
# markouts with executable (bid/ask-at-decision-time) markouts to expose
# the gap §3.3 warns about between statistical signal and trading P&L.
#
# ## Learning Objectives
#
# After completing this notebook, you will be able to:
# - Build minute-resolution OFI from the LOB-reconstruction parquet that
# `08_databento_lob_reconstruction` writes.
# - Measure predictability at multiple horizons (1, 5, 10, 30 min) and
# recognize that the correlations are tiny and noisy on a sample of this
# length.
# - Distinguish midprice markouts from executable markouts and explain why
# spread and queue dynamics dominate P&L at sub-minute horizons.
#
# ## Book reference
#
# Section §3.3, *From Raw Messages to the Limit Order Book* — tradability
# caveat under "LOB Stylized Facts: Predictive Patterns".
#
# ## Prerequisites
#
# - DataBento XNAS-ITCH MBO parquets at
# `data/equities/market/microstructure/market_by_order/NVDA/` (10 trading
# days, November 2024).
#
# ## Data Source
#
# DataBento provides institutional-quality **Market-by-Order (MBO)** data:
# - Every order add, cancel, modify, and fill
# - Nanosecond timestamps (both exchange and receipt time)
# - Pre-filtered by symbol (no ITCH stock_locate mapping needed)
#
# ---
# %% [markdown]
# ## Setup
# %%
"""DataBento MBO: Order Flow Predictability — analyzing OFI signals in tick data."""
from pathlib import Path
import matplotlib.pyplot as plt
import numpy as np
import pandas as pd
import polars as pl
from tqdm.auto import tqdm
from data import load_mbo_data
from utils.paths import display_path, get_output_dir
from utils.reproducibility import set_global_seeds
from utils.style import COLORS, show_with_alt
# %% tags=["parameters"]
SEED = 42
# %%
set_global_seeds(SEED)
# Two extra shades on top of the repository palette: a lighter blue for a second series
# on the same axis, and a brown that stays distinguishable from both in greyscale.
COLORS = {**COLORS, "accent": "#4A90A4", "warm": "#8B4513"}
# %%
OUTPUT_DIR = get_output_dir(3, "databento")
SYMBOL = "NVDA"
# Get file paths from canonical loader
data_files = load_mbo_data(symbols=[SYMBOL], list_files=True)
SYMBOL_DIR = data_files[0].parent if data_files else None
# Production: use all files; testing: limit
MAX_FILES = None
MAX_ROWS = None
if SYMBOL_DIR and SYMBOL_DIR.exists():
data_files = sorted(SYMBOL_DIR.glob("*.parquet"))
if MAX_FILES:
data_files = data_files[:MAX_FILES]
print(f"Symbol: {SYMBOL}")
print(f"Data files: {len(data_files)} days")
if data_files:
dates = [f.stem.split("-")[-1].split(".")[0] for f in data_files]
print(f"Date range: {dates[0]} to {dates[-1]}")
else:
print(f"DataBento data not found at {display_path(SYMBOL_DIR)}")
data_files = []
# %% [markdown]
# ## 1. Understanding MBO Data
#
# Market-by-Order data captures every change to the order book:
#
# | Action | Meaning | Order Flow Signal |
# |--------|---------|-------------------|
# | **Add (A)** | New limit order placed | Passive liquidity supply |
# | **Cancel (C)** | Order withdrawn | Liquidity withdrawal |
# | **Fill (F)** | Order executed | Active demand/supply |
# | **Trade (T)** | Trade report | Price discovery |
#
# The **side** field tells us whether the order was on the bid (buy) or ask (sell).
# Combining action and side reveals the order flow dynamics.
# %%
ACTION_MAP = {"A": "Add", "C": "Cancel", "F": "Fill", "M": "Modify", "R": "Clear", "T": "Trade"}
SIDE_MAP = {"A": "Ask", "B": "Bid", "N": "None"}
def load_databento_day(file_path: Path, max_rows: int | None = None) -> pl.DataFrame:
"""Load one day of DataBento MBO data."""
df = pl.read_parquet(file_path)
if max_rows and len(df) > max_rows:
df = df.head(max_rows)
# Handle fixed-point prices if needed
if "price" in df.columns and df["price"].max() > 1_000_000:
df = df.with_columns((pl.col("price") / 1e9).alias("price"))
return df.with_columns(pl.col("timestamp").cast(pl.Datetime("ns")))
# %%
sample_df = None
multi_day = None
if data_files:
sample_df = load_databento_day(data_files[0], max_rows=MAX_ROWS)
date_str = data_files[0].stem.split("-")[-1].split(".")[0]
print(f"=== {SYMBOL} on {date_str} ===")
print(f"Total messages: {len(sample_df):,}")
# Message composition
print("\nMessage composition:")
action_counts = sample_df.group_by("action").len().sort("len", descending=True)
total = len(sample_df)
for row in action_counts.iter_rows():
pct = 100 * row[1] / total
print(f" {ACTION_MAP.get(row[0], row[0]):8s} {row[1]:>12,} ({pct:5.1f}%)")
# %% [markdown]
# Adds and cancels take almost all of that table between them and trades take a sliver.
# The ratio of the two is the point: a liquid name's book is re-quoted many times between
# consecutive trades, so most of what an MBO feed carries is the book rearranging itself
# rather than anything changing hands. That is what makes MBO data large and what makes
# it informative - the rearranging is the part a trade-only feed cannot see.
# %% [markdown]
# ## 2. Book Pressure: Measuring Order Flow Momentum
#
# Book pressure sums every message in a window into one signed number, weighting each by
# how much it should count:
#
# $$P = \sum_i s_i \, w_i \, q_i \, e^{-\lambda d_i}$$
#
# where $q_i$ is the size the message moved and:
#
# - $s_i$ is the side, positive on the bid and negative on the ask.
# - $w_i$ is the action. An add puts size on the book and counts positively; a cancel takes
# it off and counts negatively; a fill counts at half weight, because it both removes
# resting size and reveals someone who wanted to trade, and those pull in opposite
# directions.
# - $d_i$ is how far the message was from the midpoint, and $\lambda$ how quickly distance
# stops mattering. An order five cents away can be cancelled without anyone noticing; an
# order at the touch is what the next trade will hit.
#
# A positive sum means the bid side grew relative to the ask, weighted towards what is
# close enough to matter. Whether that anticipates the next price move is the question the
# rest of the notebook puts to the data rather than assumes.
# %%
def compute_book_pressure(
df: pl.DataFrame,
decay_lambda: float = 0.01,
ema_halflife: int = 100,
) -> pl.DataFrame:
"""Compute book pressure from MBO messages with EMA smoothing."""
book_actions = df.filter(pl.col("action").is_in(["A", "C", "F"]))
if len(book_actions) == 0:
return df.with_columns(pl.lit(0.0).alias("book_pressure"))
# Rolling midprice from recent trades (avoids lookahead)
trades = df.filter(pl.col("action") == "T").select(["timestamp", "price"])
if len(trades) > 0:
rolling_mid = trades.with_columns(
pl.col("price").rolling_mean(window_size=100, min_samples=1).alias("mid_price")
).select(["timestamp", "mid_price"])
df = df.sort("timestamp").join_asof(
rolling_mid.sort("timestamp"), on="timestamp", strategy="backward"
)
df = df.with_columns(pl.col("mid_price").fill_null(rolling_mid["mid_price"][0]))
else:
df = df.with_columns(pl.lit(book_actions["price"].median()).alias("mid_price"))
# Compute pressure components
result = (
df.with_columns(
[
pl.when(pl.col("side") == "B")
.then(1.0)
.when(pl.col("side") == "A")
.then(-1.0)
.otherwise(0.0)
.alias("side_sign"),
pl.when(pl.col("action") == "A")
.then(1.0)
.when(pl.col("action") == "C")
.then(-1.0)
.when(pl.col("action") == "F")
.then(0.5)
.otherwise(0.0)
.alias("action_weight"),
(pl.col("price") - pl.col("mid_price")).abs().alias("dist_from_mid"),
]
)
.with_columns((-decay_lambda * pl.col("dist_from_mid")).exp().alias("decay_weight"))
.with_columns(
(
pl.col("side_sign")
* pl.col("action_weight")
* pl.col("size")
* pl.col("decay_weight")
).alias("raw_pressure")
)
.with_columns(
pl.col("raw_pressure").ewm_mean(half_life=ema_halflife).alias("book_pressure")
)
)
return result.select(
["timestamp", "action", "side", "price", "size", "order_id", "book_pressure"]
)
# %%
pressure_df = None
if sample_df is not None and len(sample_df) > 0:
# One liquid mid-morning hour (10:00-11:00 America/New_York). The stored
# timestamp is naive UTC, so convert before filtering — a raw 10:00-11:00 UTC
# window would land on 05:00-06:00 ET, i.e. thin pre-market.
_et = pl.col("timestamp").dt.replace_time_zone("UTC").dt.convert_time_zone("America/New_York")
sample_hour = sample_df.filter((_et.dt.hour() >= 10) & (_et.dt.hour() < 11))
if len(sample_hour) > 1000:
pressure_df = compute_book_pressure(sample_hour)
fig, axes = plt.subplots(2, 1, figsize=(14, 8), sharex=True)
# Top: trade prices
trades = pressure_df.filter(pl.col("action") == "T").to_pandas()
if len(trades) > 0:
axes[0].plot(
trades["timestamp"],
trades["price"],
".",
markersize=1,
alpha=0.5,
color=COLORS["blue"],
)
axes[0].set_ylabel("Trade Price ($)")
axes[0].set_title(f"Book Pressure vs Price - {SYMBOL} ({date_str})")
# Bottom: book pressure (clipped for readability)
pressure_pd = pressure_df.to_pandas()
pressure_clipped = pressure_pd["book_pressure"].clip(
lower=pressure_pd["book_pressure"].quantile(0.01),
upper=pressure_pd["book_pressure"].quantile(0.99),
)
axes[1].plot(
pressure_pd["timestamp"], pressure_clipped, linewidth=0.5, color=COLORS["accent"]
)
axes[1].axhline(0, color="black", linestyle="--", linewidth=0.5)
axes[1].set_ylabel("Book Pressure (clipped at 1%/99%)")
axes[1].set_xlabel("Time (UTC)")
axes[1].fill_between(
pressure_pd["timestamp"],
pressure_clipped,
0,
where=pressure_clipped > 0,
alpha=0.3,
color="green",
label="Buy pressure",
)
axes[1].fill_between(
pressure_pd["timestamp"],
pressure_clipped,
0,
where=pressure_clipped < 0,
alpha=0.3,
color="red",
label="Sell pressure",
)
axes[1].legend(loc="upper right")
show_with_alt(
fig,
"Two stacked panels sharing a time axis over one hour of trading. The upper plots individual trade prices as small points. The lower plots book pressure as a line about a zero line, with the area above it shaded for buy pressure and the area below shaded in red for sell pressure.",
)
# Statistics
print("=== Book Pressure Statistics ===")
print(f"Mean: {pressure_pd['book_pressure'].mean():>10.2f}")
print(f"Std: {pressure_pd['book_pressure'].std():>10.2f}")
print(f"Skewness: {pressure_pd['book_pressure'].skew():>10.2f}")
# %% [markdown]
# **What the chart reveals**:
#
# Plotted at message resolution the pressure is dense and oscillates tightly
# around zero: green (net buying) and red (net selling) alternate almost
# continuously, with only occasional sustained excursions to one side. The
# intuition is that positive pressure precedes upticks and negative pressure
# precedes drops, but the visual signal-to-noise ratio is clearly low — no
# stable "buy zone / sell zone" structure jumps out of the raw trace.
#
# That is exactly why the rest of the notebook stops eyeballing tick pressure
# and instead aggregates to minute bars and measures the order-flow imbalance
# against forward returns directly: whether any of this apparent information is
# real has to be settled by the correlation and markout analysis below, not by
# reading the chart.
# %% [markdown]
# ## 3. Building Minute Bars with Microstructure Features
#
# Tick data is too noisy for most analysis. We aggregate to minute bars while
# preserving the microstructure information that matters for prediction.
# %%
def reconstruct_bbo_bars(df: pl.DataFrame, bar_freq: str = "1m") -> pl.DataFrame:
"""Reconstruct the best bid/offer from the MBO stream and snapshot it per bar.
A markout that reflects *tradable* P&L needs the quotes a marketable order
would actually hit — the best bid and best ask at decision time — not a
trade-price proxy. We replay the order stream to maintain top of book
(resting size by price on each side), then take the last quote in each bar.
Add (A) inserts resting size; Cancel (C) removes it; Fill (F) reduces it;
Modify (M) re-prices an order; Clear (R) wipes the side. We key on
``order_id`` so cancels, fills, and modifies adjust the level the order
actually rested on.
"""
orders: dict[int, tuple[str, float, int]] = {} # order_id -> (side, price, size)
bids: dict[float, int] = {} # price -> resting size
asks: dict[float, int] = {} # price -> resting size
# The inside quote is tracked incrementally; a full scan fires only when the
# prevailing best level empties.
best_bid: float | None = None
best_ask: float | None = None
def _book(side: str) -> dict[float, int]:
return bids if side == "B" else asks
def _add(side: str, price: float, size: int) -> None:
nonlocal best_bid, best_ask
book = _book(side)
book[price] = book.get(price, 0) + size
if side == "B":
if best_bid is None or price > best_bid:
best_bid = price
elif best_ask is None or price < best_ask:
best_ask = price
def _reduce(side: str, price: float, size: int) -> None:
nonlocal best_bid, best_ask
book = _book(side)
if price in book:
book[price] -= size
if book[price] <= 0:
del book[price]
# Only when the inside level vacates do we rescan that side.
if side == "B" and price == best_bid:
best_bid = max(bids) if bids else None
elif side == "A" and price == best_ask:
best_ask = min(asks) if asks else None
ts_out: list = []
bid_out: list[float | None] = []
ask_out: list[float | None] = []
for ts, action, side, price, size, oid in (
df.sort("timestamp")
.select(["timestamp", "action", "side", "price", "size", "order_id"])
.iter_rows()
):
if action == "A" and side in ("B", "A"):
orders[oid] = (side, price, size)
_add(side, price, size)
elif action == "C":
if oid in orders:
s, p, sz = orders.pop(oid)
_reduce(s, p, sz)
elif side in ("B", "A"):
_reduce(side, price, size)
elif action == "F":
if oid in orders:
s, p, sz = orders[oid]
_reduce(s, p, size)
if sz - size > 0:
orders[oid] = (s, p, sz - size)
else:
orders.pop(oid, None)
elif side in ("B", "A"):
_reduce(side, price, size)
elif action == "M":
if oid in orders:
s, p, sz = orders.pop(oid)
_reduce(s, p, sz)
if side in ("B", "A"):
orders[oid] = (side, price, size)
_add(side, price, size)
elif action == "R":
orders.clear()
bids.clear()
asks.clear()
best_bid = best_ask = None
# Trade (T) prints do not change the book; the matching Fill (F) does.
ts_out.append(ts)
bid_out.append(best_bid)
ask_out.append(best_ask)
bbo = pl.DataFrame(
{"timestamp": ts_out, "best_bid": bid_out, "best_ask": ask_out},
schema={"timestamp": pl.Datetime("ns"), "best_bid": pl.Float64, "best_ask": pl.Float64},
)
# Snapshot the quote prevailing at each bar close (last update within the bar).
return bbo.group_by_dynamic("timestamp", every=bar_freq, period=bar_freq).agg(
pl.col("best_bid").last(), pl.col("best_ask").last()
)
# %%
def process_day_to_bars(
file_path: Path, bar_freq: str = "1m", max_rows: int | None = None
) -> pl.DataFrame:
"""Process one day to minute bars with OHLCV and order flow metrics."""
df = load_databento_day(file_path, max_rows=max_rows)
# Reconstruct top of book from the full stream (orders resting before the
# open carry into the session), then restrict bars to regular trading hours.
bbo_bars = reconstruct_bbo_bars(df, bar_freq=bar_freq)
# Regular trading hours are defined on the exchange's clock, so convert before
# filtering: a window fixed in UTC is an hour wrong for half the year.
_et = pl.col("timestamp").dt.replace_time_zone("UTC").dt.convert_time_zone("America/New_York")
df = df.filter(
((_et.dt.hour() > 9) | ((_et.dt.hour() == 9) & (_et.dt.minute() >= 30)))
& (_et.dt.hour() < 16)
)
if len(df) == 0:
return pl.DataFrame()
trades = df.filter(pl.col("action") == "T")
if len(trades) == 0:
return pl.DataFrame()
# Aggregate to bars
bars = trades.group_by_dynamic("timestamp", every=bar_freq, period=bar_freq).agg(
[
pl.col("price").first().alias("open"),
pl.col("price").max().alias("high"),
pl.col("price").min().alias("low"),
pl.col("price").last().alias("close"),
pl.col("size").sum().alias("volume"),
pl.len().alias("trade_count"),
# Cast to a signed dtype: ``size`` is unsigned, and buy_volume -
# sell_volume (below, for OFI) underflows to a huge positive value on
# net-selling bars if these stay unsigned. Polars wraps silently.
(pl.when(pl.col("side") == "B").then(pl.col("size")).otherwise(0))
.sum()
.cast(pl.Int64)
.alias("buy_volume"),
(pl.when(pl.col("side") == "A").then(pl.col("size")).otherwise(0))
.sum()
.cast(pl.Int64)
.alias("sell_volume"),
(pl.col("price") * pl.col("size")).sum().alias("notional"),
]
)
# Derived features
bars = bars.with_columns(
[
(pl.col("close") / pl.col("open") - 1).alias("return"),
pl.when(pl.col("buy_volume") + pl.col("sell_volume") > 0)
.then(
(pl.col("buy_volume") - pl.col("sell_volume"))
/ (pl.col("buy_volume") + pl.col("sell_volume"))
)
.otherwise(0.0)
.alias("ofi"),
(pl.col("notional") / pl.col("volume")).alias("vwap"),
]
)
bars = bars.join(bbo_bars, on="timestamp", how="left")
bars = bars.with_columns(((pl.col("best_bid") + pl.col("best_ask")) / 2).alias("mid_quote"))
return bars
# %%
if data_files:
all_bars = []
n_days = MAX_FILES or len(data_files)
for file in tqdm(data_files[:n_days], desc="Processing days"):
bars = process_day_to_bars(file, max_rows=MAX_ROWS)
if len(bars) > 0:
all_bars.append(bars)
if all_bars:
multi_day = pl.concat(all_bars).sort("timestamp")
print("\n=== Multi-Day Dataset ===")
print(f"Minute bars: {len(multi_day):,}")
print(f"Trading days: {multi_day['timestamp'].dt.date().n_unique()}")
print(
f"Date range: {multi_day['timestamp'].min().date()} to {multi_day['timestamp'].max().date()}"
)
# %% [markdown]
# ## 4. The Key Question: Does OFI Predict Returns?
#
# Order Flow Imbalance (OFI) measures the balance between buyer and seller
# aggression in each minute:
#
# **OFI = (Buy Volume - Sell Volume) / Total Volume**
#
# Values range from -1 (all selling) to +1 (all buying). If markets are
# efficient, OFI should predict short-term returns—aggressive buyers push
# prices up, aggressive sellers push prices down.
#
# But how strong is this relationship, and how long does it last?
# %% [markdown]
# A signal computed at the close of a bar cannot be acted on within that bar. `LATENCY_BARS`
# is how many bars pass between the two, and it is bound once here because both the markout
# computation and the plotting cell below read it; two copies of the same assumption would
# eventually disagree.
# %%
LATENCY_BARS = 1
def compute_markouts(
bars: pl.DataFrame, horizons: list | None = None, latency_bars: int = LATENCY_BARS
) -> pl.DataFrame:
"""Compute forward returns at multiple horizons within each trading session."""
if horizons is None:
horizons = [1, 5, 10, 30]
result = bars.clone().with_columns(pl.col("timestamp").dt.date().alias("session_date"))
for h in horizons:
# Standard markout: price change from now to h minutes ahead
result = result.with_columns(
(pl.col("mid_quote").shift(-h).over("session_date") / pl.col("mid_quote") - 1).alias(
f"markout_{h}"
)
)
# Latency-adjusted: assumes 1-bar execution delay
p1 = pl.col("mid_quote").shift(-latency_bars).over("session_date")
p2 = pl.col("mid_quote").shift(-h).over("session_date")
result = result.with_columns((p2 / p1 - 1).alias(f"markout_{h}_adj"))
# Executable markout: a long pays the ask at decision time and exits at
# the bid h bars later, the round trip a marketable order realizes. Its gap to
# the midpoint markout is the half-spread paid on entry plus the one paid on exit.
ask_now = pl.col("best_ask")
bid_future = pl.col("best_bid").shift(-h).over("session_date")
result = result.with_columns((bid_future / ask_now - 1).alias(f"markout_{h}_exec"))
return result
# %%
if multi_day is not None and len(multi_day) > 0:
multi_day = compute_markouts(multi_day, horizons=[1, 5, 10, 30], latency_bars=LATENCY_BARS)
# Add lagged OFI within sessions
multi_day = multi_day.with_columns(
[
pl.col("ofi").shift(1).over("session_date").alias("ofi_lag1"),
pl.col("ofi")
.rolling_mean(window_size=5)
.shift(1)
.over("session_date")
.alias("ofi_ma5"),
]
)
# Predictor only: dropping the markouts here would truncate every horizon to the
# longest one's sample, leaving the per-horizon dropna() below nothing to select.
pdf = multi_day.drop_nulls(["ofi", "ofi_lag1"]).to_pandas()
print("Correlation of OFI(t-1) with the return over the following window\n")
print(
"The standard error is roughly 1/sqrt(n) for independent draws and larger here, "
"because the forward windows overlap. Each horizon has its own n: a longer "
"window loses more bars at the end of each session.\n"
)
print(f"{'Horizon':<10}{'Bars':>10}{'Correlation':>14}{'In naive SEs':>15}")
print("-" * 49)
for h in [1, 5, 10, 30]:
col = f"markout_{h}"
if col not in pdf.columns:
continue
pair = pdf[["ofi_lag1", col]].dropna()
if len(pair) < 2:
continue
corr = pair["ofi_lag1"].corr(pair[col])
naive_se = 1 / np.sqrt(len(pair))
print(f"{h:>3} min {len(pair):>10,}{corr:>14.4f}{corr / naive_se:>15.1f}")
# %% [markdown]
# Read that table against its last column rather than its middle one. A Pearson
# correlation is a number whatever the data does, and the question is whether it is
# distinguishable from what noise would produce on this many observations.
#
# Three things make the honest standard error larger than the naive one printed above.
# The correlations are computed on the minute-bar panel, not on the tick stream, so the
# sample is thousands of bars rather than millions of messages. The forward-return
# windows overlap - the return over the next ten minutes shares nine minutes with the
# one starting a minute later - so successive rows carry much of the same information.
# And returns are serially correlated within a session. Each of those shrinks the
# effective number of independent observations below the bar count.
#
# So a sign that flips between horizons is not a finding about horizons. And a
# correlation inside a standard error of zero does not establish that there is nothing
# there - an imprecise estimate is what a real but small effect also looks like on this
# much data. What it establishes is that this sample does not measure a linear
# relationship; whether one exists, and in which direction, is left open.
#
# That is why the question worth asking of a microstructure signal is economic rather
# than statistical: how many basis points the conditional return is worth, against what
# a round trip costs. The two questions come apart in both directions. A correlation
# indistinguishable from zero on this sample can still carry a conditional return large
# enough to pay for itself, and a correlation measured precisely can be worth a fraction
# of a basis point. The next panels put the forward return in basis points on the axis,
# which is the quantity a cost can be compared against.
# %%
if multi_day is not None and len(multi_day) > 0:
fig, axes = plt.subplots(1, 2, figsize=(14, 5))
# Both panels measure the same thing, so they use the same rows: the bars a
# 5-minute markout exists for. pdf keeps every horizon's own sample, so the
# pairing has to be made here rather than assumed.
paired = pdf[["ofi_lag1", "markout_5"]].dropna()
# Sample for scatter plot (too many points otherwise)
sample = paired.sample(min(10000, len(paired)), random_state=42)
# Left panel: scatter OFI vs 5-min return + trend
axes[0].scatter(
sample["ofi_lag1"],
sample["markout_5"] * 10000,
alpha=0.1,
s=2,
color=COLORS["blue"],
)
axes[0].set_xlabel("OFI (t-1)")
axes[0].set_ylabel("5-Minute Return (bps)")
axes[0].set_title("OFI vs 5-Minute Forward Returns")
axes[0].axhline(0, color="black", linestyle="--", linewidth=0.5)
axes[0].axvline(0, color="black", linestyle="--", linewidth=0.5)
z = np.polyfit(sample["ofi_lag1"], sample["markout_5"] * 10000, 1)
p = np.poly1d(z)
x_line = np.linspace(-1, 1, 100)
axes[0].plot(
x_line, p(x_line), color=COLORS["warm"], linewidth=2, label=f"Trend (slope={z[0]:.2f})"
)
axes[0].legend()
# Right panel: binned decile analysis
paired = paired.assign(
ofi_bin=pd.qcut(paired["ofi_lag1"], q=10, labels=False, duplicates="drop")
)
binned = paired.groupby("ofi_bin")["markout_5"].mean() * 10000
colors = [COLORS["warm"] if v < 0 else COLORS["accent"] for v in binned.values]
axes[1].bar(range(len(binned)), binned.values, color=colors, edgecolor="black", linewidth=0.5)
axes[1].set_xlabel("OFI Decile (1=Most Selling, 10=Most Buying)")
axes[1].set_ylabel("Mean 5-Min Return (bps)")
axes[1].set_title("Mean 5-Min Return by OFI Decile (Within Noise)")
axes[1].axhline(0, color="black", linestyle="--", linewidth=0.5)
spread = binned.iloc[-1] - binned.iloc[0]
axes[1].annotate(
f"Spread: {spread:.1f} bps",
xy=(8.5, binned.iloc[-1]),
fontsize=10,
ha="center",
)
fig.suptitle(f"{SYMBOL}: order-flow imbalance against the following return", fontsize=12)
show_with_alt(
fig,
f"Two panels for {SYMBOL}. The left scatters each bar's order-flow imbalance against the return over the following window, one point per bar. The right sorts the bars into ten equal groups by imbalance and draws the mean return of each group as a bar, with an annotation giving the difference between the top and bottom groups in basis points.",
)
# %% [markdown]
# What to look for in the decile means, in order of how much it would take to convince
# you:
#
# - **A ramp.** If flow carried direction, the means would rise monotonically from the
# heaviest-selling decile to the heaviest-buying one. Alternating signs across the
# interior bins are what noise looks like when it is sorted into ten buckets.
# - **Tails larger than the middle.** A signal that lives in extreme flow shows up as
# the two end bins standing away from the rest. Interior bins as large as the extreme
# ones say the sort found nothing.
# - **A top-minus-bottom spread larger than it costs to capture.** Round-trip execution
# in liquid US equities runs on the order of a basis point or two, so a spread of that
# size is not an edge whatever its sign - it is the fee.
#
# That third point is the one that decides it, and it is why the next panel makes the
# cost frictions explicit and shows
# why even a marginally non-zero correlation would not convert to tradable P&L.
# %% [markdown]
# ## 5. From Price Response to Tradable P&L: Three Markout Types
#
# The correlation above is a *price-response* signal measured on the midpoint.
# What a desk actually keeps is smaller, because two frictions sit between the
# signal and the fill. We separate them with three markouts at each horizon:
#
# - **Midpoint (p₂ − p₀)**: the raw price response — return of the quote
# midpoint from signal time to *h* bars ahead. This is the number the OFI
# correlation is built on, and the most generous reading of the edge.
# - **Latency-adjusted (p₂ − p₁)**: the same midpoint return after a one-bar
# execution delay. The gap to the midpoint markout is the alpha that decays
# while the order is in flight.
# - **Executable (bid/ask)**: a long pays the **ask** at decision time and
# exits at the **bid** *h* bars later. The gap to the midpoint markout is the
# half-spread paid on entry plus the half-spread paid on exit, the friction §3.3
# warns turns a statistical signal into negative trading P&L at sub-minute horizons.
#
# Reading the three together shows where the edge goes: latency erodes it, and
# the spread can erase it outright.
# %% [markdown]
# The three markout definitions differ only in what they subtract. The midpoint markout is
# the move in the quote midpoint, which is the signal with no frictions at all. The
# latency-adjusted one starts a bar later, which is the earliest a signal computed at a
# bar close could have been acted on. The executable one also pays the spread, which is
# what an order that crosses actually gives up. Comparing the three distributions locates
# where a return would be lost rather than asserting that it is.
# %%
if multi_day is not None and len(multi_day) > 0:
summary_rows = []
fig, axes = plt.subplots(2, 2, figsize=(14, 10))
for i, h in enumerate([1, 5, 10, 30]):
ax = axes[i // 2, i % 2]
mid_col, adj_col, exec_col = f"markout_{h}", f"markout_{h}_adj", f"markout_{h}_exec"
pdf_clean = multi_day.drop_nulls([mid_col, exec_col]).to_pandas()
def _clip_bps(series):
return series.clip(lower=series.quantile(0.01), upper=series.quantile(0.99)) * 10000
ax.hist(
_clip_bps(pdf_clean[mid_col]),
bins=50,
alpha=0.5,
label="Midpoint (p₂-p₀)",
color=COLORS["blue"],
density=True,
)
ax.hist(
_clip_bps(pdf_clean[exec_col]),
bins=50,
alpha=0.5,
label="Executable (bid/ask)",
color=COLORS["warm"],
density=True,
)
ax.axvline(0, color="black", linestyle="--", linewidth=0.5)
ax.set_xlabel("Markout (bps)")
ax.set_ylabel("Density")
ax.set_title(f"{h}-Minute Horizon")
ax.legend()
mid_mean = pdf_clean[mid_col].mean() * 10000
exec_mean = pdf_clean[exec_col].mean() * 10000
ax.annotate(
f"mean: {mid_mean:+.1f} vs {exec_mean:+.1f} bps\nspread cost: {mid_mean - exec_mean:.1f} bps",
xy=(0.95, 0.95),
xycoords="axes fraction",
ha="right",
va="top",
fontsize=9,
)
# At h == LATENCY_BARS the latency-adjusted markout compares a price with itself,
# so it is zero by construction rather than by measurement; report it past that.
adj_mean = (
multi_day.drop_nulls([adj_col])[adj_col].mean() * 10000 if h > LATENCY_BARS else None
)
summary_rows.append((h, mid_mean, adj_mean, exec_mean))
fig.suptitle(f"{SYMBOL}: midpoint and executable markouts at four horizons", fontsize=12)
show_with_alt(
fig,
f"Four panels in a two-by-two grid, one per forward horizon, for {SYMBOL}. Each overlays two distributions of returns in basis points: the markout measured on the quote midpoint, and the markout an order that crossed the spread would have realised. The gap between the two distributions is the half-spread paid on entry plus the one paid on exit.",
)
print("=== Mean markout by type (bps) ===")
print(f"{'Horizon':<10}{'Midpoint':>12}{'Latency-adj':>14}{'Executable':>14}")
print("-" * 50)
for h, mid_mean, adj_mean, exec_mean in summary_rows:
adj_str = f"{adj_mean:>14.1f}" if adj_mean is not None else f"{'≡ 0':>14}"
print(f"{h:>3} min {mid_mean:>12.1f}{adj_str}{exec_mean:>14.1f}")
# %% [markdown]
# **Key insights from the markout analysis**:
#
# 1. **Midpoint markout is the generous view**: it credits the signal with the
# full price response and ignores both execution delay and the spread.
#
# 2. **Latency erodes the edge**: a single bar of delay (p₂ − p₁) removes the
# fraction of the move that lands in the first minute — largest at short
# horizons, negligible by 30 minutes once total movement dominates.
#
# 3. **The spread can erase it**: paying the ask and exiting at the bid shifts
# the executable markout left of the midpoint markout by roughly one spread.
# At sub-minute horizons that gap routinely exceeds the signal itself, which
# is why a positive midpoint correlation is necessary but not sufficient for
# tradable P&L.
# %% [markdown]
# ## 6. Save Results
#
# Export the processed minute bars for use in downstream analysis.
# %%
if multi_day is not None and len(multi_day) > 0:
output_file = OUTPUT_DIR / f"{SYMBOL}_minute_bars.parquet"
multi_day.write_parquet(output_file)
print(f"Saved: {output_file}")
print(f"Shape: {multi_day.shape}")
# %% [markdown]
# ## Key Takeaways
#
# 1. **Read a correlation against its standard error, not against zero.** The table
# prints the sample size, the estimate and their ratio, and it is the ratio that says
# whether the estimate is a measurement. An estimate inside one standard error leaves
# the question open rather than answering it in the negative: a real effect too small
# for this sample looks exactly the same. What it does rule out is reading the sign,
# or treating a flip between horizons as a horizon effect.
#
# 2. **Overlapping forward windows are not independent observations.** A ten-minute return
# measured every minute shares nine minutes with its neighbour, so the effective sample
# is well below the bar count and the naive standard error understates the real one.
# Any test that treats these rows as independent overstates its own confidence.
#
# 3. **Sort into deciles to look for a shape, not for a number.** A directional signal
# shows as a monotone ramp across the bins or as tails standing away from the middle.
# Alternating signs across the interior is what a sort of noise produces.
#
# 4. **Judge a microstructure signal against what trading it costs.** A decile spread on
# the order of a round trip is not a small edge; it is the fee. That comparison, not the
# statistical one, is what decides whether a signal is worth anything, which is why the
# three markout columns exist: the midpoint markout is the signal before frictions, the
# latency-adjusted one subtracts the delay, and the executable one subtracts the spread.
# On a signal with no edge to begin with they demonstrate the method rather than measure
# a loss.
#
# 5. **Order flow is more defensible as a filter than as a forecast.** Declining to trade
# against heavy one-sided flow asks much less of the data than predicting direction from
# it, and this notebook's evidence supports the first and not the second.
#
# ### Known limitations
#
# - One symbol over a short slice of days. Nothing here establishes what holds for other
# names, other periods, or other venues.
# - The correlations are unconditional. This notebook does not split a return into
# permanent and transient components, nor condition on spread, volatility or time of day,
# any of which could carry structure the unconditional estimate averages away.
# - The cost figures used for comparison are typical magnitudes for liquid US equities,
# not measurements of what these trades would have cost.
#
# ### DataBento vs ITCH: When to Use Each
#
# | Aspect | ITCH (Free) | DataBento (~$10/symbol/mo) |
# |--------|-------------|----------------------------|
# | Format | Binary, all NASDAQ stocks | Parquet, per-symbol |
# | Preprocessing | Significant | Minimal |
# | Multi-exchange | No | Yes |
# | Best for | Learning, backtesting | Production, research |
#
# ---
#
# **Next**: The notebooks in Chapter 8 build on these microstructure insights
# to engineer alpha factors for machine learning models.
```출처의 라이선스에 따라 출처를 표시하고 전문을 공개합니다. 라이선스: MIT
이 요약은 원문을 바탕으로 Stratmill의 리서치 에이전트가 작성했으며, 원문을 복사한 것이 아닙니다.