الانتقال إلى المحتوى
جميع مستندات المكتبة

اختلال تدفق الأوامر وإشارات العوائد وتكاليف التداول في NVDA

دفتر ملاحظات Machine Learning for Trading

الملخص

تبني هذه الدراسة مقياسًا على مستوى الدقيقة لضغط تدفق الأوامر باستخدام بيانات دفتر الأوامر على مستوى كل أمر NVDA. وتحدد اتجاه إضافات الأوامر وإلغائها وتنفيذها بحسب جانب الشراء أو البيع، وتزن الأحداث وفق حجمها وبعدها عن نقطة المنتصف، ثم تنعّم سلسلة الضغط الناتجة. وتفحص بعد ذلك ما إذا كان المقياس يتنبأ بالعوائد عبر آفاق مستقبلية عدة، مستخدمةً أسعار التقييم عند نقطة المنتصف والأسعار القابلة للتنفيذ وقت القرار للتمييز بين الارتباط الإحصائي والأداء القابل للتداول.

الأدلة المبلغ عنها ضعيفة: توصف الارتباطات بأنها ضئيلة ومليئة بالضوضاء في عينة قصيرة، وتحذر الوثيقة من أن نوافذ العوائد المتداخلة تجعل الأخطاء المعيارية العادية مفرطة في التفاؤل. ويمكن للتصنيف إلى عشيرات أن يساعد على تمييز نمط اتجاهي متسق من إشارات مشوشة، بينما قد تطغى تكاليف الفروق والتأخير على آثار صغيرة حول نقطة المنتصف. ويرى المؤلف أن تدفق الأوامر قد يكون أنفع كمرشح منه كتنبؤ اتجاهي. وتقتصر النتائج على سهم واحد وفترة قصيرة وارتباطات غير مشروطة وتكاليف تداول معتادة لا مقاسة مباشرة.

الأفكار الرئيسية

  • يجمع ضغط تدفق الأوامر بين أحداث الأوامر محددة الاتجاه وحجم الحدث وبعده عن نقطة المنتصف والتنعيم.
  • قد تبالغ أسعار التقييم عند نقطة المنتصف في تقدير فائدة الإشارة مقارنة بأسعار العرض والطلب القابلة للتنفيذ.
  • تقلل نوافذ العوائد المستقبلية المتداخلة حجم العينة الفعال، وقد تجعل تقديرات الدلالة الساذجة غير موثوقة.
  • النمط الرتيب عبر عشيرات الإشارة يوحي بوجود نمط اتجاهي أكثر من الإشارات المتناوبة.
  • تدعو العينة القصيرة التي تقتصر على سهم واحد إلى تفسير حذر، ولا تثبت الأداء في سياقات أخرى.

الوسوم

النص الكامل
# DataBento MBO: Order Flow Predictability


# 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)

---

## Setup

```python
"""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
```

```python
SEED = 42
```

```python
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"}
```

```python
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 = []
```

## 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.

```python
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")))
```

```python
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}%)")
```

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.

## 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.

```python
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"]
    )
```

```python
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}")
```

**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.

## 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.

```python
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()
    )
```

```python
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
```

```python
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()}"
        )
```

## 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?

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.

```python
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
```

```python
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}")
```

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.

```python
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.",
    )
```

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.

## 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.

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.

```python
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}")
```

**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.

## 6. Save Results

Export the processed minute bars for use in downstream analysis.

```python
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}")
```

## 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.
![notebook output](figures/p1_1.png)
![notebook output](figures/p1_2.png)
![notebook output](figures/p1_3.png)

يُعرض النص كاملًا مع نسبه إلى مصدره وفقًا لترخيصه. الترخيص: MIT

أعدّ وكيل الأبحاث في Stratmill هذا الملخص استنادًا إلى المصدر الأصلي؛ وهو ليس نسخة منه.