Medição de spreads intradiários, profundidade e desequilíbrio do fluxo de ordens
Resumo
Este notebook analisa livros de ofertas limitadas reconstruídos do NASDAQ para descrever spreads intradiários e a profundidade no topo do livro, e então examina se o desequilíbrio do fluxo de ordens está associado aos retornos subsequentes por intervalo. Ele expressa os spreads em pontos-base para comparar ações com diferentes níveis de preço e plota como os spreads e a liquidez exibida mudam ao longo de uma sessão. Também calcula o desequilíbrio a partir de adições e remoções de ordens, agrega o fluxo em intervalos de tempo e compara correlações com retornos posteriores para ações individuais e para o corte transversal. Uma transformação logarítmica com sinal é usada junto ao fluxo bruto para reduzir a influência de valores extremos.
O gráfico de spread ilustra um padrão intradiário comumente descrito, mas abrange um símbolo e um dia, portanto não é uma evidência ampla desse padrão. As correlações de desequilíbrio descrevem a mesma sessão usada para calculá-las; nenhum período de holdout testa o desempenho preditivo. A medida manual de fluxo ignora mensagens de substituição, ao contrário da medida baseada na reconstrução, e a análise usa apenas a profundidade nos melhores preços de compra e venda. Símbolos com poucas cotações são excluídos quando fornecem poucos intervalos, condicionando ainda mais o corte transversal.
Ideias principais
- Expresse os spreads cotados em pontos-base para tornar mais significativas as comparações entre diferentes níveis de preço.
- Examine tanto o spread quanto a profundidade exibida de compra e venda para caracterizar a liquidez intradiária nos melhores preços.
- Agrupe o fluxo de ordens em intervalos e compare valores brutos com valores em log com sinal, pois caudas pesadas podem dominar as correlações.
- Examine a distribuição das correlações transversais em vez de tirar conclusões a partir de uma única ação.
- Trate as correlações da mesma sessão como descritivas, pois a análise não usa um período de holdout.
Tags
Texto completo
# LOB Analysis: Stylized Facts and Predictive Patterns
# LOB Analysis: Stylized Facts and Predictive Patterns
**Chapter 3: Market Microstructure**
**Docker image**: `ml4t`
## Book reference
Section §3.3, *From raw messages to the limit order book*, and Figure 3.3.
## Purpose
The reconstructed books from `02_itch_lob_reconstruction` are the input. This notebook
reads them for three things: how the spread moves through a session, how the shares
resting at the touch move with it, and whether the imbalance between buying and selling
pressure says anything about the return over the bucket that follows. The last of those
is measured
across a cross-section of stocks rather than one, because a single symbol cannot
distinguish a signal from a coincidence.
## Pipeline Position
```
nasdaq_itch_parser (01)
↓
order_book_reconstruction (02) → LOB snapshots → THIS NOTEBOOK (03)
↓ ↓
order_lifecycle_analysis (04) Ch8: Feature Engineering
```
**Upstream**: LOB snapshots from `02_itch_lob_reconstruction` in
`03_market_microstructure/output/nasdaq_itch/order_book/{SYMBOL}/`
See the **README.md** for the full notebook inventory and learning paths.
## Learning Objectives
After completing this notebook, you will be able to:
- Read the bid-ask spread off a reconstructed book and express it in basis points, so
that a wide spread on a $300 stock and a wide spread on a $3 stock are comparable.
- Plot how the spread and the depth at the touch move through a trading session, and
say what shape they take.
- Build order-flow imbalance - shares added to one side minus shares taken off it -
from raw ITCH messages, and say which messages that construction ignores.
- Correlate that imbalance against the following bucket's return for one stock, then
for a cross-section of fifty, and read the spread of correlations rather than any
single value.
## Prerequisites
- LOB snapshots from `02_itch_lob_reconstruction`, one directory per symbol. The
cross-section in Section 5 needs that notebook run once per symbol.
- Parsed ITCH messages from `01_itch_parser`, for the imbalance built in Section 4.
---
## Setup
```python
"""LOB analysis: spread, depth and order-flow imbalance over reconstructed ITCH books."""
from pathlib import Path
import matplotlib.dates as mdates
import matplotlib.patches as mpatches
import matplotlib.pyplot as plt
import numpy as np
import pandas as pd
import polars as pl
from data.equities.loader import load_nasdaq_itch
from utils.paths import display_path, get_output_dir
from utils.style import COLORS, show_with_alt
```
### Declared parameters
`TRADING_DATE` is the session the books were reconstructed for; it labels the figures
and nothing else selects on it, because each symbol directory holds one day.
`BUCKET_FREQ` is the interval the imbalance is summed over and the return measured
across. One minute is short enough that order flow and the next price move are plausibly
related and long enough that a bucket holds many messages.
`MIN_BUCKETS` is how many buckets a symbol must contribute before its
correlation is reported. A correlation over a handful of points is noise with a value
attached, so thinly traded symbols drop out rather than widening the cross-section with
estimates nobody should read.
```python
TRADING_DATE = "2020-01-30"
BUCKET_FREQ = "1m"
MIN_BUCKETS = 50
```
```python
NASDAQ_ITCH_OUTPUT = get_output_dir(3, "nasdaq_itch")
OUTPUT_DIR = get_output_dir(3, "lob_analysis")
OUTPUT_DIR.mkdir(parents=True, exist_ok=True)
ORDER_BOOK_DIR = NASDAQ_ITCH_OUTPUT / "order_book"
MESSAGES_DIR = load_nasdaq_itch(get_base_path=True)
```
### Which symbols have a reconstructed book
`02_itch_lob_reconstruction` writes one directory per symbol it was run for, so the
directory listing is the inventory: whatever is there is what this notebook can read.
```python
def discover_available_symbols(base_dir: Path) -> list[str]:
"""Find all symbols with reconstructed LOB data."""
symbols = []
if not base_dir.exists():
return symbols
for subdir in sorted(base_dir.iterdir()):
if subdir.is_dir() and (subdir / "lob_snapshots.parquet").exists():
symbols.append(subdir.name)
return symbols
AVAILABLE_SYMBOLS = discover_available_symbols(ORDER_BOOK_DIR)
print(f"Data directory: {display_path(ORDER_BOOK_DIR)}")
print(
f"Available symbols: {AVAILABLE_SYMBOLS if AVAILABLE_SYMBOLS else 'None (run notebook 02 first)'}"
)
```
## 1. Load Reconstructed LOB Data
LOB snapshots were generated by `02_itch_lob_reconstruction.py` and saved to
`03_market_microstructure/output/nasdaq_itch/order_book/{SYMBOL}/`. Each symbol has its own directory
containing snapshots and order-level data.
If data is missing, run notebook 02 first.
```python
def load_lob_snapshots(symbol: str, base_dir: Path) -> pl.DataFrame | None:
"""Load LOB snapshots for a symbol from the order book directory."""
symbol_dir = base_dir / symbol
snapshot_file = symbol_dir / "lob_snapshots.parquet"
if snapshot_file.exists():
df = pl.read_parquet(snapshot_file)
print(f" {symbol}: {len(df):,} snapshots loaded")
return df
print(f" {symbol}: Not found at {snapshot_file}")
return None
```
```python
# Load available LOB data (dynamically discovered)
print("Loading LOB snapshots from order_book_reconstruction output:")
lob_data = {}
for symbol in AVAILABLE_SYMBOLS:
df = load_lob_snapshots(symbol, ORDER_BOOK_DIR)
if df is not None:
lob_data[symbol] = df
assert lob_data, (
f"No LOB snapshots found at {ORDER_BOOK_DIR}.\n"
"This notebook needs LOB data from the ITCH pipeline:\n"
" 1. Download: uv run python data/equities/market/microstructure/nasdaq_itch_download.py\n"
" 2. Parse: Run 01_itch_parser.py\n"
" 3. Reconstruct LOB: Run 02_itch_lob_reconstruction.py\n"
f" Expected output: {ORDER_BOOK_DIR}/{{SYMBOL}}/lob_snapshots.parquet"
)
print(f"\nLoaded {len(lob_data)} symbols for analysis.")
```
## 2. Spread Dynamics: Intraday Patterns
The spread is the distance between the highest bid and the lowest ask: what it costs to
buy and immediately sell. Quoted in dollars it is not comparable across stocks, so it is
converted to basis points - hundredths of a percent of the mid price - which is what
makes a penny spread on a $3 stock and a penny spread on a $300 stock tell different
stories.
Spreads are widely reported to trace a U through the session: wide at the open while
overnight information is still being priced, narrower through the middle of the day,
and wide again into the close as positions are squared. The figure below is one symbol
on one day, which is an illustration of the shape rather than evidence for it.
Chapter 8 builds time-of-day features on the same idea.
```python
def compute_spread_metrics(lob_df: pl.DataFrame) -> pl.DataFrame:
"""Compute spread-related metrics from LOB snapshots."""
result = lob_df.clone()
# Ensure spread exists
if "spread" not in result.columns:
if "best_ask" in result.columns and "best_bid" in result.columns:
result = result.with_columns((pl.col("best_ask") - pl.col("best_bid")).alias("spread"))
# Mid price if not present
if "mid_price" not in result.columns and "best_ask" in result.columns:
result = result.with_columns(
((pl.col("best_ask") + pl.col("best_bid")) / 2).alias("mid_price")
)
# Spread in basis points (only meaningful for positive spreads)
if "spread" in result.columns and "mid_price" in result.columns:
result = result.with_columns(
(pl.col("spread") / pl.col("mid_price") * 10000).alias("spread_bps")
)
return result
```
```python
# Select first available stock for spread dynamics analysis
SPREAD_SYMBOL = list(lob_data.keys())[0] if lob_data else None
spread_1m = None
mid_1m = None
if lob_data and SPREAD_SYMBOL:
spread_df = compute_spread_metrics(lob_data[SPREAD_SYMBOL])
# Filter to valid spreads only
spread_df = spread_df.filter(pl.col("spread") > 0)
if len(spread_df) > 10:
# Convert timestamp to pandas for time-series operations
spread_pd = spread_df.to_pandas()
if "timestamp" in spread_pd.columns:
# Polars hands back microseconds; pandas resampling and plotting want
# nanoseconds, so make the unit explicit rather than inferred.
spread_pd["timestamp"] = spread_pd["timestamp"].astype("datetime64[ns]")
spread_pd = spread_pd.set_index("timestamp")
# Resample to 1-minute for cleaner visualization
if "spread_bps" in spread_pd.columns:
spread_1m = spread_pd["spread_bps"].resample("1min").mean()
mid_1m = spread_pd["mid_price"].resample("1min").last()
else:
print(f"Not enough valid spread data for {SPREAD_SYMBOL}")
```
```python
ALT_SPREAD_DYNAMICS = "Two stacked line charts sharing a clock-time axis across one trading session. The upper panel plots the mid price as a single line. The lower panel plots the bid-ask spread in basis points as a thin noisy line with a heavier five-minute rolling mean over it, both on the same axis."
if spread_1m is not None and mid_1m is not None:
fig, axes = plt.subplots(2, 1, figsize=(14, 8), sharex=True)
mid_clean = mid_1m.dropna()
axes[0].plot(mid_clean.index, mid_clean.values, color=COLORS["blue"], linewidth=1)
axes[0].set_title(f"{SPREAD_SYMBOL} mid price, one-minute samples")
axes[0].set_ylabel("Price ($)")
axes[0].grid(True, alpha=0.3)
spread_clean = spread_1m.dropna()
axes[1].plot(
spread_clean.index,
spread_clean.values,
color=COLORS["slate"],
linewidth=1,
alpha=0.7,
label="Spread",
)
if len(spread_clean) > 5:
rolling_mean = spread_clean.rolling(window=5, min_periods=1).mean()
axes[1].plot(
rolling_mean.index,
rolling_mean.values,
color=COLORS["amber"],
linewidth=2,
label="5-minute rolling mean",
)
axes[1].set_title(f"{SPREAD_SYMBOL} bid-ask spread in basis points")
axes[1].set_ylabel("Spread (bps)")
axes[1].set_xlabel("Time (US/Eastern)")
axes[1].legend(loc="upper right")
axes[1].grid(True, alpha=0.3)
axes[1].xaxis.set_major_formatter(mdates.DateFormatter("%H:%M"))
fig.autofmt_xdate()
show_with_alt(fig, ALT_SPREAD_DYNAMICS)
print(f"Spread for {SPREAD_SYMBOL} on {TRADING_DATE}:")
print(f" Mean: {spread_clean.mean():.2f} bps")
print(f" Median: {spread_clean.median():.2f} bps")
```
## 3. Top-of-Book Evolution
The previous figure shows what the touch cost; this one shows where it sat and how much
was resting there. Each minute contributes one bid point and one ask point, positioned
at that side's price, and shaded by the log of the shares behind it - darker means more
depth. The mid price runs through the middle as a line.
The log scaling on the shading matters: depth at the touch spans several orders of
magnitude within a session, so a linear shade would leave every minute but the busiest
few indistinguishable.
These snapshots carry the top of the book only. `08_databento_lob_reconstruction`
builds the same picture with several price levels a side.
```python
def _resample_side_depth(pdf, price_col: str, size_col: str, freq: str) -> list[dict]:
"""Resample one side's depth at a single level into scatter records."""
if price_col not in pdf.columns or size_col not in pdf.columns:
return []
resampled = (
pdf[[price_col, size_col]].resample(freq).agg({price_col: "last", size_col: "sum"}).dropna()
)
return [
{"timestamp": ts, "price": row[price_col], "depth": row[size_col]}
for ts, row in resampled.iterrows()
if row[size_col] > 0
]
```
### Build Depth Scatter Data
Resample order book levels to a fixed frequency for scatter visualization.
```python
def build_depth_scatter_data(
lob_df: pl.DataFrame, n_levels: int = 1, resample_freq: str = "1min"
) -> tuple:
"""Build scatter plot data for top-of-book / inside-quote visualization.
Resamples to the specified frequency and extracts bid/ask prices and depths
at each requested level. ``n_levels>1`` only returns rows where per-level
price columns (``bid_price_{i}``, ``ask_price_{i}``) are present in
``lob_df``; on ITCH inside-quote frames levels >0 degrade silently to
empty.
Args:
lob_df: LOB snapshot DataFrame with bid/ask price and size columns
n_levels: Number of price levels to include (1 = top-of-book only)
resample_freq: Resample frequency (e.g., "1min", "5min")
Returns:
Tuple of (mid_prices_df, bid_scatter_df, ask_scatter_df)
Each scatter df has columns: [timestamp, price, depth, log_depth]
"""
pdf = lob_df.to_pandas()
pdf["timestamp"] = pd.to_datetime(pdf["timestamp"])
pdf = pdf.set_index("timestamp")
# The ITCH snapshots name the touch best_bid/best_ask; the resampler asks for the
# level-0 names, so alias them.
if "bid_price_0" not in pdf.columns and "best_bid" in pdf.columns:
pdf["bid_price_0"] = pdf["best_bid"]
if "ask_price_0" not in pdf.columns and "best_ask" in pdf.columns:
pdf["ask_price_0"] = pdf["best_ask"]
# Resample mid prices
if "mid_price" not in pdf.columns:
pdf["mid_price"] = (pdf["best_bid"] + pdf["best_ask"]) / 2
mid_prices = pdf["mid_price"].resample(resample_freq).last().dropna()
# Build bid/ask scatter data across all levels
bid_records, ask_records = [], []
for level in range(n_levels):
bid_records.extend(
_resample_side_depth(pdf, f"bid_price_{level}", f"bid_size_{level}", resample_freq)
)
ask_records.extend(
_resample_side_depth(pdf, f"ask_price_{level}", f"ask_size_{level}", resample_freq)
)
bid_df = pd.DataFrame(bid_records)
ask_df = pd.DataFrame(ask_records)
# Add log depth for color scaling
if not bid_df.empty:
bid_df["log_depth"] = np.log1p(bid_df["depth"])
if not ask_df.empty:
ask_df["log_depth"] = np.log1p(ask_df["depth"])
return mid_prices, bid_df, ask_df
```
```python
DEPTH_SYMBOL = list(lob_data.keys())[0] if lob_data else None
mid_prices = None
bid_scatter = None
ask_scatter = None
if lob_data and DEPTH_SYMBOL:
depth_lob = lob_data[DEPTH_SYMBOL]
# These snapshots carry one level a side, so n_levels=1 is what there is to draw.
mid_prices, bid_scatter, ask_scatter = build_depth_scatter_data(
depth_lob, n_levels=1, resample_freq="1min"
)
```
```python
if (
mid_prices is not None
and len(mid_prices) > 10
and bid_scatter is not None
and not bid_scatter.empty
):
fig, ax = plt.subplots(figsize=(14, 7))
# Cast both series to nanoseconds before plotting: the resampled index keeps polars'
# microseconds, a frame built from dicts defaults to nanoseconds.
mid_ts_numeric = mid_prices.index.astype("datetime64[ns]").astype(np.int64)
bid_ts_numeric = bid_scatter["timestamp"].astype("datetime64[ns]").astype(np.int64)
ask_ts_numeric = ask_scatter["timestamp"].astype("datetime64[ns]").astype(np.int64)
# Plot bid depth (blue) - below mid price
ax.scatter(
bid_ts_numeric,
bid_scatter["price"],
c=bid_scatter["log_depth"],
cmap="Blues",
alpha=0.4,
s=15,
label="Bid",
)
# Plot ask depth (red) - above mid price
ax.scatter(
ask_ts_numeric,
ask_scatter["price"],
c=ask_scatter["log_depth"],
cmap="Reds",
alpha=0.4,
s=15,
label="Ask",
)
# Plot mid price as black line
ax.plot(mid_ts_numeric, mid_prices.values, color="black", linewidth=1.5, label="Mid Price")
# Format x-axis with clock-time labels (avoid raw nanosecond ticks)
n_ticks = 8
tick_indices = np.linspace(0, len(mid_prices) - 1, n_ticks, dtype=int)
tick_positions = mid_ts_numeric[tick_indices]
tick_labels = [mid_prices.index[i].strftime("%H:%M") for i in tick_indices]
ax.set_xticks(tick_positions)
ax.set_xticklabels(tick_labels)
ax.set_ylabel("Price ($)", fontsize=12)
ax.set_xlabel("Time (ET)", fontsize=12)
# Add legend with colored patches
blue_patch = mpatches.Patch(color="royalblue", alpha=0.6, label="Bid Depth")
red_patch = mpatches.Patch(color="indianred", alpha=0.6, label="Ask Depth")
black_line = plt.Line2D([0], [0], color="black", linewidth=1.5, label="Mid Price")
ax.legend(handles=[blue_patch, red_patch, black_line], loc="upper left")
ax.set_title(f"{DEPTH_SYMBOL} best bid and ask with resting depth, {TRADING_DATE}")
ax.grid(True, alpha=0.3)
show_with_alt(
fig,
"A scatter chart over one trading session. Blue points trace the best bid and red points the best ask, one of each per minute, positioned at that side's price and shaded by the logarithm of the shares resting there, so darker points carry more depth. A black line runs between the two bands showing the mid price. The horizontal axis is labelled in clock time and the vertical axis in dollars.",
)
else:
print(
"Not enough snapshots for the depth figure; reconstruct a fuller session in "
"02_itch_lob_reconstruction."
)
```
```python
# Summary statistics
if (
mid_prices is not None
and len(mid_prices) > 10
and bid_scatter is not None
and not bid_scatter.empty
):
print(f"Top of book for {DEPTH_SYMBOL}:")
print(f" Session covered: {mid_prices.index.min()} to {mid_prices.index.max()}")
print(f" Mid price range: ${mid_prices.min():.2f} to ${mid_prices.max():.2f}")
print(f" Bid points drawn: {len(bid_scatter):,}")
print(f" Ask points drawn: {len(ask_scatter):,}")
if not bid_scatter.empty:
print(f" Mean bid depth: {bid_scatter['depth'].mean():,.0f} shares")
if not ask_scatter.empty:
print(f" Mean ask depth: {ask_scatter['depth'].mean():,.0f} shares")
```
## 4. Order-flow imbalance, built from raw messages
Depth says how much is resting; order-flow imbalance says which way it is moving. Within
a time bucket, count the shares added to the bid and subtract the shares taken off it,
do the same for the ask, and take the difference of the two:
$$\text{OFI} = (\text{bid adds} - \text{bid removes}) - (\text{ask adds} - \text{ask removes})$$
A positive value means the bid side grew relative to the ask, which is the footprint of
buying pressure; a negative value is the reverse.
**What this construction counts.** Adds come from `A` and `F` messages, removals from
`D`, `X`, `E` and `C`. It does not read `U` replaces, so an order moved from one price to
another contributes nothing here. That is a deliberate simplification for building the
quantity by hand: `02_itch_lob_reconstruction` computes an imbalance inside the
reconstruction loop that does process replaces, and Section 5 uses that one. The two
will not agree to the share, and the section that reports a cross-section says which it
is reading.
### Build the order registry
A removal names an order, not a side. So the first step is a registry built from the
add messages: every order reference the day created, with the side, size and price it
was created at.
```python
def load_order_registry(messages_dir: Path, symbol: str) -> pl.DataFrame:
"""Build order registry: order_reference_number → (side, shares, price, timestamp)."""
orders = []
for msg_type in ["A", "F"]:
path = messages_dir / msg_type
if path.exists():
df = (
pl.scan_parquet(path / "*.parquet")
.filter(pl.col("stock") == symbol)
.select(
["order_reference_number", "buy_sell_indicator", "shares", "price", "timestamp"]
)
.collect()
)
if len(df) > 0:
orders.append(df)
if not orders:
return pl.DataFrame()
return pl.concat(orders).rename({"buy_sell_indicator": "side"})
```
### Load the removals
Four message types take shares off the book: `D` deletes what is left of an order, `X`
cancels part of one, and `E` and `C` execute against one. `C` adds an execution price
and a printable flag that `E` does not carry, but neither changes `executed_shares`, so
both take the same shares off the book and leaving `C` out undercounts executions.
`D` carries no share count, because it removes whatever remained: its size is the
original add less the cancels and executions that came before it. Charging a delete the
full original size would count a partially filled order twice, once for the fill and
again for the remainder.
```python
def load_order_removals(messages_dir: Path, order_refs: set) -> pl.DataFrame:
"""Load cancel/delete/execute messages for orders in our registry."""
removals = []
# Delete messages (D) - full order deletion
d_path = messages_dir / "D"
if d_path.exists():
d_df = (
pl.scan_parquet(d_path / "*.parquet")
.filter(pl.col("order_reference_number").is_in(order_refs))
.select(["timestamp", "order_reference_number"])
.collect()
)
if len(d_df) > 0:
d_df = d_df.with_columns(pl.lit("delete").alias("event_type"))
removals.append(d_df)
# Cancel messages (X) - partial cancellation
x_path = messages_dir / "X"
if x_path.exists():
x_df = (
pl.scan_parquet(x_path / "*.parquet")
.filter(pl.col("order_reference_number").is_in(order_refs))
.select(["timestamp", "order_reference_number", "cancelled_shares"])
.collect()
)
if len(x_df) > 0:
x_df = x_df.rename({"cancelled_shares": "shares_removed"})
x_df = x_df.with_columns(pl.lit("cancel").alias("event_type"))
removals.append(x_df)
# Execute messages (E) - shares executed
e_path = messages_dir / "E"
if e_path.exists():
e_df = (
pl.scan_parquet(e_path / "*.parquet")
.filter(pl.col("order_reference_number").is_in(order_refs))
.select(["timestamp", "order_reference_number", "executed_shares"])
.collect()
)
if len(e_df) > 0:
e_df = e_df.rename({"executed_shares": "shares_removed"})
e_df = e_df.with_columns(pl.lit("execute").alias("event_type"))
removals.append(e_df)
# Execute-with-price messages (C) - same shares off the book, priced separately
c_path = messages_dir / "C"
if c_path.exists():
c_df = (
pl.scan_parquet(c_path / "*.parquet")
.filter(pl.col("order_reference_number").is_in(order_refs))
.select(["timestamp", "order_reference_number", "executed_shares"])
.collect()
)
if len(c_df) > 0:
c_df = c_df.rename({"executed_shares": "shares_removed"})
c_df = c_df.with_columns(pl.lit("execute").alias("event_type"))
removals.append(c_df)
if not removals:
return pl.DataFrame()
return pl.concat(removals, how="diagonal")
```
### Aggregate into buckets
Adds and removals are summed within each time bucket and each side, laid out one column
per side, and the four columns combined into the imbalance.
```python
def _pivot_by_side(
df: pl.DataFrame, freq: str, value_col: str, bid_name: str, ask_name: str
) -> pl.DataFrame:
"""Aggregate a DataFrame by time bucket and side, then pivot to wide format."""
agg = (
df.with_columns(pl.col("timestamp").dt.truncate(freq).alias("bucket"))
.group_by(["bucket", "side"])
.agg(pl.col(value_col).sum().alias(value_col))
)
pivot = agg.pivot(on="side", index="bucket", values=value_col).fill_null(0)
if "B" in pivot.columns:
pivot = pivot.rename({"B": bid_name})
if "S" in pivot.columns:
pivot = pivot.rename({"S": ask_name})
return pivot
```
### Enrich Removals with Registry Data
Join order removals with the registry to recover side and share count for each removal event.
```python
def _enrich_removals(removals: pl.DataFrame, registry: pl.DataFrame) -> pl.DataFrame:
"""Join removals with the registry and size each delete at the shares still resting."""
removals = removals.join(
registry.select(["order_reference_number", "side", "shares"]),
on="order_reference_number",
how="left",
)
if "shares_removed" not in removals.columns:
removals = removals.with_columns(pl.lit(None, dtype=pl.Int64).alias("shares_removed"))
# Shares already taken off this order by earlier cancels and executions. A delete has
# to sort after a partial that shares its timestamp, and polars does not keep input
# order for tied keys, so the sort ranks deletes last explicitly.
sized = pl.col("shares_removed").fill_null(0).cast(pl.Int64)
removals = removals.sort(
["order_reference_number", "timestamp", (pl.col("event_type") == "delete")]
).with_columns((sized.cum_sum() - sized).over("order_reference_number").alias("removed_before"))
return removals.with_columns(
pl.when(pl.col("shares_removed").is_null())
.then((pl.col("shares").cast(pl.Int64) - pl.col("removed_before")).clip(lower_bound=0))
.otherwise(pl.col("shares_removed").cast(pl.Int64))
.alias("shares_removed")
).drop("removed_before")
```
### Compute Order Flow Imbalance
Aggregate order book changes into a directional imbalance measure at a given frequency.
```python
def compute_ofi(messages_dir: Path, symbol: str, freq: str) -> pl.DataFrame:
"""Compute Order Flow Imbalance from raw ITCH messages.
Args:
messages_dir: Parsed ITCH message store, one directory per message type.
symbol: Stock to build the imbalance for.
freq: Bucket width, as a polars duration string. The caller passes
`BUCKET_FREQ` so the single stock and the cross-section are measured
over the same interval.
"""
# Build order registry
registry = load_order_registry(messages_dir, symbol)
if registry.is_empty():
return pl.DataFrame()
print(f" {symbol}: {len(registry):,} orders")
# Load removals and join with registry to get sides
removals = load_order_removals(messages_dir, set(registry["order_reference_number"].to_list()))
if not removals.is_empty():
removals = _enrich_removals(removals, registry)
print(f" {symbol}: {len(removals):,} removals")
# Aggregate additions and removals by time bucket
add_pivot = _pivot_by_side(registry, freq, "shares", "bid_adds", "ask_adds")
# The empty case needs the bucket column typed. `pl.DataFrame({"bucket": []})` gives it
# Null, and a join against a datetime key then fails on the schema rather than
# returning the adds unchanged, which is what a session with no removals means.
rem_pivot = (
_pivot_by_side(removals, freq, "shares_removed", "bid_removes", "ask_removes")
if not removals.is_empty()
else pl.DataFrame(schema={"bucket": add_pivot.schema["bucket"]})
)
# A bucket may hold adds with no removals, or removals with no adds, so the join keeps
# both sides and coalesces the key rather than dropping either.
ofi_df = add_pivot.join(rem_pivot, on="bucket", how="full", coalesce=True).fill_null(0)
for col in ["bid_adds", "ask_adds", "bid_removes", "ask_removes"]:
if col not in ofi_df.columns:
ofi_df = ofi_df.with_columns(pl.lit(0).alias(col))
# Share counts arrive as unsigned integers, and imbalance is a signed quantity.
# Cast before subtracting: 100 - 200 in UInt32 wraps to 4,294,967,196.
ofi_df = ofi_df.with_columns(
(
(pl.col("bid_adds").cast(pl.Int64) - pl.col("bid_removes").cast(pl.Int64))
- (pl.col("ask_adds").cast(pl.Int64) - pl.col("ask_removes").cast(pl.Int64))
).alias("ofi")
).sort("bucket")
return ofi_df
```
### Imbalance against the next bucket's return
Line the buckets up with the mid price at the end of each, take the return over the
following bucket, and correlate. The return is shifted backwards by one bucket so that
each imbalance is matched with what happened *after* it, which is the only alignment
that could be predictive rather than contemporaneous.
```python
# Initialize OFI analysis variables
OFI_SYMBOL = list(lob_data.keys())[0] if lob_data else None
ofi_arr = None
ret_arr = None
corr_raw = None
corr_log = None
```
```python
# Compute OFI and align with price data
if MESSAGES_DIR.exists() and OFI_SYMBOL:
print(f"Computing OFI for {OFI_SYMBOL}...")
ofi_df = compute_ofi(MESSAGES_DIR, OFI_SYMBOL, freq=BUCKET_FREQ)
if not ofi_df.is_empty() and OFI_SYMBOL in lob_data:
lob_df = lob_data[OFI_SYMBOL]
if "mid_price" in lob_df.columns:
# Get price buckets
prices = (
lob_df.with_columns(pl.col("timestamp").dt.truncate(BUCKET_FREQ).alias("bucket"))
.group_by("bucket")
.agg(pl.col("mid_price").last())
)
# Align datetime units and join
ofi_df = ofi_df.with_columns(pl.col("bucket").dt.cast_time_unit("us"))
prices = prices.with_columns(pl.col("bucket").dt.cast_time_unit("us"))
ofi_prices = ofi_df.join(prices, on="bucket", how="inner")
# Compute next-period returns
ofi_prices = ofi_prices.sort("bucket").with_columns(
(pl.col("mid_price").pct_change().shift(-1) * 10000).alias("next_return_bps")
)
# Filter valid observations
valid = ofi_prices.filter(
pl.col("ofi").is_not_null() & pl.col("next_return_bps").is_finite()
)
if len(valid) >= 5:
ofi_arr = valid["ofi"].to_numpy()
ret_arr = valid["next_return_bps"].to_numpy()
# Signed log transform: sign(x) * log(1 + |x|)
ofi_log = np.sign(ofi_arr) * np.log1p(np.abs(ofi_arr))
corr_raw = np.corrcoef(ofi_arr, ret_arr)[0, 1]
corr_log = np.corrcoef(ofi_log, ret_arr)[0, 1]
else:
print(
f"Only {len(valid)} buckets survived the return filter; a correlation "
f"needs more than that to mean anything."
)
else:
print("No parsed ITCH messages found. Run 01_itch_parser first.")
```
```python
# Visualize OFI vs returns
if ofi_arr is not None and ret_arr is not None:
ofi_log = np.sign(ofi_arr) * np.log1p(np.abs(ofi_arr))
fig, axes = plt.subplots(1, 2, figsize=(14, 5))
# Panel 1: OFI time series (raw values, symlog scale)
ofi_scaled = ofi_arr / 1000 # Scale to thousands for display
axes[0].bar(range(len(ofi_arr)), ofi_scaled, color=COLORS["slate"], alpha=0.7)
axes[0].axhline(0, color="black", lw=0.5)
axes[0].set_xlabel(f"Time bucket ({BUCKET_FREQ})")
axes[0].set_ylabel("OFI (thousands of shares)")
axes[0].set_title(f"Imbalance per {BUCKET_FREQ} bucket")
axes[0].set_yscale("symlog", linthresh=1) # Log scale for signed data
axes[0].grid(True, alpha=0.3)
# Panel 2: Scatter plot using log-transformed OFI
axes[1].scatter(ofi_log, ret_arr, alpha=0.7, s=80, color=COLORS["slate"])
# Trend line on log-transformed data
z = np.polyfit(ofi_log, ret_arr, 1)
x_line = np.linspace(ofi_log.min(), ofi_log.max(), 100)
axes[1].plot(x_line, np.polyval(z, x_line), color=COLORS["amber"], lw=2, label="Trend")
axes[1].axhline(0, color="black", lw=0.5)
axes[1].axvline(0, color="black", lw=0.5)
axes[1].set_xlabel("OFI (signed log scale)")
axes[1].set_ylabel("Next-bucket return (bps)")
axes[1].set_title("Next-bucket return against order-flow imbalance")
axes[1].legend()
axes[1].grid(True, alpha=0.3)
fig.suptitle(
f"Order-flow imbalance and next-bucket return over {BUCKET_FREQ}, {OFI_SYMBOL}",
fontsize=12,
)
show_with_alt(
fig,
f"Two panels side by side. The left is a bar chart of order-flow imbalance per {BUCKET_FREQ} bucket "
"against bucket number, on a symmetric logarithmic vertical scale with a line at zero, so bars run "
"both above and below. The right is a scatter of the next bucket's return in basis points against the "
"signed log of the same imbalance, with a straight fitted trend line through it and reference lines "
"at zero on both axes.",
)
print(f"Order-flow imbalance for {OFI_SYMBOL}, {BUCKET_FREQ} buckets:")
print(f" Buckets: {len(ofi_arr)}")
print(f" Correlation with next-bucket return, raw imbalance: {corr_raw:.4f}")
print(f" Correlation with next-bucket return, signed log: {corr_log:.4f}")
print(f" Imbalance range: {ofi_arr.min():,.0f} to {ofi_arr.max():,.0f} shares")
```
## 5. The same measurement across fifty stocks
One stock's correlation is one number, and a single number cannot say whether order flow
carries information or whether this stock happened to trend. The cross-section can: run
the same measurement on fifty stocks and read the distribution of correlations rather
than any one of them.
The imbalance used here is the one the reconstruction computed per second, summed into
one-minute buckets - not the hand-built version from Section 4, which ignores replaces.
A symbol reaches the figure only if it contributed at least `MIN_BUCKETS` buckets, so
the number of points is at most fifty and the cell prints what it was on this run
rather than leaving the reader to count them.
The fifty symbols are fixed rather than discovered, so the figure redraws to the same
cross-section on every run. They are drawn in five strata of ten by daily message
count, because a cross-section of only heavily traded names would answer a narrower
question than the one being asked. Reconstructing another symbol is a run of
`02_itch_lob_reconstruction` with `ITCH_SYMBOL` set to it.
### One stock's correlation
Sum the per-second imbalance into `BUCKET_FREQ` buckets, take the mid price at the end of
each bucket, and correlate the bucket's imbalance against the return over the following
bucket. The imbalance enters as a signed log - the sign kept, the magnitude compressed -
because a handful of enormous buckets would otherwise decide a Pearson correlation on
their own.
```python
def compute_ofi_correlation(lob_df: pl.DataFrame, symbol: str) -> dict | None:
"""Correlate bucketed order-flow imbalance against the next bucket's return.
Args:
lob_df: LOB snapshots carrying `ofi` and `mid_price`.
symbol: Stock symbol, for reporting.
Returns:
Dictionary with symbol, bucket count, correlation and imbalance dispersion,
or None when the symbol contributes fewer than MIN_BUCKETS buckets.
"""
if "ofi" not in lob_df.columns:
print(f" {symbol}: snapshots carry no ofi column")
return None
# Aggregate to 1-minute buckets
df = (
lob_df.with_columns(pl.col("timestamp").dt.truncate(BUCKET_FREQ).alias("bucket"))
.group_by("bucket")
.agg(
pl.col("ofi").sum().alias("ofi_bucket"),
pl.col("mid_price").last().alias("mid_price"),
)
.sort("bucket")
)
# Compute next-minute returns
df = df.with_columns(
(pl.col("mid_price").pct_change().shift(-1) * 10000).alias("next_return_bps")
)
# Filter to valid observations
valid = df.filter(pl.col("ofi_bucket").is_not_null() & pl.col("next_return_bps").is_finite())
if len(valid) < MIN_BUCKETS:
return None
ofi_arr = valid["ofi_bucket"].to_numpy()
returns = valid["next_return_bps"].to_numpy()
ofi_log = np.sign(ofi_arr) * np.log1p(np.abs(ofi_arr))
correlation = np.corrcoef(ofi_log, returns)[0, 1]
return {
"symbol": symbol,
"buckets": len(valid),
"corr": correlation,
"ofi_std": float(np.std(ofi_arr)),
}
```
### How active each stock was
The number of add messages a stock received over the day is the plainest measure of how
actively it was quoted, and it is what the cross-section figure uses for its horizontal
axis.
```python
def load_order_counts() -> dict[str, int]:
"""Count each stock's add messages, A and F together.
F is an Add Order carrying the market participant's identifier, so it creates an
order reference exactly as A does. It is 1% of the day's adds overall but a third of
them for some thinly quoted names, and load_order_registry above reads both.
"""
# Scanned per type, not as one file list: F carries an `attribution` column that A
# does not, and a single scan over both raises on the schema mismatch.
counts: dict[str, int] = {}
for msg_type in ("A", "F"):
files = list((MESSAGES_DIR / msg_type).glob("*.parquet"))
if not files:
continue
by_stock = pl.scan_parquet(files).select("stock").group_by("stock").len().collect()
for row in by_stock.iter_rows(named=True):
counts[row["stock"]] = counts.get(row["stock"], 0) + row["len"]
return counts
```
### The cross-section
Run the single-stock measurement over every symbol that has a book, attach its daily
order count, and drop the ones too thinly quoted to estimate.
```python
def analyze_all_stocks_ofi(
lob_data: dict, order_counts: dict[str, int] | None = None
) -> list[dict]:
"""Compute OFI → return correlation for all available stocks."""
if order_counts is None:
order_counts = load_order_counts()
results = []
for symbol, df in lob_data.items():
result = compute_ofi_correlation(df, symbol)
if result is not None:
# Add order count (more intuitive measure of activity)
result["order_count"] = order_counts.get(symbol, 0)
results.append(result)
return sorted(results, key=lambda x: x["corr"], reverse=True)
```
### Reading the cross-section
Each stock is one point: how many orders it received that day against its imbalance-to-
return correlation. The horizontal axis is logarithmic because daily order counts span
several orders of magnitude across the cross-section. What to look at is the spread of
the points around zero and whether it narrows as activity rises.
```python
def plot_ofi_vs_returns(results: list[dict], ax=None):
"""Scatter each stock's daily order count against its imbalance-return correlation.
Args:
results: One dict per stock with 'symbol', 'order_count' and 'corr'.
ax: Optional axis; a new figure is created when None.
"""
if not results:
print("No results to plot")
return
if ax is None:
fig, ax = plt.subplots(figsize=(10, 6))
else:
fig = ax.get_figure()
# Use order_count as x-axis (more intuitive than buckets)
order_counts = np.array([r.get("order_count", r.get("buckets", 0)) for r in results])
corrs = np.array([r["corr"] for r in results])
symbols = [r["symbol"] for r in results]
# Blue for negative, red for positive, on a scale centred on zero and set by the
# widest correlation actually observed, so the neutral colour always means neutral.
color_bound = float(np.max(np.abs(corrs))) if len(corrs) else 1.0
scatter = ax.scatter(
order_counts,
corrs,
s=150,
c=corrs,
cmap="coolwarm",
vmin=-color_bound,
vmax=color_bound,
edgecolor="black",
linewidth=1.5,
)
# Labels for each symbol
for sym, oc, c in zip(symbols, order_counts, corrs, strict=False):
ax.annotate(
sym, (oc, c), textcoords="offset points", xytext=(8, 5), fontsize=9, fontweight="bold"
)
ax.set_xscale("log")
ax.set_xlabel("Add messages received that day (log scale)", fontsize=12)
ax.set_ylabel("Correlation, imbalance to next-bucket return", fontsize=12)
ax.set_title("Imbalance-return correlation by trading activity")
ax.axhline(0, color="gray", linestyle="-", alpha=0.5, linewidth=1.5)
ax.grid(True, alpha=0.3)
cbar = plt.colorbar(scatter, ax=ax, label="Correlation")
cbar.ax.axhline(0, color="gray", linewidth=1)
return fig
```
The five strata below were drawn from the day's stocks by daily add-message count, ten
from each, and then fixed. Fixing them is what makes the figure reproducible: a
cross-section re-sampled on every run would move for reasons that have nothing to do
with order flow.
```python
HIGH_AND_MID_ACTIVITY = [
# Stratum 1: high activity
"QQQ",
"SPY",
"TQQQ",
"IWM",
"AMD",
"DIA",
"AAPL",
"XLK",
"SH",
"MSFT",
# Stratum 2: medium-high activity
"DBO",
"NWS",
"AAWW",
"TAP",
"AGRX",
"RUBI",
"RETA",
"FPX",
"FUT",
"PRTY",
# Stratum 3: medium activity
"CDMO",
"DWX",
"RYF",
"OVID",
"BGR",
"MITK",
"PSCU",
"NWPX",
"ITRN",
"AVGR",
]
LOW_ACTIVITY = [
# Stratum 4: medium-low activity
"AUG",
"ISR",
"ELAT",
"AFMC",
"GHG",
"SBR",
"PBE",
"UBP-K",
"CIK",
"AIRI",
# Stratum 5: low activity
"VGI",
"PMM",
"WINS",
"RCON",
"JOYY",
"ISIG",
"BAC-A",
"CMRE-E",
"LOAC",
"BRN",
]
CROSS_SECTION_SYMBOLS = HIGH_AND_MID_ACTIVITY + LOW_ACTIVITY
```
```python
cross_section_data = {
sym: lob_data[sym]
for sym in CROSS_SECTION_SYMBOLS
if sym in lob_data and "ofi" in lob_data[sym].columns
}
print(
f"Cross-section: {len(cross_section_data)} of {len(CROSS_SECTION_SYMBOLS)} symbols have "
f"a reconstructed book with imbalance"
)
cross_section_results = analyze_all_stocks_ofi(cross_section_data)
print(
f"Of those, {len(cross_section_results)} contributed at least {MIN_BUCKETS} buckets "
f"and carry a correlation"
)
if len(cross_section_results) >= 3:
# The book's figure script rebuilds this chart from the saved cross-section rather
# than re-running the chapter's ITCH pipeline.
results_df = pl.DataFrame(
[
{
"symbol": r["symbol"],
"order_count": r.get("order_count", 0),
"corr": r["corr"],
"buckets": r.get("buckets", 0),
}
for r in cross_section_results
]
)
results_path = OUTPUT_DIR / "ofi_correlation_cross_section.parquet"
results_path.parent.mkdir(parents=True, exist_ok=True)
results_df.write_parquet(results_path)
print(f"Saved cross-section to {display_path(results_path)}")
print(f"Cross-section: {len(cross_section_results)} NASDAQ stocks")
fig = plot_ofi_vs_returns(cross_section_results)
show_with_alt(
fig,
"A scatter chart with one labelled point per stock. The horizontal axis is the number of add messages the stock received that day, on a logarithmic scale spanning several orders of magnitude; the vertical axis is that stock's correlation between bucketed order-flow imbalance and the next bucket's return, with a horizontal line at zero. Points are shaded blue through red by the same correlation, on a scale centred on zero, and a colour bar to the right gives the mapping.",
)
else:
print(
f"Only {len(cross_section_results)} symbol(s) cleared the bucket threshold; the "
f"cross-section needs at least three. Run 02_itch_lob_reconstruction for more "
f"symbols with ITCH_SYMBOL."
)
```
The cross-section as a table, sorted from the most positive correlation to the most
negative, with the spread of the whole set below it.
```python
cross_section_table = pl.DataFrame(
[
{
"symbol": r["symbol"],
"add_messages": r.get("order_count", 0),
"buckets": r.get("buckets", 0),
"correlation": r["corr"],
}
for r in cross_section_results
]
)
cross_section_table
```
```python
if len(cross_section_results) >= 3:
corrs = np.array([r["corr"] for r in cross_section_results])
print(f"Correlations across {len(corrs)} stocks:")
print(f" Range: {corrs.min():.3f} to {corrs.max():.3f}")
print(f" Mean: {corrs.mean():.3f}")
print(f" Standard deviation: {corrs.std():.3f}")
print(f" Positive: {(corrs > 0).sum()} of {len(corrs)}")
```
## 6. Summary Statistics
Compute summary metrics for the LOB data we analyzed.
```python
def compute_liquidity_summary(lob_df: pl.DataFrame) -> dict:
"""Compute liquidity summary metrics from LOB data."""
metrics = {}
# Spread metrics
if "spread" in lob_df.columns:
valid = lob_df.filter(pl.col("spread") > 0)["spread"]
if len(valid) > 0:
metrics["mean_spread"] = valid.mean()
metrics["median_spread"] = valid.median()
# Price
if "mid_price" in lob_df.columns:
metrics["mean_price"] = lob_df["mid_price"].mean()
# Depth at best level
if "bid_size_0" in lob_df.columns:
metrics["mean_bid_depth"] = lob_df["bid_size_0"].mean()
if "ask_size_0" in lob_df.columns:
metrics["mean_ask_depth"] = lob_df["ask_size_0"].mean()
return metrics
```
```python
# Print summary for available stocks
if lob_data:
print("\n" + "=" * 60)
print("LOB Analysis Summary")
print("=" * 60)
for symbol, df in lob_data.items():
metrics = compute_liquidity_summary(df)
print(f"\n{symbol}:")
print(f" Snapshots: {len(df):,}")
if "mean_price" in metrics:
print(f" Avg price: ${metrics['mean_price']:.2f}")
if "mean_spread" in metrics:
spread_bps = metrics["mean_spread"] / metrics.get("mean_price", 1) * 10000
print(f" Avg spread: ${metrics['mean_spread']:.4f} ({spread_bps:.1f} bps)")
if "mean_bid_depth" in metrics and "mean_ask_depth" in metrics:
print(
f" Avg depth: {metrics['mean_bid_depth']:,.0f} bid / {metrics['mean_ask_depth']:,.0f} ask"
)
```
## Key Takeaways
1. **Quote the spread in basis points, not dollars.** A dollar spread mixes the cost of
trading with the price level, so it cannot be compared across stocks or across a
period in which the price moved.
2. **Compress before correlating.** Bucketed order flow is heavy-tailed, and a Pearson
correlation on the raw values is decided by a handful of buckets. The signed log
keeps the direction and takes the scale out. The notebook reports both, and the gap
between them is the size of that effect.
3. **A cross-section answers what one stock cannot.** A single correlation has no
reference; a cross-section has a distribution, and the spread of that distribution
is the quantity to read.
4. **Say which imbalance.** This chapter has two: one built by hand from adds and
removals, which ignores replaces, and one computed inside the reconstruction, which
does not. They are different numbers under one name, so every section names its
source.
5. **Fix the cross-section to make the figure reproducible.** Symbols discovered at run
time move the chart for reasons unconnected to the question.
### Known limitations
- One venue and one day. Everything here is NASDAQ-routed activity on a single session,
so nothing about it establishes what holds over time or across venues.
- The correlations are contemporaneous with the data that produced them: no split, no
holdout, nothing withheld. They describe the day; they do not forecast another one.
- Thinly quoted stocks drop out at the bucket threshold, so the cross-section is
conditioned on having been quoted enough to measure.
- The books carry the top of the book only, so depth here means depth at the touch.
### Bridge to Later Chapters
| Chapter | How These Patterns Connect |
|---------|---------------------------|
| **Chapter 6** | An order-flow strategy that conditions on extreme flow and spread |
| **Chapter 8** | Feature engineering: spread, imbalance, time-of-day |
| **Chapter 9** | Evaluating microstructure signal decay |
| **Chapter 19** | Price impact modeling using liquidity |
---
## Reference
Bouchaud, J.-P., Bonart, J., Donier, J., & Gould, M. (2018).
*Trades, Quotes and Prices: Financial Markets Under the Microscope*.
Cambridge University Press.
[https://doi.org/10.1017/9781009028943](https://doi.org/10.1017/9781009028943)
```python
```



Exibido na íntegra, com atribuição conforme a licença da fonte. Licença: MIT
Este resumo foi escrito pelo agente de pesquisa da Stratmill com base no original; não é uma cópia da fonte.