sunchao commented on code in PR #5027:
URL: https://github.com/apache/datafusion-comet/pull/5027#discussion_r4178060865
##########
spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java:
##########
@@ -299,4 +357,511 @@ private static boolean isOneOf(ValueVector result,
ValueVector[] inputs) {
}
return false;
}
+
+ /**
+ * Moves the result's buffer accounting from the task allocator to the root
allocator and returns
+ * the recorded bytes of the chunks the result owns exclusively (see {@link
#chargedOutputSize}),
+ * which drops their Spark task charge. Neither the ownership transfer nor
the eventual FFI
+ * release is observed by the task allocator's {@link AllocationListener}
(Arrow only notifies the
+ * allocator owning a chunk when the chunk is destroyed), so the charge must
be released here.
+ * When Spark granted less than was recorded, {@link TaskState#release}
settles these bytes
+ * against the shortfall first, so the export never hands Spark back more
than it granted.
+ *
+ * <p>The returned vector shares the original buffers and must be closed by
the caller after
+ * export; the exported FFI array keeps the buffers alive until native
execution releases them.
+ */
+ private static FieldVector transferForExport(
+ TaskState state, BufferAllocator outputAllocator, FieldVector result) {
+ long charged = chargedOutputSize(result, outputAllocator);
+ TransferPair transferPair = result.getTransferPair(result.getField(),
ROOT_ALLOCATOR);
+ transferPair.transfer();
+ state.releaseExportedCharge(charged);
+ return (FieldVector) transferPair.getTo();
+ }
+
+ /**
+ * Bytes of the result's buffers currently accounted against the task
allocator, i.e. the recorded
+ * bytes that move to native ownership on export. {@code getAccountedSize()}
is non-zero only on
+ * the ledger that owns a chunk, so pass-through input buffers (owned by the
root allocator) and
+ * chunks whose ownership already moved on a previous export contribute
nothing.
+ *
+ * <p>Buffers are enumerated recursively from each vector's physical field
buffers rather than
+ * through {@code getBuffers(false)}: that method omits the allocated
buffers of zero-length
+ * children (an empty list's data vector still holds its allocated
capacity), but {@code
+ * TransferPair.transfer()} moves those chunks to the root allocator all the
same, and a charge
+ * not counted here would never be released.
+ *
+ * <p>A chunk is counted only when every live reference to its ledger comes
from the result tree
+ * ({@code getRefCount()} equals the number of distinct result buffers on
that ledger). A chunk
+ * shared with a retained scratch vector (the documented per-task
scratch-buffer contract, e.g. an
+ * aligned {@code splitAndTransfer} slice) keeps its Spark charge: when
native execution releases
+ * the FFI result, Arrow hands ownership back to the surviving scratch
ledger without any listener
+ * callback, and the eventual scratch close fires {@link
TaskState#onRelease}, which must then
+ * release a charge exactly once. Dropping the charge at export as well
would release it twice. If
+ * the scratch side instead closes while native still holds the buffers, the
retained charge is
+ * dropped wholesale at task completion, matching the pre-export accounting
model.
+ */
+ private static long chargedOutputSize(FieldVector result, BufferAllocator
outputAllocator) {
+ Set<ArrowBuf> seenBuffers = Collections.newSetFromMap(new
IdentityHashMap<>());
+ IdentityHashMap<ReferenceManager, Integer> resultRefs = new
IdentityHashMap<>();
+ collectPhysicalBuffers(result, seenBuffers, resultRefs);
+ long charged = 0L;
+ for (Map.Entry<ReferenceManager, Integer> entry : resultRefs.entrySet()) {
+ ReferenceManager referenceManager = entry.getKey();
+ if (referenceManager.getAllocator() == outputAllocator
Review Comment:
[P2] [P2] Release export charges from descendant allocators
Could this recognize buffers owned by descendants sharing the task
allocation listener? A custom UDF using `allocator.newChildAllocator(...)`
inherits that listener, so its allocations are charged to Spark. This identity
check excludes those buffers from `chargedOutputSize`, however, and
transferring them to the root bypasses release callbacks. Their Spark charges
therefore accumulate until task completion even after every output is freed.
Repeated batches can exhaust the task's apparent budget and force native
reservations to fail despite no live UDF buffers. Track ownership through the
shared task listener or allocator ancestry, preserving the exclusive-reference
check, and add a descendant-allocator regression test.
Evidence: A bounded public-bridge probe compiled against this exact head
used a 64 MiB off-heap pool and a custom UDF with one child allocator, closed
in `close()`. It evaluated and released 1,024 outputs of 8,192 rows. With
direct task allocation: live Arrow bytes=0, Spark charge=0, subsequent 1 MiB
CometTaskMemoryManager grant=1,048,576. With child allocation: live Arrow
bytes=0, Spark charge=67,108,864, subsequent grant=0. Arrow 18.3.0
BaseAllocator.newChildAllocator inherits the listener, while ownership transfer
does not invoke its release callback. A disposable diagnostic change matching
the owning allocator's listener restored zero charge and the full native grant
in an equivalent control. Reproduction source and output:
/tmp/comet-5027-current-review-fyd55u4u/ChildAllocatorProbe.java and
child-default-batches.log. Native SparkMemory::try_acquire rejects short
grants, confirming the downstream reservation consequence.
--
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]