feat(batch): add batch_size_bytes flush trigger - #385
Open
enochtangg wants to merge 4 commits into
Open
enochtangg wants to merge 4 commits into
enochtangg wants to merge 4 commits into
Conversation
The batch step closes a window on message count or elapsed time. When a consumer falls behind it reads at broker-fetch speed rather than production rate, so a 3s window collects far more than 3s of messages. On items-span this took batches from ~29,700 to 315,727 messages and the pods crossed their memory limit. The count cap was configured ~55x above anything that occurs, so nothing arrested it. Add a third trigger: close the window once accumulated payload bytes reach a threshold, regardless of arrival rate. Measuring payload bytes needs the GIL and the submit path deliberately does not hold it, so bytes are accumulated in chunks. Each element is measured exactly once and the chunk shrinks as the batch approaches the cap, bounding overshoot without paying a GIL acquisition per message. Pipelines without a byte cap are unaffected: the count and deadline checks run first and stay GIL-free. Only RawMessage and bytes-like PyAnyMessage payloads are counted, so batch_size_bytes does not on its own satisfy Batch.validate(). A batch bounded only by an unmeasurable byte cap would never close. Adds streams.pipeline.batch.flush_reason to show which limit closed a window, and streams.pipeline.batch.byte_cap_unmeasurable to surface a configured cap that can never fire. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The first byte check ran at the full measure chunk (256 elements). There is no mean payload size before that check, so the adaptive chunk cannot kick in, and a batch of large payloads could overshoot the byte cap many times over. Take the first measurement at MIN_BYTE_MEASURE_CHUNK (32) instead. Also add a test that the Python Arroyo adapter rejects batch_size_bytes. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
evanh
reviewed
Oct 2, 2026
evanh
left a comment
Member
There was a problem hiding this comment.
This functionally makes sense to me but is a very complicated change to add `flush if bytes > max_bytes). I would maybe ask the AI to see if it can make the change more minimal.
Replace the chunked, adaptive byte measurement with a running total updated as each message is appended. The consumer thread holds the GIL for as long as it runs, so measuring each message does not wait on anything, and a local benchmark showed no throughput change. The byte limit is now exact: a batch goes over by at most one message. Also drops BatchLimits, FlushReason, the flush_reason metric and the unmeasurable byte cap warning to keep the change small. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
enochtangg
marked this pull request as ready for review
October 2, 2026 18:20
Contributor
Author
|
@evanh Thanks, I managed to simplify it quite a bit. The extra complexity came from avoiding GIL acquisition per message, so we was only measuring bytes every 32 to 256 messages. In practice this probably isn't needed today since the consumer thread already holds the GIL while it runs. To simply, we now check each message and the size is added to a running total, and the batch flushes once the total reaches the limit. |
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Adds
batch_size_bytestoBatch, so a batch also closes once its payloads add up to a byte limit. Today a batch can only be closed by count or time, so its size in bytes depends entirely on message size. Whichever of count, bytes or time is hit first closes the batch.steps_config, likebatch_sizeandbatch_timedelta.How it works
Each message's payload size is added to a running total as it's appended to the batch, and the batch flushes once the total reaches the limit. A batch goes over the limit by at most one message.
Measuring needs the GIL, but the consumer thread already holds it for as long as it runs, so this doesn't wait on anything. Nothing is measured when no byte limit is set.
Limitation
Only raw bytes payloads are counted. If
Batchcomes after a Parser or Map, the payloads are Python objects, they count as 0, and the byte limit never fires. The existingstreams.pipeline.batch.size_bytesmetric shows 0 in that case.Testing
PyAnyMessage, payloads that aren't bytes, no measuring without a limit, and the step flushing throughsubmit/poll.🤖 Generated with Claude Code