polars_online.prep¶
Frame preparation: streams whose labels arrive late, and series that tick at their own times (docs/ENHANCEMENTS.md E47, E58).
embargo() turns a frame into the doubled stream
that a forward-looking target needs: every row appears twice, once as a
prediction at its own clock with zero weight, and once as a lesson at
clock + delay, the two merged back into clock order. It is the recipe a
spec’s label_delay runs natively, written out in Polars – useful for
seeing what the delay does, for a model that has no label_delay, and as
the oracle the native path is tested against.
refresh_time() puts asynchronous series on a common grid by
Barndorff-Nielsen, Hansen, Lunde & Shephard’s refresh-time rule: a grid point
wherever every series has ticked at least once since the last one. The scan
is a Rust operator, wrapped here as a lazy source.
Everything here is lazy and streaming – merge_sorted on two sorted halves
of the same frame, a chunk-fed operator for the grid – so a stream too long
to hold is still too long to hold and this does not change that.
- polars_online.prep.embargo(lf: LazyFrame | DataFrame, *, clock: str, delay: float, weight: str | None = None, role: str = '_online_role') LazyFrame[source]¶
The doubled stream for a target that is only known
delaylater.Every row comes back twice, in clock order:
a predict row at
clock, with its weight forced to 0, so the model scores it and learns nothing from it;a learn row at
clock + delay, carrying the same features and target at full weight.
A
rolecolumn says which is which ("predict"/"learn"), so the output is filtered back down without.filter(pl.col(role) == "predict").Why bother: a target that is a forward quantity over
delayclock units is not known at the row it sits on. A stream that learns it there has seendelayof the future before predicting the rows in between, and every “out-of-sample” number after that is contaminated – with an autocorrelated feature, even a pure noise column will show a correlation with its target. Zero-weight rows are legal and mean “advance the clock, learn nothing”, so the doubled stream says exactly what is wanted: predict here, learn later.weightnames an existing weight column; without it the function adds one (namedrole + "_weight") that is 1 on learn rows and 0 on predict rows – pass that name to the spec’sweight=.The frame must already be in
clockorder, as a stream must be. The result is sorted byclockwith learn rows before predict rows at the same clock value: a label whosedelayhas just run out is known at that instant, so a prediction made then may use it. A spec’slabel_delayreleases in the same order, which is what lets the two be compared row for row.delaymust be finite and positive;0would be the undoubled stream, and negative is a label from the past, which is not what this is for.A spec’s
label_delay=does the same thing in the stream with no doubling and no filtering, which is cheaper and does not need the frame rewritten. Reach for this when a delay has to be visible in the data – an oracle, a demonstration, or an engine that is not this one.
- polars_online.prep.refresh_time(lf: LazyFrame | DataFrame, *, series: str, names: Sequence[str], time: str, value: str, by: str | None = None, pairs: bool = False, keep: Sequence[str] = (), chunk_rows: int | None = None) LazyFrame[source]¶
Asynchronous series on a common grid, by refresh time (E58).
Series observed at their own times cannot be correlated directly: the Epps effect attenuates a correlation computed over a fine grid, and filling forward invents observations. Barndorff-Nielsen, Hansen, Lunde & Shephard’s rule places a grid point at the first instant by which every series has ticked at least once since the previous point, and takes each series’ last value there:
tau_0 = max_i (first tick of series i) tau_{j+1} = max_i (first tick of series i strictly after tau_j)Nothing is interpolated – every value in the output was observed – and the grid adapts to the slowest series rather than carrying a stale value across an interval.
The input is long: one row per tick, with a
seriescolumn naming it, atimeand avalue. A wide frame is already synchronised;lf.unpivot(index=[time], variable_name="series", value_name="value")is the line that makes one from the other.namesis required and gives the series in output order. It is not discovered from the data because a lazy plan has to declare its schema before a row is read, and the output columns are named after the series. A row whoseseriesis not innamesis an error naming it: dropping it would hide a misspelling.Output, one row per grid point:
time_refreshThe completing tick’s time – the max over series of their last update, their Definition 1.
<s>_valueEach series’ last value at that instant.
n_obs_<s>Ticks of
ssince the previous point, the first of which is the one on the grid.retained_fractionm / sum(n_obs): how many of the interval’s ticks the grid kept. Look at it before trusting a correlation computed on the result.
plus the
bycolumn – in the dtype it came in as – and anykeepcolumns, at their value on the completing tick.pairs=Trueruns an independent two-series grid per unordered pair instead – which keeps far more of the data when one series is slow – and returns the long frame(by?, pair, time_refresh, a_value, b_value, n_obs_a, n_obs_b, retained_fraction)withpair = "a|b"innamesorder.The staleness caveat (their §2.1): the output looks synchronous and is not. A refresh vector is treated as observed at
time_refresh, but each series’ value is up to one of its own inter-tick intervals old.n_obs_<s>is that staleness made visible: the series with the largest count is the one holding the grid up, and the one whose value is freshest.Rows must be in
timeorder within eachbykey, as a stream must be; a time below the previous row’s is aValueErrornaming the row. A nullvalueis a tick that observed nothing, so it does not update the series. Feeding the input in one chunk or a thousand gives the same grid: a point is a property of the ticks up to it.Ties are broken by row order: “strictly after
tau_j” is read against the row sequence, so a tick carrying the same timestamp as the one that just closed a point, but later in the frame, belongs to the next interval. That is what lets a point be emitted the moment its last series ticks, which is what makes the result chunk-invariant. Sort the input bytimeand by the order you want within a timestamp.ValueErrorfor fewer than twonamesor a duplicate, and for a column the frame has not got.