From 0ed391191ffcf2614517566bcbc2db8a5d1e6e7e Mon Sep 17 00:00:00 2001 From: Kynan Rilee Date: Tue, 14 Jul 2026 12:26:42 -0400 Subject: [PATCH 1/2] batch diff_pairs iterator, preparing for an async traversal --- src/interchange/src/envelopes.rs | 282 +++++++++++++++++++++++++------ 1 file changed, 232 insertions(+), 50 deletions(-) diff --git a/src/interchange/src/envelopes.rs b/src/interchange/src/envelopes.rs index 23b8c1d53b6b3..7cd3470800954 100644 --- a/src/interchange/src/envelopes.rs +++ b/src/interchange/src/envelopes.rs @@ -19,80 +19,125 @@ use mz_ore::cast::CastFrom; use mz_repr::{ CatalogItemId, ColumnName, Datum, Diff, Row, RowPacker, SqlColumnType, SqlScalarType, }; +use timely::progress::Antichain; use crate::avro::DiffPair; /// Walks `batch` and invokes `on_diff_pair` for each `DiffPair` at each /// `(key, timestamp)`. /// +/// Thin wrapper around `iter_diff_pairs`. +pub fn for_each_diff_pair(batch: &B, mut on_diff_pair: F) +where + B: BatchReader