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]

Reply via email to