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/rust-arroyo/src/processing/strategies/reduce.rs b/rust-arroyo/src/processing/strategies/reduce.rs index 351c45f7..3d28a50e 100644 --- a/rust-arroyo/src/processing/strategies/reduce.rs +++ b/rust-arroyo/src/processing/strategies/reduce.rs @@ -191,14 +191,26 @@ 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; - if !batch_complete && !force { + if !size_trigger_complete && !time_trigger_complete && !force { return Ok(()); } - timer!("arroyo.strategies.reduce.batch_time.ms", batch_time); + let flush_reason = if size_trigger_complete { + "size" + } else if time_trigger_complete { + "time" + } else { + "force" + }; + + timer!( + "arroyo.strategies.reduce.batch_time.ms", + batch_time, + "flush_reason" => flush_reason + ); let batch_state = mem::replace( &mut self.batch_state, 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)