Skip to content

feat(batch): add batch_size_bytes flush trigger - #385

Open
enochtangg wants to merge 4 commits into
mainfrom
enochtang/batch-size-bytes-flush-trigger
Open

enochtangg wants to merge 4 commits into
mainfrom
enochtang/batch-size-bytes-flush-trigger

Conversation

@enochtangg

@enochtangg enochtangg commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Adds batch_size_bytes to Batch, 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.

  • Only the Rust adapter supports it. The Python Arroyo adapter raises an error if it's set in code. (A value set only in deployment YAML isn't caught there, because that adapter never applies step overrides. That gap already existed.)
  • It can be set in the deployment config's steps_config, like batch_size and batch_timedelta.
  • Setting it alone isn't valid. A count or time limit is still required.

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 Batch comes after a Parser or Map, the payloads are Python objects, they count as 0, and the byte limit never fires. The existing streams.pipeline.batch.size_bytes metric shows 0 in that case.

Testing

  • Rust unit tests for the byte limit, bytes inside PyAnyMessage, payloads that aren't bytes, no measuring without a limit, and the step flushing through submit/poll.
  • Python tests for config override, validation, and the Python adapter rejecting the setting.
  • Ran it locally against Kafka. Batches closed exactly where expected, for example at 878 messages / 200,184 bytes with 228-byte messages and a 200,000-byte limit.
  • Throughput: 500k messages through a DevNull sink, 7 runs each, about 90,700 msg/s without a limit and 90,400 msg/s with one. The difference is within the noise between runs.

🤖 Generated with Claude Code

enochtangg and others added 3 commits September 23, 2026 15:41
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 evanh left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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
enochtangg marked this pull request as ready for review October 2, 2026 18:20
@enochtangg
enochtangg requested a review from a team as a code owner October 2, 2026 18:20
@enochtangg

Copy link
Copy Markdown
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

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants