Streams, values, and alignment#
screamer's single-stream functors (RollingMean, RollingCorr, ...) assume
lockstep alignment: row i of one input pairs with row i of another. Real
multi-stream data breaks that assumption - feeds tick at different rates, arrive
out of step, and drop samples. screamer adds a small, composable set of
operators for combining, splitting, and filtering streams that
do not tick together, while keeping every existing functor unchanged.
The whole design rests on five principles.
1. A stream is values with an optional index#
A stream is a sequence of values (a 1-D or 2-D NumPy array), optionally
annotated with an index (a 1-D array of the same length). The index is an
ordering coordinate - a datetime64 timestamp, an int64 tick count, a
float64 second - that locates each value in a shared timeline. It is not a
dict entry; screamer only ever orders and compares index values.
index=None means positional. A positional stream has no explicit ordering
coordinate; its position (row number) is its logical place in time. Two
positional streams are assumed to tick together, so their rows pair by position
and must have equal length. To model streams on different clocks, provide an
index for each stream.
A stream value is either a bare ndarray (positional) or a (values, index) 2-tuple
(indexed). from_pandas / to_pandas convert between these and pandas objects:
from screamer import to_pandas, from_pandas
values, index = from_pandas(df) # DataFrame -> (values, index) tuple
series = to_pandas(values, index) # (values, index) -> Series or DataFrame
See the tuple convention guide for the full contract and the pandas round-trip.
2. Stream operators are polymorphic#
Most stream operators (CombineLatest, Dropna, Filter, Select,
Resample) dispatch on the type of their inputs and mirror that type on return:
Input type |
Return type |
|---|---|
bare ndarray(s), optional |
|
|
|
graph |
a |
Batch return is always a (values, index) 2-tuple: vals, idx = CombineLatest()(a, b) always
works; idx is None is a checkable flag meaning "no real ordering here." A bare
array is treated as a positional stream with index=None.
Merge and split follow the same input convention but their outputs
differ because they carry a per-event sources tag: Merge returns
(values, sources, index).
split returns per-source (values, index) pairs.
3. Compute functors preserve cardinality; stream operators may change it#
Layer |
Cardinality |
Examples |
|---|---|---|
Compute functors |
preserved (output length == input length) |
|
Stream operators |
may change it |
|
Compute functors handle NaN internally via their nan_policy (see
NaN and warmup) and never add or drop rows. Stream operators own
all time alignment and stream shaping. Dropna / Filter / split are the
cardinality-changing tools; fillna / ffill are shape-preserving and belong
to both worlds.
4. Alignment is a separate layer from computation#
Time-aware stream operators do the index handling and hand aligned data to the unchanged compute functors. The idiom is:
from screamer import CombineLatest, RollingCorr
# Two async price streams, each a (values, index) pair.
aligned, idx = CombineLatest()(p_a, p_b, index=[t_a, t_b]) # as-of latest-value join
corr = RollingCorr(20)(aligned[:, 0], aligned[:, 1]) # functor, untouched
CombineLatest emits one row per distinct index (same-index events from
different streams coalesce into a single settled row). emit="when_all"
(default) waits until every input is warm; emit="on_any" emits from the first
event with NaN for inputs not yet seen. Feed the aligned columns to any
existing functor.
Other stream operators:
Merge()(*values, index=None)->(values, sources, index): one index-sorted, source-tagged stream.split(values, sources, index=None)-> the inverse ofMerge.Dropna(how="any")(values, index=None)-> drop NaN events.Filter()(data, mask)-> keep each data value whose aligned mask is nonzero.MergeandCombineLatestare unified: pass lazy iterators of events and the operator returns a lazy iterator; pass arrays or(values, index)tuples and the operator returns the batch result. No separate*_iterfunction is needed. (splithas no streaming form.)
Migration from the retired *_iter names#
The per-operator streaming variants were removed in v0.5. The unified operators handle both batch and lazy inputs:
Old (removed) |
New |
|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
5. Causality#
Causal: an output at index
tdepends only on events at indices<= t. There is no backward-fill and no lookahead operator, ever. This guarantee holds whether you run the operator on a stored dataset or feed it events one at a time live. Conformance is enforced bytests/test_streams_identity.py.
See also#
Polymorphic API - the single-stream input/output contract; lockstep is the positional (no-index) special case of this page.
Pipelines - wiring stream operators and functions into a graph you define once and run batch or live.
NaN and warmup - how compute functors treat
NaN;ffillis the same forward-fill carry thatCombineLatestuses.