Protocol Interfaces¶
Contracts that pluggable pipeline components must satisfy. Uses structural (duck) typing — implement the right method signature and it works, no inheritance required.
| Protocol | Method | Purpose |
|---|---|---|
EventLoader |
load(source) → DataFrame |
Parse raw L3 data into events |
DepthSource |
load(source) → DataFrame |
Load an L2 price-level feed into the depth frame |
TradeSource |
load(events, source) → DataFrame |
Build the canonical trades DataFrame |
DataWriter |
write(data, dest) |
Serialize pipeline outputs |
Source |
level · feed_type · trade_attribution · settings |
The shape shared by every data source, file or live |
OfflineSource |
factory methods | A Source that replays stored files (loader, trade source, writer) |
LiveSource |
snapshot · stream · shutdown_synthetic_events |
A Source that captures a live venue feed |
Metric |
compute(result) · prepare(frame) |
A measurement taken from a finished run, drawn as a level-less plot |
A Source declares two coordinates: a FeedType (matched_book vs
diff_feed; see Data quality) and a Level (L2 vs
L3; see Process L2 feeds). It also declares a
TradeAttribution: which orders of a trade its order events can name (see
Check data quality). A source that records the venue's
sequence number may also declare a SequenceKind: whether a skipped number
means a lost message. It carries typed settings and registers in the source
registry (see Sources).
Level ¶
Bases: str, Enum
Order-book resolution a feed (or a plot) works at — MBP vs MBO.
The granularity axis, orthogonal to :class:FeedType's crossing
invariant. A source declares its :attr:Source.level so the
pipeline knows which stages apply:
- :attr:
L2— Market-By-Price (MBP): aggregate volume per price level, with no persistent order identity. Price-level feeds (Binance, Kalshi, Polymarket, most CCXT sources) are L2. The per-order stages (:func:~ob_analytics.analytics.set_order_types, :func:~ob_analytics.analytics.order_aggressiveness, queue reconstruction) have nothing to key on and are skipped; depth / spread / trade analytics run directly on the price-level book. - :attr:
L3— Market-By-Order (MBO): one primitive per resting order, with stable identity (queue position recoverable). The reconstruction model ob-analytics was built for (Bitstamp, LOBSTER, Databento).
The str mixin lets members slot directly into the visualization
renderer-registry tuple keys and, via the :meth:__str__ override, render
as the bare token ("L2") in file stems and f-strings rather than
"Level.L2".
FeedType ¶
Bases: str, Enum
How a data feed represents the order book — its crossing invariant.
A source declares its feed type so downstream code can reason about crossed books by coordinate, not by source name. The distinction is a property of the source, not of the reconstruction:
- :attr:
MATCHED_BOOK— an L3 feed emitted by the venue's own matching engine (LOBSTER, exchange MBO such as Databento). Bids can never rest above asks, so an uncrossed book is a guaranteed invariant of the data. - :attr:
DIFF_FEED— an L3 feed reconstructed from a public placement/cancellation diff stream (the Bitstamp public feed). It can contain genuinely crossed resting orders (a bid resting above an ask, neither filling); :func:~ob_analytics.analytics.order_bookreplays this faithfully — a crossed book in the output is a property of the feed, not a reconstruction bug. - :attr:
UNKNOWN— a source that does not declare its feed type (the structural default for third-party sources predating this attribute).
Mixes in str so members compare and serialise as their value
(FeedType.DIFF_FEED == "diff_feed"), which keeps CLI/JSON output and
equality checks ergonomic.
TradeAttribution ¶
Bases: str, Enum
Which orders of a trade a feed's order events can name.
Every trade has a maker, the order that was resting in the book, and a taker, the order that arrived and traded against it. Whether the order events show both depends on the feed, not on the reconstruction:
- :attr:
BOTH— the feed reports every order the venue accepts, takers included, so both orders of a trade can be found. Bitstamp'slive_orderschannel does this. - :attr:
MAKER_ONLY— the feed shows resting orders only. A taker trades on arrival and never rests, so it never appears. Exchange order-by-order feeds (Nasdaq ITCH, from which LOBSTER is built, and Databento MBO) do not identify the aggressor, and a feed made of book snapshots cannot show it. - :attr:
NONE— the feed has no order identity (L2), so neither order can be named.
A source declares it so the unmatched-trades check in
:func:~ob_analytics.analytics.data_quality_summary counts only the orders
the feed can show. A source that does not declare it is read as
:attr:BOTH, the check's behaviour before this declaration existed.
Mixes in str so members compare and serialise as their value, as
:class:FeedType does.
trade_attribution_of ¶
Return what source declares as its :class:TradeAttribution.
A source that does not declare one is read as :attr:TradeAttribution.BOTH.
SequenceKind ¶
Bases: str, Enum
What a feed's venue sequence number promises.
A source declares it so a gap check knows what a skipped number means:
- :attr:
CONTIGUOUS— every message adds exactly one. A skipped number is a dropped message. This is the default, and what a per-message counter (cryptofeed's L3 channels) gives. - :attr:
MONOTONIC— the number only rises. A skip is normal and says nothing about loss; only a step that does not rise is a fault. CCXT's booknonceis this kind: on Binance it is the last update ID of a diff that covers a range of IDs, andwatch_order_bookcan apply several diffs before it returns a book. Databento'ssequenceis too: it numbers every message on the venue's channel, but a file usually holds one instrument of that channel, and the trade and fill records it carries become trades rather than book events.
Mixes in str so members compare and serialise as their value, as
:class:FeedType does.
sequence_kind_of ¶
Return what source declares as its :class:SequenceKind.
A source that does not declare one is read as
:attr:SequenceKind.CONTIGUOUS.
EventLoader ¶
Bases: Protocol
Loads raw order-book events from a data source.
The returned DataFrame must contain at least the columns required by
ob_analytics.schemas.validate_events_df.
load ¶
Load events from source and return a DataFrame.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
Any
|
Data source identifier. The canonical type is |
required |
Returns:
| Type | Description |
|---|---|
DataFrame
|
Events with at least the columns required by
|
DepthSource ¶
Bases: Protocol
Loads a price-level (L2) depth stream into the canonical depth frame.
The L2 counterpart to :class:EventLoader. Where an EventLoader
returns per-order events that the pipeline later folds into depth via
:func:~ob_analytics.depth.price_level_volume, a DepthSource returns
the depth frame directly — a price-level feed is a depth stream, so
there is nothing to reconstruct.
An :attr:Level.L2 source's :meth:OfflineSource.create_loader returns a
DepthSource; the pipeline validates its output with
:func:~ob_analytics.schemas.validate_depth_df and feeds it straight to
:class:~ob_analytics.depth.DepthMetricsEngine.
Returned DataFrame columns (see
:data:~ob_analytics.schemas.DEPTH_COLUMNS):
timestamp— pandas datetime64[ns, UTC]price— int64, the price level in integer ticks (×tick_sizefor the quote currency)volume— int64, the level's new absolute resting size in integer lots (×lot_sizefor the base asset) after the update (0removes the level); not a signed deltadirection— categoricalbid/ask
load ¶
Load a price-level depth stream from source and return the frame.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
Any
|
Data source identifier — canonically a |
required |
Returns:
| Type | Description |
|---|---|
DataFrame
|
Depth with at least the columns required by
:func: |
TradeSource ¶
Bases: Protocol
Builds the trades DataFrame for a given run.
Implementations read explicit trade records (a separate
trades.csv, LOBSTER execution rows embedded in the events
frame, etc.) and project them into the canonical trades schema.
Returned DataFrame columns:
timestamp— pandas datetime64[ns, UTC]price— int64 (integer ticks; ×tick_sizefor the quote currency)volume— int64 (integer lots; ×lot_sizefor the base asset)direction— categoricalbuy/sell(taker side)maker_event_id— integer event id of the resting ordertaker_event_id— integer event id of the aggressing ordermaker— order id of the resting ordertaker— order id of the aggressing ordermaker_og— original_number of the maker eventtaker_og— original_number of the taker event
load ¶
Build the trades DataFrame.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
events
|
DataFrame
|
The processed events frame (post-loader). |
required |
source
|
Any
|
The same |
required |
Returns:
| Type | Description |
|---|---|
DataFrame
|
|
DataWriter ¶
Bases: Protocol
Writes pipeline results to a format-specific output.
write ¶
Write pipeline DataFrames to dest.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
dict of str to DataFrame
|
Pipeline output keyed by name (e.g. |
required |
dest
|
str or Path
|
Output path (file or directory, format-dependent). |
required |
Source ¶
Bases: Protocol
Structural contract shared by every data source, file or live.
One shape covers both a file loader and a live capturer: a source states
the two coordinates downstream code reasons by — :attr:level (L2 vs L3)
and :attr:feed_type (the crossing invariant) — and carries its typed,
validated :attr:settings (a :class:~ob_analytics.config.SourceSettings)
in place of an untyped dict. There is no base class to inherit: any
object providing these members satisfies the contract (structural typing).
A source declares how it produces the shared schema by also satisfying a
capability protocol — :class:OfflineSource (replay stored files) and/or
:class:LiveSource (capture a live venue). A venue that supports both
(e.g. Bitstamp) satisfies both.
Attributes:
| Name | Type | Description |
|---|---|---|
name |
str
|
Short lowercase identifier registered in
:data: |
level |
Level
|
Order-book resolution this source produces: :attr: |
feed_type |
FeedType
|
The source's crossing invariant (:class: |
trade_attribution |
TradeAttribution
|
Which orders of a trade the source's order events can name
(:class: |
sequence_kind |
SequenceKind
|
What the venue |
settings |
SourceSettings
|
Typed per-source configuration. The empty base for a source that needs
none; a typed subclass (e.g. |
OfflineSource ¶
Bases: Source, Protocol
A :class:Source that replays stored files into the shared schema.
Bundles the per-source factories the pipeline needs to read from a path:
how to load events (or depth), how to acquire trades, and (optionally) how
to write results or compute depth directly. Pass instances to
Pipeline(source=...).
create_loader ¶
Return the loader for this source.
An :attr:Level.L3 source returns an :class:EventLoader (per-order
events); an :attr:Level.L2 source returns a :class:DepthSource
(the price-level depth frame directly).
create_trade_source ¶
Return a trade source for this source.
create_writer ¶
Return a writer for this source, or None if unsupported.
compute_depth ¶
compute_depth(
events: DataFrame,
config: Any,
source: Any,
ctx: RunContext,
) -> tuple[pd.DataFrame, pd.DataFrame] | None
Return (depth, depth_summary) to override the standard
depth pipeline, or None to use it.
config_defaults ¶
Return default :class:PipelineConfig overrides for this source.
required_context ¶
:class:RunContext field names this source requires.
E.g. LOBSTER returns ["trading_date"] because its filenames carry
no date; Bitstamp returns []. Lets the CLI/pipeline validate
required context generically instead of special-casing source names.
Callers should treat a missing method as [] (structural default).
LiveSource ¶
Bases: Source, Protocol
The live-capture capability of a :class:~ob_analytics.protocols.Source.
A live source is an async iterator of order/trade (or depth) events, plus
two bookend methods for snapshot + shutdown synthesis. It inherits the
source coordinates from :class:~ob_analytics.protocols.Source — name,
level, feed_type, and typed settings — so a venue that both
replays files and captures live agrees with itself on those.
:attr:~ob_analytics.protocols.Source.level routes the book events the
runner writes: :attr:~ob_analytics.protocols.Level.L3 → per-order events
to orders.csv; :attr:~ob_analytics.protocols.Level.L2 → price-level
depth updates to depth.csv.
Implementors only worry about parsing. Persistence, raw-frame archival,
and signal handling live in ob_analytics.live._runner.
A source does not reconnect on its own. When it loses its connection it
raises from :meth:stream: a book carried across a disconnect misses
every change made while it was down. :func:~ob_analytics.live.run_capture
then closes the segment, waits, and starts a new one from a fresh snapshot,
recording the time between them as a gap in manifest.json.
A live source MAY additionally implement :class:SupportsDiagnostics to
surface per-run counters in meta.json; that hook is a separate,
optional protocol so it is never required to conform to this one.
snapshot ¶
Yield the opening book: one event per resting order (L3) or level (L2).
Called once at startup before :meth:stream. L3: each event MUST
have action="created"; the runner writes them to orders.csv so
later changed / deleted events have matching creates. L2:
each event is an absolute-size depth row (side / price /
volume) written to depth.csv.
stream ¶
Yield (kind, event, raw_frame) for every live event.
kind is "order" (L3), "depth" (L2), or "trade".
raw_frame is the original JSON-decoded WebSocket payload (or
None); the runner writes it to raw.jsonl iff config.keep_raw.
Implementations should self-terminate after config.minutes of
wall-clock time. The runner also enforces this externally, so
cancellation must be cooperative.
shutdown_synthetic_events ¶
Yield synthetic close-out events for everything still on the book.
Called once at shutdown. L3: each event MUST have
action="deleted", giving every id in orders.csv a complete
created -> ... -> deleted lifecycle. L2: price levels have no
lifecycle to close, so an L2 capturer typically yields nothing here.
Metric ¶
Bases: Protocol
Structural contract for a measurement taken from a finished run.
A metric reads a run's tables and returns one table of its own, then says
how to turn that table into a renderer payload. There is no base class
to inherit: any object providing these members satisfies the contract
(structural typing), and registering it in
:data:~ob_analytics.metrics.METRICS is what makes it run and plot.
:attr:name is both the registry key and the level-less plot concept the
metric draws under, so a renderer registered at (name, None, backend)
is the metric's face.
Attributes:
| Name | Type | Description |
|---|---|---|
name |
str
|
Short lowercase identifier registered in
:data: |
title |
str
|
Human-readable title for the metric's gallery card. |
levels |
tuple of Level
|
The resolutions the metric applies to. A metric that reads per-order
events declares |