Skip to content

Preserve evaluation timestamps in native temporal pane merges - #569

Draft
zzylol wants to merge 1 commit into
stack/509-22-complete-workload-selectionfrom
stack/509-23-temporal-pane-merges
Draft

zzylol wants to merge 1 commit into
stack/509-22-complete-workload-selectionfrom
stack/509-23-temporal-pane-merges

Conversation

@zzylol

@zzylol zzylol commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

Problem: native compilation of a per-series pane merge drops the timestamp

#509 §Example 4 (Materialization of window summaries in physical planning), Pattern B, answers quantile_over_time(0.99, latency_ms[5m]) every minute from 1-min tumbling KLLs:

Candidate At ingestion time At query time
B1 Build one tumbling KLL per minute Merge the latest 5, read p99
B2 Nothing Read 5 min of raw samples, rebuild all 5 tumbling KLLs, merge them, read p99

#509 §Physical operator implementation says a KLL node becomes "summary build, merge and quantile estimation operators", and §4 Execution says the deployment "runs the selected plan as given". So every pane plan that #566 generates and #568 can select must compile and run natively.

That fails when the panes are per-series (PromQL). A per-series state carries a timestamp column (time_index is set). The native binding of SummaryMerge (bind_operation) groups by every column except the state and the time index:

Payload::SummaryMerge => Operator::summary_merge(
    input.clone(),
    state,
    (0..input.fields.len())
        .filter(|&i| i != state && Some(i) != input.time_index)
        .collect(),
),

So the merged native output has no timestamp column, while the Planner output schema still has one. with_output_schema compares the two and fails:

native output type differs from Planner output

The timestamp cannot just be kept as a group key either. Each retained pane carries its own build timestamp (e.g. 00:01, 00:02, … 00:05). Grouping by it would keep five rows instead of merging them, and picking one of them would give the answer a pane's build time instead of the query's evaluation time.

Scope: this PR fixes native compilation of temporal SummaryMerge for both B1 (merge retained panes) and B2 (rebuild panes at query time). It does not add a rotating pane cache, historical EH buckets, or any new planning rule.

Proposed method

The fix is in physical compilation (compile_internal in the native physical planner). When a node is SummaryMerge and its Planner output has a time_index, it lowers to two native operators instead of one:

  1. Merge by series identity. bind_operation builds the usual summary_merge, which groups by all non-state, non-time columns. Its output schema (compact) has no timestamp. It is added under a helper id, helper_id(id, 1). With more than one input, the existing code already inserts a union in front, so the merge reads all panes at once.
  2. Attach the evaluation timestamp. Operator::scope_timestamp(compact, output) maps every compact column to the Planner output schema and fills the time column from the run scope (the query's evaluation time). It is added under the node's own id, so consumers are unchanged.

This is the same pattern the native planner already uses for per-series SummaryAgg (summary_build then scope_timestamp) and for SummaryMerge in the population precompute adapter (precompute.rs: union → summary_merge → scope_timestamp). Non-temporal merges (no time_index) take the old path unchanged.

Key code interfaces

No public API changes. The change is one branch in crates/asap-physical-operators/src/physical_planner/mod.rs, compile_internal:

if matches!(node.payload, Payload::SummaryMerge) && output.time_index.is_some() {
    let merged = bind_operation(node, &schemas)?;
    let compact = merged.schema();
    let merge_id = helper_id(id, 1);
    physical_dag.add(merge_id, inputs, merged)?;
    physical_dag.add(
        id,
        vec![merge_id],
        Operator::scope_timestamp(compact, output)?,
    )?;
    continue;
}

It uses these existing crate-internal functions:

// crates/asap-physical-operators/src/operators/scope_timestamp.rs
pub(crate) fn scope_timestamp(input: SchemaRef, output: SchemaRef) -> Result<Operator, Error>;

// crates/asap-physical-operators/src/physical_planner/mod.rs
fn bind_operation(node: &PostAsapDAGNode, inputs: &[SchemaRef]) -> Result<Operator, Error>;
fn helper_id(node: NodeId, index: u64) -> NodeId;

Test helper: crates/integration-tests/tests/automatic_window_composition.rs now runs one body, fn execute_generated_panes(per_series: bool), from two tests.

Fields

Name Type Meaning
node.payload Payload The new branch matches only Payload::SummaryMerge.
output SchemaRef The Planner's declared output schema for this node. The branch runs only when output.time_index is Some, i.e. the merged state is temporal (per-series).
schemas Vec<SchemaRef> Input schemas. After the existing union step there is one: the union of all pane schemas, which must be equal.
inputs Vec<NodeId> Physical input ids; the union's id when there were several panes.
merged Operator summary_merge from bind_operation: groups by every column except the state column and time_index.
compact SchemaRef merged.schema(): the Planner output schema without the timestamp column.
merge_id NodeId helper_id(id, 1): a helper id derived only from the Planner node id, above the u32 Planner id range. It is the same in every materialization candidate, so frontier cuts need no renumbering.
id NodeId The Planner node id. It now names the scope_timestamp operator, so downstream nodes read the timestamped output.
scope_timestamp(input, output) input = compact, output = the Planner schema. Requires output.time_index to be a non-null Timestamp. Every other output column must match exactly one input column (by type and nullability, and by name for plain columns). Errors if an input column is dropped or repeated. At run time the time column is filled from the execution Scope: evaluation_time_ms for Scope::Query, window_end_ms for Scope::Ingestion.
per_series bool Test parameter. true: scan schema (ts: Timestamp, value: Float64) with time_index = Some(0) and Reduction::PerEntity. false: the old non-temporal fixture (value: Float64) with Reduction::by(vec![]).

Examples

End to end (automatic_per_series_panes_preserve_quantile)

Input:

  • SummaryAgg KLL k=200 over TimeRange(300 s) of metric events, Reduction::PerEntity.
  • Query quantile_over_time(0.99, events[5m]), repeated every 60 000 ms.
  • enumerate_window_compositions returns 2 variants; the test takes the composed one (5 panes merged by SummaryMerge) and adds SummaryEstimate { Quantile { q: 0.99 } }.
  • Each pane's raw input is 20 rows; values 0..99 across the 5 panes, timestamps (pane*20 + i) * 1000 ms.

Native lowering of the merge node after this PR:

pane states ──► union ──► summary_merge (group by series, no ts) ──► scope_timestamp (ts := evaluation time) ──► estimate p99
                          helper_id(id, 1)                           id

Two executions, both with Scope::Query { evaluation_time_ms: 300_000, revision: 1 }:

Path What happens Asserted output
B2: empty frontier, panes rebuilt from raw rows at query time build 5 panes, merge, estimate one row, ts = 300_000, p99 = 98.0
B1: frontier = all SummaryAgg panes; panes built under Scope::Ingestion { 0..300_000 }, then retained the test overwrites each retained pane's timestamp to (pane + 1) * 60_000 (60 000 … 300 000), then runs the query plan on them one row, ts = 300_000, p99 = 98.0

Before this PR, the test failed while compiling with "native output type differs from Planner output".

Accepted vs. changed cases

Merge node Lowering
SummaryMerge, output has time_index (per-series) union (if >1 input) → summary_merge → scope_timestamp (new)
SummaryMerge, no time_index union (if >1 input) → summary_merge, unchanged; automatic_panes_and_materialization_preserve_quantile still passes
Per-series SummaryAgg summary_build → scope_timestamp, unchanged (the pattern this PR follows)

Note for review: the last retained pane's overwritten timestamp (300 000) equals the evaluation time, so the B1 assertion alone does not rule out "take the latest pane timestamp". The code path itself does not read pane timestamps; it groups them away and fills the time column from the scope.

Out of scope

Stack and validation

Stack: #566 → #568 → #569. Base: stack/509-22-complete-workload-selection (#568). Completes native temporal binding for the automatically generated exact pane plans (#509 §Example 4, Pattern B). docs/develop_docs/planner-layering-status.md notes that per-series native merges ignore pane build timestamps and attach the query evaluation timestamp.

Validation: the new regression first failed with "native output type differs from Planner output"; it now passes through direct execution and an automatically enumerated materialization frontier, with different timestamps on retained panes and the query evaluation timestamp asserted. Final workspace validation: 1,638 tests/doctests pass, 2 existing ignores; all-feature Clippy, formatting and diff checks pass.

🤖 Generated with Claude Code

@zzylol

zzylol commented Oct 3, 2026

Copy link
Copy Markdown
Contributor Author

Parked as draft: PR priorities changed (see #528). Order is now (A) finish #511 operator sharing, (B) the #572 crate/module reorganization, (C) #509 end-to-end stages. This PR sits on the old stack/528-legacy-physical-base chain, and Phase B moves the files it touches. Its content will be re-scoped onto the new layout in Phase C.

🤖 Generated with Claude Code

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant