Resample#
Causal windowed downsampling. Group a stream into fixed index-interval buckets
(freq=) or fixed event-count buckets (count=), reduce each bucket with a
per-bar aggregation, and return a (values, index) tuple. A bucket emits only once a
later index proves it complete; the trailing partial bucket emits at end of input.
Usable eagerly (raw arrays or (values, index) tuples) and inside a Pipeline.
Feeding a lazy iterator of (value, index) pairs returns a lazy iterator of bar events; feeding arrays or (values, index) tuples returns the batch result.
The agg parameter accepts two forms:
String shorthand: one of first, last, min, max, sum, count,
mean, ohlc, ohlcv, ohlcv2. ohlc returns four columns
(open, high, low, close). ohlcv and ohlcv2 accept a two-column
input [price, volume]; see below.
Any EvalOp functor, e.g. ExpandingSkew(). The functor is reset() at
each bar boundary and fed every in-bar sample; its last output before the close
is emitted as the bar value. All screamer functors are valid reducers.
Several reducers at once: run one Resample per statistic over the same
bucketing and align the results with CombineLatest. Because every Resample
shares the same freq= (or count=) and clock, the bars line up exactly and
cannot drift, e.g.
CombineLatest()(Resample(freq=5, agg="first")(price, t), Resample(freq=5, agg="sum")(vol, t)).
Inside a Pipeline, place each per-stat Resample node on its upstream expression
and combine them the same way; per-tick transforms live in the expression, e.g.
Resample(freq=5, agg=ExpandingSum())(PosPart()(vol)).
Bucketing: freq= vs count=#
Exactly one of freq= or count= sets how bars are bounded, and they answer
different questions.
freq=Wbuckets along the index. Barnis the half-open interval[origin + n*W, origin + (n+1)*W), so the index values decide membership. Bars have equal width on the index but a variable number of ticks; a tick exactly on a boundary starts the later bar. Boundaries are anchored atorigin(default0, i.e. multiples ofW), not at the first tick. Setorigin=to shift the grid. Internal empty intervals are real and controlled byfill=.count=Nbuckets by arrival order. A bar closes everyNevents and never consults the index values to place boundaries. Bars have an equal number of ticks but a variable width on the index, and one bar can straddle an arbitrary index gap without noticing it.
The index argument is optional in both modes. count= does not need it to
find boundaries; freq= uses it as the timeline being bucketed. If omitted, row
position (0, 1, 2, ...) is used as the index.
Bar labels depend on the mode and on label=:
freq=: the bar's grid edge,origin + n*Wforlabel="left"(default) ororigin + (n+1)*Wforlabel="right". This is the interval boundary itself, which need not equal any actual tick's index (and a right label can sit past the last tick).count=: an actual tick index, the first tick of the bar forlabel="left"or the last forlabel="right".
Concretely, eight ticks at index [0, 1, 2, 10, 11, 20, 21, 22]:
freq=10-> bars{0,1,2} {10,11} {20,21,22}(counts 3, 2, 3), labels[0, 10, 20](grid edges).count=3-> bars{0,1,2} {10,11,20} {21,22}(counts 3, 3, 2), labels[0, 10, 21](first tick of each bar). The middle bar straddles the11 -> 20gap becausecount=measures rows, not index distance.
Empty buckets: fill=#
fill= controls what happens when a freq= bar contains no samples (an index
interval with no ticks). It applies to eager arrays, (values, index) tuples, and graphs alike;
it is not Pipeline-only.
"skip"(default): emit no row for an empty bucket (the legacy behavior)."nan": emit an all-NaN row at the empty bucket's label."carry": repeat the previous emitted row's values verbatim.
Only internal empty buckets (gaps between two events) are filled by Resample
itself. Trailing empty buckets after the last event are not synthesized here; that
needs a clock, via either pipe.live().advance(now) or a clock input
wired into a Pipeline. The pipe.live() reference has a worked example.
fill= is meaningful only under freq=. With count=, a bar is defined by
holding N events, so empty bars cannot exist by construction and fill= has no
effect.
fill= edge cases#
Leading buckets (before the first event) are never synthesized on the eager path: output starts at the bucket containing the first tick. A tick at index 47 with
freq=10starts at label 40, not 0. Leading empty buckets arise only from a clock (aPipelineclock input oradvance()) that advances time before the first data event; there"nan"emits NaN rows and"carry"skips them, since there is no previous row to repeat."carry"is verbatim and uniform across every column: it repeats the previous bar's whole row, including count-like or sum-like columns. So a genuinely empty bar carries the previous bar's count, not 0. Use"nan"where that is wrong for your columns; a per-column fill policy is out of scope in v1.label="right"composes withfill=unchanged: a filled bucket is labelled at its right grid edge, like any other bar.Functor reducers are unaffected. A filled empty bucket emits a synthetic row (a NaN row, or the carried row) without feeding or resetting the reducer; the reducer is
reset()only at the next real bar boundary. An expanding reducer therefore starts each real bar clean and never accumulates across a filled gap.
Output format and column order#
Every Resample call returns a (values, index) tuple. Unpack it as
values, index = Resample(...)(data, idx). Multi-column aggregations (ohlc,
ohlcv, ohlcv2) return a 2-D values array; columns are positional in
documented order. Use to_pandas(values, index, columns=["open","high","low","close"])
to attach names for display. See OHLC column order
for the full column listing.
ohlcv and ohlcv2 (two-column input)#
Both require values to be a (T, 2) array: column 0 is price, column 1 is
volume (unsigned for ohlcv, signed for ohlcv2).
ohlcv produces (open, high, low, close, volume). The volume column is the sum
of column-1 values inside each bar.
ohlcv2 produces (open, high, low, close, buy_vol, sell_vol). Buy volume is
sum(PosPart(signed_vol)) and sell volume is sum(NegPart(signed_vol)) per bar,
the signed-part decomposition.
Examples#
Mean bar#
Downsample tick values into buckets of width 10 and take the mean of each bucket.
import numpy as np
from screamer import Resample
idx = np.array([0, 3, 10, 12, 20])
vals = np.array([1.0, 2.0, 3.0, 4.0, 5.0])
values, index = Resample(freq=10, agg="mean")(vals, idx)
print(values)
print(index)
[1.5 3.5 5. ]
[ 0 10 20]
OHLCV bars from tick data#
Pass a two-column [price, volume] array and use agg="ohlcv" for a
labelled five-column bar stream.
import numpy as np
from screamer import Resample
np.random.seed(0)
price = np.array([100., 101., 99., 102., 98., 103., 97., 104., 96., 105.])
volume = np.array([10., 20., 15., 30., 12., 22., 18., 25., 14., 28.])
idx = np.arange(10, dtype=np.int64)
values, index = Resample(freq=5, agg="ohlcv")(np.column_stack([price, volume]), idx)
# columns positional: open(0), high(1), low(2), close(3), volume(4)
print(index)
print(values.round(2))
[0 5]
[[100. 102. 98. 98. 87.]
[103. 105. 96. 105. 107.]]
Custom per-bar statistic with a functor#
Any EvalOp functor resets at each bar boundary and accumulates within the
bar. ExpandingSkew() returns the intra-bar price skewness at bar close.
import numpy as np
from screamer import ExpandingSkew, Resample
np.random.seed(0)
price = np.random.normal(100, 1, 20)
idx = np.arange(20, dtype=np.int64)
values, index = Resample(freq=5, agg=ExpandingSkew())(price, idx)
print(values.round(4))
[-0.669 -0.1923 1.1797 0.5104]
Several statistics over the same bucketing#
Run one Resample per statistic and align them with CombineLatest. Every
Resample shares the same freq= and clock, so the bars line up exactly.
import numpy as np
from screamer import ExpandingSkew, ExpandingSlope, Resample, CombineLatest
np.random.seed(0)
price = 100 + np.cumsum(np.random.normal(0, 0.3, 20))
idx = np.arange(20, dtype=np.int64)
skew = Resample(freq=5, agg=ExpandingSkew())(price, idx)
slope = Resample(freq=5, agg=ExpandingSlope())(price, idx)
aligned, bar_idx = CombineLatest()(skew, slope)
print(bar_idx)
print(aligned.round(4))
[ 0 5 10 15]
[[ 0.7594 0.4258]
[-1.5142 0.0587]
[-1.3085 0.1933]
[-1.1415 0.0481]]
Multi-column bars in a Pipeline#
Inside a graph, place a Resample node on each per-tick expression rooted at an
Input, then align them with CombineLatest. All bars share the same freq= and
clock, so they cannot drift. You bind data to the named inputs at call time.
import numpy as np
from screamer import First, Last, ExpandingMax, ExpandingMin, Resample, CombineLatest
from screamer import Input, Pipeline
t_arr = np.arange(10, dtype=np.int64)
px = np.array([100., 101., 99., 102., 98., 103., 97., 104., 96., 105.])
price = Input("price")
open_b = Resample(freq=5, agg=First())(price)
high_b = Resample(freq=5, agg=ExpandingMax())(price)
low_b = Resample(freq=5, agg=ExpandingMin())(price)
close_b = Resample(freq=5, agg=Last())(price)
bars = CombineLatest()(open_b, high_b, low_b, close_b)
ohlc, ohlc_idx = Pipeline([price], [bars])(price=(px, t_arr))
print(ohlc_idx)
print(ohlc.round(2))
[0 5]
[[100. 102. 98. 98.]
[103. 105. 96. 105.]]
Use count= to bucket by a fixed number of events instead of an index interval.