andygrove commented on code in PR #6440:
URL: https://github.com/apache/datafusion-comet/pull/6440#discussion_r4145824364
##########
native/shuffle/src/writers/local/spill.rs:
##########
@@ -74,14 +77,21 @@ impl PartitionedSpill {
write_buffer_size: usize,
batch_size: usize,
num_partitions: usize,
+ runtime: &RuntimeEnv,
) -> Self {
+ let ranges = vec![Vec::new(); num_partitions];
+ let reservation = MemoryConsumer::new("ShuffleSpill")
Review Comment:
With this change every multi-partition local shuffle registers its memory as
`ShuffleSpill` instead of `ShuffleRepartitioner[{partition}]`. That's the name
`TrackConsumersPool` prints in the top consumers list when anything in the task
fails an allocation, and the one `spark.comet.debug.memory` logs. Almost all of
that reservation is buffered input, so an OOM report would send people looking
at spill bookkeeping, and the partition index is gone. The Celeborn path still
uses the old name.
Could the repartitioner keep registering `ShuffleRepartitioner[{partition}]`
and give the writer a sibling from `reservation.new_empty()` instead of sharing
an `Arc<MemoryReservation>`? Siblings share the consumer id, so the fair pool's
per-consumer ledger still counts input and range metadata against one share.
The repartitioner would keep `free()` and `size()` for its own input, so
`reserved_input_bytes` and the `shrink` in `spill()` could go, and the
`spill()` log line would stop reporting metadata as bytes being spilled. It
also avoids a trap in the shared version. A later `free()` on the
repartitioner's reservation would drop the metadata charge too, and the next
`release_ranges` would panic in `shrink`.
##########
native/shuffle/src/partitioners/multi_partition.rs:
##########
@@ -577,12 +584,15 @@ impl<T: PartitionWriter>
MultiPartitionShuffleRepartitioner<T> {
// A rejected reservation does not include this batch's memory, even
though the batch
// and its partition indices have already been buffered and must be
counted as spilled.
let reservation_failed =
self.reservation.try_grow(mem_growth).is_err();
+ if !reservation_failed {
+ self.reserved_input_bytes += mem_growth;
+ }
// Checking after buffering lets the writer overshoot the limit by at
most one batch,
// which is how the memory-pressure trigger already behaves.
if reservation_failed
|| self
.max_buffer_bytes
- .is_some_and(|limit| self.reservation.size() >= limit)
+ .is_some_and(|limit| self.reserved_input_bytes >= limit)
Review Comment:
Is there a test that shows `max_buffer_bytes` ignores the range metadata?
The Rust tests that set it use `FailingPartitionWriter` or
`CollectingPartitionWriter`, which don't share the reservation, and `Some(1)`,
which spills on every batch either way. The JVM tests use about 10 partitions,
where the metadata is far below the threshold. If the sibling reservation works
out, this holds by construction and the question goes away.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]