API reference: analysis interfaces

Continues the acquisition and storage reference: what operates on a loaded or replayed capture.

Data quality and resource guards

MarketTickStreamer.NON_PRICE_CONDITIONS — Constant
NON_PRICE_CONDITIONS

Sale-condition codes that disqualify a print from setting a price, keyed by tape: "A" and "B" are CTA-processed, "C" is UTP-processed, "O" is the OTC tape. The same character means different things on different tapes — the average price trade is B on CTA and W on UTP, where B is a bunched trade — which is why the lists are not shared.

The lists for A, B and C are derived from the plans' own sale-condition matrices (CTS Pillar output specification 2.11b, UTP data feed services specification 4.1): a code is listed when its consolidated update-last entry is "no", or "no unless it is the only qualifying trade of the day". That covers prints that carry no current price by construction — averages (B, W), prices fixed earlier or elsewhere (P, 4), late reports (Z, G, U), non-regular settlement (C, R, N), contingent trades (V, 7), price variation (H), market-center official open and close (Q, M) — and odd lots (I). On 40 symbol-days captured with this package the stale-price classes among them (W, 4, P, Z) sit one to two orders of magnitude further from the last regular print than regular prints do; odd lots do not, and are listed because the plans bar them. The provider delivers the plans' codes unchanged, and its condition_map returns their glossary.

Two deliberate departures from the matrices. T (extended hours) is not listed although the plans keep it from the consolidated last, which is a regular-session statistic: at tick resolution an extended-hours print is the price path of its session. 9 (corrected consolidated close) is listed although the plans let it update the last: at tick resolution it is a correction message rather than an execution, though for a daily bar it is precisely the official close.

The O list is the provider's published one and has not been checked against a specification.

This is not a per-field eligibility table — the plans track high/low, open, last and volume eligibility separately, and a print may update some and not others. It answers one question, the one a price path needs answered: may this print set a price.

source
MarketTickStreamer.check_live_heap — Method
check_live_heap(limit_mb; context = "") -> Nothing

Config-gated RAM ceiling: if the live heap exceeds limit_mb, force a garbage collection; if it still exceeds the ceiling, fail loudly with a message naming context — a graceful stop beats an OOM kill. Call from long accumulation loops (backfill pages, compaction groups).

source
MarketTickStreamer.coverage_report — Method
coverage_report(paths, provider; feed = "sip", page_limit = 10_000,
                rate_sleep_s = 0.35) -> DataFrame

Compare a live capture against the historical tape: for every symbol in the raw NDJSON paths, count captured trades (after exact-duplicate removal) and query the provider's historical trade count over the same inclusive exchange-time window. coverage = captured / reference; values below 1 quantify feed coverage and stream drops, values above 1 indicate duplicate or spurious prints that dedup did not catch.

source
MarketTickStreamer.deduplicate_trades — Method
deduplicate_trades(trades) -> Vector{Trade}

Remove exact duplicate prints (same symbol, exchange timestamp, id, price, size, exchange), keeping first occurrence and preserving order. Guards against reconnection double-delivery and overlapping backfill/live captures.

source
MarketTickStreamer.observed_round_lot — Method
observed_round_lot(trades; code = "I") -> Float64

The venue's round lot for these prints, read off the tape: one share more than the largest trade still flagged as an odd lot. NaN when no print carries the flag.

Worth reporting per capture because the round lot is neither 100 shares nor constant. Under the SEC Market Data Infrastructure rules it is tiered by share price — 100 shares at or below 250 USD, 40 up to 1000 USD, 10 up to 10 000 USD, 1 above — and reassigned semiannually per symbol from that symbol's average close over an evaluation month, effective the first business day of May and of November.

The consequence is that the odd-lot flag, and therefore the price-forming population, silently changes membership at those dates. On a 251-day AAPL capture taken with this package, the lot fell from 100 to 40 shares on 2026-05-01 and the median inter-arrival time of the price-forming population dropped by a factor of eight across the boundary, with no change in market behaviour. A symbol near a tier threshold can also oscillate: over the same span ERIE went 100, 40, then 100 again inside eight months.

