From 0757e36a2199e4d5e996850a9701a69ed1d326dc Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 25 Aug 2026 12:22:36 -0700 Subject: [PATCH 1/5] feat(reduce): track batch trigger reason on batch_time metric. --- rust-arroyo/src/processing/strategies/reduce.rs | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/rust-arroyo/src/processing/strategies/reduce.rs b/rust-arroyo/src/processing/strategies/reduce.rs index 351c45f7..783e833d 100644 --- a/rust-arroyo/src/processing/strategies/reduce.rs +++ b/rust-arroyo/src/processing/strategies/reduce.rs @@ -191,14 +191,19 @@ impl Reduce { } let batch_time = self.batch_state.batch_start_time.elapsed(); - let batch_complete = self.batch_state.message_count >= self.max_batch_size - || batch_time >= self.max_batch_time; + let size_trigger_complete = self.batch_state.message_count >= self.max_batch_size; + let time_trigger_complete = batch_time >= self.max_batch_time; + let batch_complete = size_trigger_complete || time_trigger_complete; if !batch_complete && !force { return Ok(()); } - timer!("arroyo.strategies.reduce.batch_time.ms", batch_time); + timer!( + "arroyo.strategies.reduce.batch_time.ms", + batch_time, + "trigger_reason" => if size_trigger_complete { "size" } else { "time" } + ); let batch_state = mem::replace( &mut self.batch_state, From b514ec26c955e2cb78a02f9dd3151296f7f0db2f Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 25 Aug 2026 14:53:59 -0700 Subject: [PATCH 2/5] Support 'force' trigger reason. --- rust-arroyo/src/processing/strategies/reduce.rs | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/rust-arroyo/src/processing/strategies/reduce.rs b/rust-arroyo/src/processing/strategies/reduce.rs index 783e833d..d76de98b 100644 --- a/rust-arroyo/src/processing/strategies/reduce.rs +++ b/rust-arroyo/src/processing/strategies/reduce.rs @@ -193,16 +193,23 @@ impl Reduce { let batch_time = self.batch_state.batch_start_time.elapsed(); let size_trigger_complete = self.batch_state.message_count >= self.max_batch_size; let time_trigger_complete = batch_time >= self.max_batch_time; - let batch_complete = size_trigger_complete || time_trigger_complete; - if !batch_complete && !force { + if !size_trigger_complete && !time_trigger_complete && !force { return Ok(()); } + let trigger_reason = if size_trigger_complete { + "size" + } else if time_trigger_complete { + "time" + } else { + "force" + }; + timer!( "arroyo.strategies.reduce.batch_time.ms", batch_time, - "trigger_reason" => if size_trigger_complete { "size" } else { "time" } + "trigger_reason" => trigger_reason ); let batch_state = mem::replace( From a69fef4952d29c2d0d7b25d81540a989c2982dd5 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 25 Aug 2026 14:56:13 -0700 Subject: [PATCH 3/5] Rename trigger_reason -> flush_reason. --- rust-arroyo/src/processing/strategies/reduce.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/rust-arroyo/src/processing/strategies/reduce.rs b/rust-arroyo/src/processing/strategies/reduce.rs index d76de98b..a710a03d 100644 --- a/rust-arroyo/src/processing/strategies/reduce.rs +++ b/rust-arroyo/src/processing/strategies/reduce.rs @@ -198,7 +198,7 @@ impl Reduce { return Ok(()); } - let trigger_reason = if size_trigger_complete { + let flush_reason = if size_trigger_complete { "size" } else if time_trigger_complete { "time" @@ -209,7 +209,7 @@ impl Reduce { timer!( "arroyo.strategies.reduce.batch_time.ms", batch_time, - "trigger_reason" => trigger_reason + "flush_reason" => flush_reason ); let batch_state = mem::replace( From d42797c29b3de1a4e361b206effab414e6b32e44 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 25 Aug 2026 15:34:09 -0700 Subject: [PATCH 4/5] Add python support through new BufferProtocol prop. --- arroyo/processing/strategies/buffer.py | 11 +++++++++++ arroyo/processing/strategies/reduce.py | 6 ++++++ tests/processing/strategies/test_buffer.py | 4 ++++ 3 files changed, 21 insertions(+) diff --git a/arroyo/processing/strategies/buffer.py b/arroyo/processing/strategies/buffer.py index a92a5d27..8182e998 100644 --- a/arroyo/processing/strategies/buffer.py +++ b/arroyo/processing/strategies/buffer.py @@ -42,6 +42,14 @@ def is_ready(self) -> bool: """ ... + @property + def readiness_reason(self) -> str: + """Returns why is_ready returned True. + + Only meaningful when is_ready is True. Used for observability. + """ + ... + def append(self, message: BaseValue[TPayload]) -> None: """Accept a TPayload mutating the internal state of the batch builder.""" ... @@ -115,9 +123,12 @@ def __flush(self, force: bool) -> None: ) ) self.__next_step.submit(buffer_msg) + + flush_reason = self.__buffer.readiness_reason if self.__buffer.is_ready else "force" self.__metrics.timing( "arroyo.strategies.reduce.batch_time", time.time() - self.__init_time, + tags={"flush_reason": flush_reason}, ) # Reset to the empty state. diff --git a/arroyo/processing/strategies/reduce.py b/arroyo/processing/strategies/reduce.py index b8d00fa8..bf926e57 100644 --- a/arroyo/processing/strategies/reduce.py +++ b/arroyo/processing/strategies/reduce.py @@ -48,6 +48,12 @@ def is_ready(self) -> bool: or time.time() >= self._buffer_until ) + @property + def readiness_reason(self) -> str: + if self._buffer_size >= self.max_batch_size: + return "size" + return "time" + def append(self, message: BaseValue[TPayload]) -> None: self._buffer = self.accumulator(self._buffer, message) if self.compute_batch_size: diff --git a/tests/processing/strategies/test_buffer.py b/tests/processing/strategies/test_buffer.py index e9e5c7fd..04ba24ae 100644 --- a/tests/processing/strategies/test_buffer.py +++ b/tests/processing/strategies/test_buffer.py @@ -22,6 +22,10 @@ def is_empty(self) -> bool: def is_ready(self) -> bool: return len(self._buffer) >= 3 + @property + def readiness_reason(self) -> str: + return "size" + def append(self, message: BaseValue[int]) -> None: self._buffer.append(message.payload) From dbd44e61b9eb381af60a0ffb6bb1beb9c5ee8c08 Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 25 Aug 2026 16:21:17 -0700 Subject: [PATCH 5/5] lint --- rust-arroyo/src/processing/strategies/reduce.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rust-arroyo/src/processing/strategies/reduce.rs b/rust-arroyo/src/processing/strategies/reduce.rs index a710a03d..3d28a50e 100644 --- a/rust-arroyo/src/processing/strategies/reduce.rs +++ b/rust-arroyo/src/processing/strategies/reduce.rs @@ -200,7 +200,7 @@ impl Reduce { let flush_reason = if size_trigger_complete { "size" - } else if time_trigger_complete { + } else if time_trigger_complete { "time" } else { "force"