Skip to content

fix: charge RepartitionExec coalescer buffers to the memory pool - #26167

Open
pmoust wants to merge 1 commit into
apache:mainfrom
pmoust:repartition-coalescer-reservation
Open

pmoust wants to merge 1 commit into
apache:mainfrom
pmoust:repartition-coalescer-reservation

Conversation

@pmoust

@pmoust pmoust commented Oct 9, 2026

Copy link
Copy Markdown

Which issue does this PR close?

Rationale for this change

When RepartitionExec coalesces its output, each output partition's coalescer copies incoming rows into its own buffers until it has a full batch. Only completed batches are charged to the output partition's MemoryReservation, at send. The rows the coalescer holds before that are charged nowhere, so the memory pool does not see them. The amount grows with the number of output partitions and with batch_size, and string columns make it larger, because the coalescer copies string view data into new buffers.

What changes are included in this PR?

  • LimitedBatchCoalescer::size() passes through arrow's BatchCoalescer::size().
  • SharedCoalescer keeps, next to the coalescer and under the same mutex, the number of bytes it has charged to the output partition's reservation.
  • After each push, once completed batches are drained, that charge is set to the coalescer's current size. When the coalescer is finalized, the charge is released.

Completed batches are still charged in OutputChannel::send as before. The coalescer's charge drops by the bytes a completed batch took before send charges that batch, so each byte is charged once.

What is the testing strategy for this PR?

New test repartition_reserves_coalescer_buffered_rows in repartition/mod.rs. With batch_size = 100, it pushes 150 rows as five small batches through a RepartitionExec and holds the input open after the last one. The first 100 rows form a completed batch; the last 50 stay in the coalescer. The test:

  • receives the completed batch, then checks that the pool reserves exactly what a LimitedBatchCoalescer fed the same batches reports as its size(). This shows the buffered rows are charged and the completed batch is not counted a second time.
  • releases the input, receives the 50-row residual batch, and checks that the pool is back to 0, so the coalescer's charge is released.

Without the change, the first check fails: the pool reports 0 bytes reserved while the coalescer holds the 50 rows.

Are there any user-facing changes?

No API changes. Memory pools now see more of the memory that RepartitionExec uses. Under a tight memory limit, repartition spills completed batches a little sooner, and other consumers see less free memory, because the buffered rows are now counted. A query that only just fit under its memory limit before can now fail with ResourcesExhausted in an operator that cannot spill.

Open questions

  • The charge uses grow, not try_grow. The buffered rows cannot be spilled, so there is nothing to do when the pool refuses. As a result, the reservation can go over the pool limit by up to one coalescer's size per output partition (about batch_size rows each). Other consumers then see the pool as full and spill or fail. An alternative is to flush the partial batch early when try_grow fails and send it through the existing spill path. That puts a hard limit on the memory, but it produces smaller batches under memory pressure. It also needs a way to flush LimitedBatchCoalescer without finishing it. I kept the smaller change. I am happy to switch if reviewers prefer the hard limit.
  • BatchCoalescer::size() counts the capacity of the in-progress buffers, and for primitive columns the coalescer reserves a full batch_size on the first write. So a partition holding a few rows is charged for a full batch's capacity. I think that is right, because that memory is allocated, but it is more than the bytes of the rows themselves.
  • If the receiver hangs up (for example under a LIMIT), the input task drops that output channel without finalizing the coalescer. The charge then stays until the reservation is dropped with the last input task. The buffered memory has the same lifetime, so the accounting stays correct, but nothing releases it earlier.
  • The reservation is updated while the coalescer mutex is held, so concurrent input tasks pushing to the same output partition keep the stored charge and the reservation in step. MemoryReservation updates are atomic and no code path takes the coalescer mutex from inside a pool call, so I do not expect a lock-order problem. A custom MemoryPool whose grow blocks would now block input tasks of that output partition while they hold the mutex.

The output coalescer in RepartitionExec copies incoming rows into its own
buffers until it has a full batch. Only completed batches were charged to the
output partition's memory reservation, at send. The buffered rows were
outside every reservation, so a memory pool limit did not see them.

Charge the coalescer's size to the output partition's reservation after each
push, once completed batches have been drained, and release it when the
coalescer is finalized. A completed batch is still charged at send, after the
coalescer's charge has dropped by the bytes it took, so each byte is charged
once.

Signed-off-by: Panagiotis Moustafellos <2493339+pmoust@users.noreply.github.com>
@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Oct 9, 2026
@pmoust
pmoust marked this pull request as ready for review October 10, 2026 00:54
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

RepartitionExec output coalescer buffers are not charged to any memory reservation

1 participant