Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions src/compute/src/render/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1569,6 +1569,48 @@ mod tests {
assert!(!err.is_empty());
}

/// The passthrough forwards every input record, including those whose key
/// evaluation errored, so a consumer of the bundle's collection sees the
/// unarranged input rather than the ok side of the arrangement.
#[mz_ore::test]
fn arrange_collection_passthrough_forwards_input() {
let rows = test_rows();
let mut expected: Vec<(Row, Timestamp, Diff)> = rows
.iter()
.map(|(row, t)| (row.clone(), Timestamp::from(*t), Diff::ONE))
.collect();
expected.sort();

let key = vec![LirScalarExpr::literal(
Err(EvalError::DivisionByZero),
ReprScalarType::Int32,
)];
let captured = timely::execute_directly(move |worker| {
worker.dataflow::<Timestamp, _, _>(|scope| {
let (mut input, collection) = scope.new_collection();
let (_arranged, _errs, passthrough) =
CollectionBundle::<Timestamp>::arrange_collection(
&"col".to_string(),
vec_to_columnar(collection),
key,
vec![0, 1],
ArrangementBatcher::Columnation,
);
let captured = columnar_to_vec(passthrough).inner.capture();

let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
for (row, time) in rows {
input.update_at(row, Timestamp::from(time), Diff::ONE);
}
input.advance_to(Timestamp::from(max_time + 1));
input.flush();
captured
})
});

assert_eq!(extract_row_updates(captured), expected);
}

fn extract_row_updates(
captured: Captured<(Row, Timestamp, Diff)>,
) -> Vec<(Row, Timestamp, Diff)> {
Expand Down
77 changes: 63 additions & 14 deletions test/sqllogictest/temporal_bucketing.slt
Original file line number Diff line number Diff line change
Expand Up @@ -278,12 +278,12 @@ Target cluster: quickstart
EOF

# -----------------------------------------------------------------------------
# Runtime tests. With `enable_compute_temporal_bucketing` on, the
# bucketed dataflow edges are re-encoded to the columnar representation instead
# of re-wrapping `Vec`. These tests turn the flag on and assert the bucketed
# dataflows still produce the correct logical results, and that a consolidating
# `Union` concatenating a bucketed input with a `Direct` input yields the right
# output, which reaches `concat_many` as a columnar edge like the other leg.
# Runtime tests. These turn `enable_compute_temporal_bucketing` on and assert
# that the bucketed dataflows still produce the correct logical results. There
# is one fixture per site that applies bucketing, because each site hands the

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tiny nit: maybe say "arrangement loop" here, so this doesn't read as covering both ensure_collections sites?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, now reads "the arrangement loop in ensure_collections".

Posted by Claude Code.

# bucketer a differently shaped stream: the keyed `(key, val)` stream in
# `render_reduce`, the TopK input, a consolidating `Union` whose legs carry
# different strategies, and the arrangement loop in `ensure_collections`.
# -----------------------------------------------------------------------------

simple conn=mz_system,user=mz_system
Expand All @@ -295,22 +295,26 @@ statement ok
CREATE TABLE rt_events (k INT NOT NULL, event_time TIMESTAMP NOT NULL)

# Far-future event times so the temporal predicates hold regardless of the
# wall clock at query time, keeping the results deterministic.
# wall clock at query time, keeping the results deterministic. The k=4 row
# expired long ago, so every fixture must drop it: its retraction time precedes
# the dataflow's as_of, which is the case a leaking bucketer would get wrong.
statement ok
INSERT INTO rt_events VALUES
(1, '2999-01-01 00:00:00'),
(1, '2999-01-02 00:00:00'),
(2, '2999-01-03 00:00:00'),
(3, '2999-01-04 00:00:00')
(3, '2999-01-04 00:00:00'),
(4, '2000-01-01 00:00:00')

statement ok
CREATE TABLE rt_other (k INT NOT NULL)

statement ok
INSERT INTO rt_other VALUES (2), (99)

# Bucketed Reduce: temporal filter above a GROUP BY. Hits the arrangement
# re-encode in `context.rs`.
# Bucketed Reduce: temporal filter above a GROUP BY. `render_reduce` buckets the
# keyed `(key, val)` stream itself, so this reaches `reduce.rs` rather than
# either arrangement site.
statement ok
CREATE MATERIALIZED VIEW rt_reduce AS
SELECT k, count(*)
Expand All @@ -325,8 +329,8 @@ SELECT * FROM rt_reduce
2 1
3 1

# Bucketed TopK: temporal filter under ORDER BY ... LIMIT. Hits the re-encode
# in `top_k.rs`.
# Bucketed TopK: temporal filter under ORDER BY ... LIMIT. Buckets the TopK
# input in `top_k.rs`.
statement ok
CREATE MATERIALIZED VIEW rt_topk AS
SELECT k
Expand All @@ -344,8 +348,8 @@ SELECT * FROM rt_topk

# Mixed Union: `EXCEPT ALL` of a temporal-filtered leg (bucketed) against a
# plain relation (Direct) lowers to a consolidating `Union` with per-input
# strategies `[TemporalBucketing, Direct]`. Both legs reach
# `concat_many` as columnar edges.
# strategies `[TemporalBucketing, Direct]`, so `concat_many` sees one bucketed
# leg and one that skipped the bucketer.
query T multiline
EXPLAIN PHYSICAL PLAN AS VERBOSE TEXT FOR
CREATE MATERIALIZED VIEW rt_union AS
Expand Down Expand Up @@ -387,3 +391,48 @@ SELECT * FROM rt_union
1
1
3

# Bucketed arrangement: an index on a temporal-filtered view. `ensure_collections`

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If I'm reading this right, the index is the terminal consumer here and the SELECT reads the arrangement, so the passthrough itself might not have a runtime consumer in this shape? Not blocking, just wondering if that part is covered.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

You are reading it right. The index is the terminal consumer, so the passthrough has no reader in this dataflow, and the unit tests in context.rs all discarded it too. I reworded this comment to claim only the bucketer feeding the arrangement, and added arrange_collection_passthrough_forwards_input, which captures the passthrough under an always-erroring key and checks every input record comes through unchanged.

Posted by Claude Code.

# buckets in the loop that builds the requested arrangements, and hands the result
# to `arrange_collection`. The index is the terminal consumer, so this covers the
# bucketer feeding the arrangement and not the passthrough collection, which has
# no reader in this shape. The plan is pinned because the site is only reached
# while the `ArrangeBy` carries `strategy=TemporalBucketing`.
statement ok
CREATE VIEW rt_indexed AS
SELECT k FROM rt_events WHERE event_time + INTERVAL '45 day' > mz_now()

query T multiline
EXPLAIN PHYSICAL PLAN AS VERBOSE TEXT FOR
CREATE DEFAULT INDEX ON rt_indexed
----
materialize.public.rt_indexed_primary_idx:
ArrangeBy
strategy=TemporalBucketing
raw=true
arrangements[0]={ key=[#0{k}], permutation=id, thinning=() }
Get::PassArrangements materialize.public.rt_indexed
raw=true

materialize.public.rt_indexed:
Get::Collection materialize.public.rt_events
raw=true

Source materialize.public.rt_events
project=(#0)
filter=((mz_now() < timestamp_to_mz_timestamp((#1{event_time} + 45 days))))

Target cluster: quickstart

EOF

statement ok
CREATE DEFAULT INDEX ON rt_indexed

query I rowsort
SELECT k FROM rt_indexed

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think all four rows sit inside the window, so this would catch dropped records but maybe not a leaked expired row? Might not be worth the churn to the earlier pinned results though.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Right on the gap, but there is no churn. Every fixture filters on the same event_time + 45 day > mz_now() predicate below its operator, so a row with event_time in the year 2000 is dropped by all four and none of the pinned results move. Added (4, '2000-01-01') to rt_events. Its retraction time precedes every fixture's as_of, which is the case a leaking bucketer would get wrong.

Posted by Claude Code.

----
1
1
2
3
Loading