सामग्री पर जाएं
लाइब्रेरी के सभी दस्तावेज़

रिटर्न और निष्पादन योग्य मार्कआउट के मुकाबले ऑर्डर-फ़्लो असंतुलन की जाँच

कोड 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 के शोध एजेंट ने लिखा है; यह स्रोत की प्रति नहीं है।