Reporting it turns an invisible redefinition into an observable, and it costs one pass over prints already in memory.

source
MarketTickStreamer.price_forming — Method
price_forming(t::Trade; non_price = NON_PRICE_CONDITIONS) -> Bool

Whether t may set a price, i.e. whether none of its sale conditions is listed for its tape in non_price. A print carrying several conditions is disqualified by any one of them, which is the precedence rule both tape plans state: a single "does not update" overrides every permissive condition beside it.

A print whose tape has no entry in non_price is kept. An unknown tape means no basis on which to exclude, and silently discarding prints on that ground would be the worse error; session_report counts what survives, so the effect stays visible.

This is an analysis-time predicate. Capture is never filtered — the raw layer records every print the venue reported, and which subset constitutes "a trade" is a decision each analysis makes for itself. It is also not one decision: a price path wants price-forming prints only, whereas an arrival process or a waiting-time distribution counts every execution, odd lots and contingent trades included. Filtering the capture would foreclose the second question to answer the first.

source
MarketTickStreamer.session_report — Method
session_report(paths; gap_threshold_s = 60.0) -> DataFrame

Audit raw session files and return one row per symbol:

  • n_trades, n_duplicates (exact duplicate prints)
  • n_price_forming (prints that may set a price, see price_forming; the difference from n_trades is odd lots, contingent and derivatively priced trades, corrections and the like — present in the tape, and in the arrival process, but not in a price path)
  • first_time, last_time (UTC, ms precision, from exchange timestamps)
  • n_out_of_order (exchange timestamps decreasing in arrival order)
  • max_gap_s / n_gaps (largest / count of exchange-time gaps exceeding gap_threshold_s — judge against the instrument's typical activity; quiet symbols gap naturally)
  • median_latency_ms / n_negative_latency (receive minus exchange time, live-captured rows only; negative values indicate clock skew. NaN when the file is pure backfill)
  • round_lot (the venue's round lot, read off the tape by observed_round_lot; it is price-tiered and reassigned semiannually, so it is not 100 and does not hold still)

Inspect this before trusting any captured session.

source

Resampling onto activity clocks

MarketTickStreamer.dollar_bars — Method
dollar_bars(trades, value; keep_partial = false) -> Vector{Bar}

Resample trades onto a traded-value clock: one Bar per value of price times size.

This is the clock that survives changes of scale. It is invariant to splits, and roughly invariant to price-level drift, so a threshold chosen on one sample stays meaningful on another — the property neither tick_bars nor volume_bars has, and the reason dollar bars are the usual default for cross-sample work.

value is in the price's own units; no currency conversion is performed.

source
MarketTickStreamer.tick_bars — Method
tick_bars(trades, n; keep_partial = false) -> Vector{Bar}

Resample trades onto a transaction clock: one Bar per n prints.

The count is of prints as recorded, so whether odd lots and other non-price-forming records count towards it is decided before the call — pass filter_price_forming(trades) to exclude them. On a high-priced name that choice changes the bar count by a factor of three, so it is not incidental.

Sampling by transaction count is the oldest of these clocks: Mandelbrot and Taylor (1967) doi:10.1287/opre.15.6.1057 proposed price changes as a subordinated process running on transaction time, and Ané and Geman (2000) doi:10.1111/0022-1082.00286 showed returns sampled this way are close to normal where calendar-time returns are heavy-tailed.

Prints are ordered by exchange timestamp, not arrival. The print that completes a bar belongs to it. Its timestamp is the bar's close_ns, so consecutive bars satisfy bars[i].close_ns <= bars[i+1].time_ns. A trailing incomplete bar is dropped unless keep_partial, since its threshold — and therefore its comparability to the others — is not met.

tick_bars(trades, 500)                      # one bar per 500 prints
tick_bars(filter_price_forming(trades), 500)  # price-forming prints only
source
MarketTickStreamer.volume_bars — Method
volume_bars(trades, volume; keep_partial = false) -> Vector{Bar}

Resample trades onto a volume clock: one Bar per volume shares.

Clark (1973) doi:10.2307/1913889 introduced volume as the directing process for speculative prices, giving a finite-variance alternative to the stable-Paretian account of heavy tails: the unconditional return distribution is heavy-tailed because it mixes over a random volume clock, not because the underlying increments are.

The caveat is mechanical: a share is not a fixed unit of economic activity across time or across instruments. Share prices drift, splits reset the scale, and a fixed share threshold therefore samples a different amount of value in January than in December. dollar_bars is the usual remedy.

source

Replay

MarketTickStreamer.replay_source — Method
replay_source(paths; pace = "recorded", speed = 1.0, clock = "auto",
              capacity = 100_000) -> Channel{Trade}

Create a channel that replays recorded trades, ordered by the replay clock. paths are either raw NDJSON files (.jsonl) or processed day files (YYYY-MM-DD.csv|.arrow, e.g. from processed_files); the two kinds cannot be mixed.

  • Raw files are read whole and sorted, since arrival order is not chronological across venues and parts. This is the mode for a session.

  • Processed files are streamed one trading day at a time: the files of a date — one per symbol — are loaded, merged on the replay clock and emitted before the next date is opened, so memory is bounded by the busiest day however long the span. This is the mode for a corpus. Recorded pace runs on one schedule across the days, closures included: replaying a week at speed = 1 takes a week, so pass speed, or pace = "max".

  • pace = "recorded": honor the original inter-arrival times of the replay clock against an absolute schedule, compressed by speed (e.g. speed = 60.0 replays an hour in a minute), so the timing error does not accumulate over a session.

  • pace = "max": emit as fast as the consumer takes them (throughput mode).

The replay clock is selected by clock:

  • "recv": local receive time — what a live consumer experienced, network jitter included. Throws if any record lacks it.
  • "exchange": the venue's own timestamps — the only clock a backfilled recording has.
  • "auto" (default): "recv" when every record carries a receive timestamp, otherwise "exchange", logged at @info. For processed files the choice is made on the first day and held.

The channel closes when the recording is exhausted, which cleanly terminates any consumer written against run_sink!-style loops. An error while streaming a later day (an unreadable file, a missing receive clock) closes the channel with that error, and the consumer's iteration rethrows it.

source

Diagnostic figures

Declared and documented in the package, implemented by a package extension that loads with CairoMakie: using MarketTickStreamer, CairoMakie. Without it a call ends in a MethodError that names the package to load.

MarketTickStreamer.overview_figure — Function
overview_figure(symbol, days; tz = tz"America/New_York",
                non_price = NON_PRICE_CONDITIONS,
                session = (9.5, 16.0),
                price_unit = "USD", size_unit = "shares") -> Figure

Multi-day diagnostic figure for one symbol from per-day processed tables (days is a vector of (date, DataFrame) pairs, sorted internally by date):

  • price path on a concatenated trading-time axis (overnight gaps removed, session boundaries dashed, days labeled at their centers), drawn from the prints that price_forming admits under non_price,
  • day × session-minute activity heatmap,
  • pooled intra-session inter-arrival CCDF (overnight gaps excluded by construction — they would contaminate the waiting-time tail; with a 24-hour session midnight is not a boundary, and the wait across it is kept whenever the next calendar day is present),
  • pooled trade-size CCDF with fitted tail exponent.

session is the venue's regular session as exchange-local hours (open, close), 0 <= open < close <= 24; it sets the width of a day on the concatenated axis and the span of the activity heatmap. The default is the US equity session; a venue that never closes takes (0.0, 24.0). price_unit and size_unit name the units on the price and size axes. All three are properties of the venue and come with its ProviderSpec.

Implemented by the CairoMakie extension: using CairoMakie first.

source
MarketTickStreamer.save_overview_figures — Function
save_overview_figures(processed_dir, out_dir;
                      symbols = nothing, from = nothing, to = nothing,
                      formats = ("pdf", "png"), tz = tz"America/New_York",
                      session = (9.5, 16.0), min_days = 2,
                      price_unit = "USD", size_unit = "shares") -> Vector{String}

Render one multi-day overview_figure per symbol from the processed tree (processed_dir/SYMBOL/YYYY-MM-DD.csv|.arrow; safesave _#N backups and .partial files are ignored — the base file per day is authoritative and always the newest). symbols, from, and to restrict the sweep. Output: out_dir/SYMBOL_<from>_<to>.pdf|png (safesave). Days are loaded one symbol at a time, so memory stays bounded by one symbol's span. session is passed to overview_figure.

Implemented by the CairoMakie extension: using CairoMakie first.

source
MarketTickStreamer.save_session_figures — Function
save_session_figures(raw_paths, out_dir;
                     formats = ("pdf", "png"), tz = tz"America/New_York",
                     min_trades = 10,
                     price_unit = "USD", size_unit = "shares") -> Vector{String}

Render one diagnostic figure per (symbol, trading day) found in the raw NDJSON raw_paths, saved as out_dir/SYMBOL_YYYY-MM-DD.pdf|png (safesave — the new figure takes the canonical name, a displaced one is kept as _#N). Groups with fewer than min_trades ticks are skipped with an @info (distribution panels are meaningless). PNG output is written at px_per_unit = 4. Returns the files written.

Implemented by the CairoMakie extension: using CairoMakie first.

source
MarketTickStreamer.session_figure — Function
session_figure(trades; tz = tz"America/New_York",
               non_price = NON_PRICE_CONDITIONS,
               price_unit = "USD", size_unit = "shares") -> Figure

Build the 2×2 diagnostic figure for one symbol's single-day ticks:

  • price path over exchange-local time (HH:MM axis, min–max decimated), drawn from the prints that price_forming admits under non_price, with the admitted share annotated,
  • activity (trades per minute, every print),
  • inter-arrival time survival function P(Δt > x) (log-log, decade ticks) with the median annotated,
  • trade-size survival function P(S > s) with a Hill tail exponent.

trades must be non-empty and single-symbol (as produced by the grouping in save_session_figures); they are sorted internally by exchange time. price_unit and size_unit name the units on the price and size axes; the provider knows them (ProviderSpec).

Implemented by the CairoMakie extension: using CairoMakie first.

source
MarketTickStreamer.tick_theme — Function
tick_theme() -> Theme

Publication defaults: Computer Modern fonts, 26 pt labels over 22 pt ticks, boxed axes with 1.5-wide spines, inward ticks, no minor ticks, faint dashed grey grid, 3-wide data lines, frameless horizontal legends.

Implemented by the CairoMakie extension: using CairoMakie first.

source

Numerics behind the figures

Plotting-free, and usable on their own.

MarketTickStreamer._pooled_gaps — Method
_pooled_gaps(days; continuous) -> (gaps, excluded)

Waiting times [s] pooled over per-day processed tables (date, DataFrame), sorted by date, and the number of between-day gaps left out.

Within a day every gap counts. Between two days the gap spans a closure and is not a waiting time of the arrival process, so it is dropped. On a venue that never closes (continuous = true) midnight is not a boundary: the wait across it is kept whenever the next calendar day is present, and only gaps across missing days are dropped.

source
MarketTickStreamer._tail_fit — Method
_tail_fit(x; frac = 0.1) -> Union{Nothing,NamedTuple}

Hill estimator of the tail exponent α in P(X > x) ∝ x^(-α) from the k largest of the ascending-sorted positive samples x, with k the top frac of the sample and at least 30:

α̂ = k / Σᵢ ln(x₍ₙ₋ᵢ₊₁₎ / x₍ₙ₋ₖ₎),    σ = α̂ / √k.

This is the maximum-likelihood estimator of a Pareto tail above the threshold x₍ₙ₋ₖ₎ (Hill 1975). A least-squares line through the log-log survival function is not a substitute: its points are cumulative and therefore strongly correlated, so the slope is biased and its regression standard error is too small by a large factor (Clauset, Shalizi & Newman 2009).

The threshold is a fixed fraction rather than a fitted x_min, and trade sizes are discrete with mass at round numbers, so the value is a diagnostic of tail weight, not a measurement of a power law. Returns nothing for fewer than 50 samples, for a tail spanning less than a factor of two, or when the tail is degenerate.

source