This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-6048-67c23a219451d76496c9be19f3fb61b280a1f51e in repository https://gitbox.apache.org/repos/asf/datafusion-comet.git
commit 32fae06d24a1489a1b1d729372585211b4f64185 Author: Andy Grove <[email protected]> AuthorDate: Sun Sep 20 15:24:34 2026 +0000 feat: trace Arrow memory held on the JVM side (#6048) * feat: trace Arrow memory held on the JVM side Comet's tracing counters cover native allocations, memory pool reservations and the JVM heap, but Arrow buffers allocated on the JVM are off-heap, so none of the existing counters see them. A single gauge over the root allocator would not answer the question either. Comet imports batches from native over the Arrow C Data Interface, and Arrow charges an imported buffer to whichever allocator wraps it, so the root's total mixes memory the JVM allocated with native memory that native_allocated already counts. Hold imported buffers in a dedicated child allocator and report two counters. The child reserves nothing, so every byte still escalates to the parent and the root keeps reporting the total, which makes jvm_arrow_allocated minus jvm_arrow_imported the Arrow memory the JVM allocated itself. * refactor: move the Arrow memory counters into Tracing Tracing already owns the JVM tracing seam and its own Native handle, so it is a better home for the counters than the package object, which holds the allocators themselves. Emitting from there also drops the Seq of tuples that the single caller built only to destructure one line later, along with the boxing and the capturing lambda it cost on every traced batch. Fold the two new NativeUtilSuite tests into one. They shared their whole fixture and asserted two views of a single invariant: that imports are charged to the import allocator and still roll up into the root. Assert allocator identity in the UDF probe rather than a byte-count delta, which is the property actually under test and removes the sentinel, the manual reset and the ran-check. * fix: narrow the imported-memory claim and keep UDF output off the import allocator The import allocator does not hold foreign memory exclusively. Arrow's importer allocates the owning ArrowArray struct from it, and loadValidityBuffer allocates a validity bitmap there when an imported vector is all-valid or all-null and carries no validity buffer. Both are bytes the JVM allocated, so the difference between the two counters is a close lower bound on the JVM's own Arrow memory rather than an exact split. Say so in the tracing guide, the allocator scaladoc and the emission scaladoc, and note that the counters are separate reads of process-wide state, so they are neither an atomic per-query balance nor a measure of RSS. A UDF may also allocate its result from the allocator it finds on its inputs, which is the import allocator, and Data.exportVector does not re-own the buffers. That output would then be reported as imported native memory for as long as the export holds it. Transfer the result to the root before exporting when the UDF allocated it elsewhere; TransferPair moves ownership without copying the payload. Add a characterization test pinning the import path's own overhead, and a regression test for a UDF output allocated from the input allocator. Record the new child in the memory management guide's allocator inventory. * fix: drop a redundant string interpolator scalafix's RedundantSyntax rule rejects an `s` prefix on a string part that interpolates nothing. This failed the syntactic lint job and, since the full rule set also runs there, all four Java lint jobs. It went unnoticed locally because `make format` cannot run here: it dies resolving semanticdb-scalac, so scalafix never ran over these changes. The syntactic rules need no semanticdb, so they can be checked directly with the same invocation CI uses. * fix: report the JVM Arrow counters as allocator charges, and keep a returned UDF input on the import allocator Allocator accounting tracks which allocator is accountable for a buffer, not where its bytes were allocated, and the two come apart in both directions. The import allocator is charged for the importer's own `ArrowArray` struct and for a validity bitmap synthesized when an imported vector carries none, while an ownership transfer re-parents a charge without moving the payload. So `jvm_arrow_allocated - jvm_arrow_imported` is not a bound on the Arrow memory the JVM allocated itself, and `jvm_arrow_imported` is not guaranteed to be included in `native_allocated`: a buffer Comet exported to native and that native passed back by reference is imported without the Rust allocator ever having handed it out. Tracking true origin instead is not reachable without vendoring. Arrow routes `wrapForeignAllocation` through the same `allocateBytes` / `onAllocation` / `releaseBytes` path as ordinary allocation, so an `AllocationListener` cannot tell the two apart, and both `ForeignAllocation.release0()` and `memoryAddress()` are protected, so a delegating wrapper that could hook the release would have to live in `org.apache.arrow.memory`. Both counters are therefore described as allocator charges, in the tracing guide, the memory management guide and the two scaladocs. One case is worth fixing rather than documenting. `CometUdfBridge` transfers a result off the import allocator when the UDF allocated it there, but a UDF that returns one of its inputs took that same branch, moving the whole batch's imported charge onto the root: with a 4096-row int input the import allocator was left holding 128 bytes of the 16 KiB. The transfer now skips a result that is reference-identical to an input, and the input loop's close covers it, so the result branch skips it too. Reference identity does not cover a result that merely shares buffers with an input, such as a slice, which is why the contract above is charge-based rather than a corrected list of exceptions. `CometUdfBridgeSuite` gains the regression: an identity UDF, asserting the charge stays on the import allocator, that importing the export back reads the payload, and that the import allocator returns to its prior charge. It fails at 128 of 16384 bytes without the skip. --- .github/workflows/pr_build_linux.yml | 1 + .github/workflows/pr_build_macos.yml | 1 + docs/source/contributor-guide/memory_management.md | 8 + docs/source/contributor-guide/tracing.md | 40 +++- .../java/org/apache/comet/udf/CometUdfBridge.java | 51 ++++- .../scala/org/apache/comet/CometExecIterator.scala | 1 + .../src/main/scala/org/apache/comet/Tracing.scala | 12 ++ .../src/main/scala/org/apache/comet/package.scala | 26 ++- .../scala/org/apache/comet/vector/NativeUtil.scala | 4 +- .../org/apache/comet/udf/CometUdfBridgeSuite.scala | 207 +++++++++++++++++++++ .../org/apache/comet/vector/NativeUtilSuite.scala | 82 +++++++- 11 files changed, 419 insertions(+), 14 deletions(-) diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 17c9da123b..a9c184a20b 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -564,6 +564,7 @@ jobs: org.apache.spark.sql.comet.util.UtilsSuite org.apache.comet.vector.NativeUtilSuite org.apache.comet.vector.CometVectorUtilsSuite + org.apache.comet.udf.CometUdfBridgeSuite org.apache.comet.objectstore.NativeConfigSuite org.apache.comet.serde.operator.CometIcebergNativeScanSuite org.apache.comet.serde.operator.CometNativeScanSuite diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 5a1abca8cf..b47ed5a46f 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -212,6 +212,7 @@ jobs: org.apache.spark.sql.comet.util.UtilsSuite org.apache.comet.vector.NativeUtilSuite org.apache.comet.vector.CometVectorUtilsSuite + org.apache.comet.udf.CometUdfBridgeSuite org.apache.comet.objectstore.NativeConfigSuite org.apache.comet.serde.operator.CometIcebergNativeScanSuite org.apache.comet.serde.operator.CometNativeScanSuite diff --git a/docs/source/contributor-guide/memory_management.md b/docs/source/contributor-guide/memory_management.md index 91199dec35..e7916f9c78 100644 --- a/docs/source/contributor-guide/memory_management.md +++ b/docs/source/contributor-guide/memory_management.md @@ -103,6 +103,14 @@ off-heap bytes in container RSS that neither Spark's `TaskMemoryManager` nor Com pool sees. In practice the volume is modest, a batch at a time per stream, but there is no ceiling and no backpressure. +One further child, `CometArrowImportAllocator` (`comet-ffi-imports`), is what the Arrow C Data +Interface import path allocates from, so that tracing can report those charges apart from the rest +of Comet's Arrow memory. Unlike the others it is process-wide and never closed, because imported +buffers are reference counted and routinely outlive the task that imported them. Its reservation +is zero, so every byte still escalates to the root and the inventory above is unchanged by it. +Being charged there is not the same as having been allocated there; see the scaladoc on the +allocator before reading anything into the split. + **The JVM shuffle allocator is an ordinary Spark consumer.** `CometShuffleMemoryAllocator.getInstance` returns `CometUnifiedShuffleMemoryAllocator`, a Spark `MemoryConsumer` drawing from `spark.memory.offHeap.size`, so shuffle pages are arbitrated against Spark's other consumers in the diff --git a/docs/source/contributor-guide/tracing.md b/docs/source/contributor-guide/tracing.md index 9757bd8a57..bc5690bc59 100644 --- a/docs/source/contributor-guide/tracing.md +++ b/docs/source/contributor-guide/tracing.md @@ -119,12 +119,38 @@ Some excess is expected (allocator metadata and fragmentation for `jemalloc_allo allocations like Arrow IPC buffers for either counter). Large or growing excess may indicate memory that is not being tracked by the pool. +Arrow memory on the JVM side is reported separately, because it is off-heap and so invisible to +`jvm_heap_used`. Comet imports batches from native over the Arrow C Data Interface, and Arrow +charges a buffer to whichever allocator owns it, so those imports are taken against a dedicated +child allocator and reported as `jvm_arrow_imported`, within the `jvm_arrow_allocated` total. + +Both are allocator charges. They report what each allocator is accountable for, not where the bytes +were allocated, and their difference is not a bound on the Arrow memory the JVM allocated itself. +Ownership and allocation come apart in both directions: + +- Bytes the JVM allocated get charged to the import allocator. Arrow's importer allocates the + owning `ArrowArray` struct there, and `BitVectorHelper.loadValidityBuffer` allocates a validity + bitmap there when an imported vector is all-valid or all-null and carries no validity buffer + (512 bytes per 4096 rows). +- Imported bytes get charged to the root. An ownership transfer re-parents a charge without moving + the payload, so a vector that shares buffers with an import, such as a slice of a UDF input, can + leave the root accountable for memory the producer allocated. + +For the same reason neither counter is a count of unique physical bytes, and `jvm_arrow_imported` +is not guaranteed to be included in `native_allocated`. Usually the producer is Rust and the bytes +are counted in both, but a buffer that Comet exported to native and that native passed back by +reference is imported without the Rust allocator ever having handed it out. Finally, the two +counters are separate reads of process-wide state, so concurrent tasks can change them between +samples: they are not an atomic per-query balance, and neither is a measure of RSS. + ## Definition of Labels -| Label | Meaning | -| -------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| jvm_heap_used | JVM heap memory usage of live objects for the executor process | -| jemalloc_allocated | Native memory usage for the executor process (requires `jemalloc` feature) | -| native_allocated | Bytes handed out by the Rust global allocator, process-wide (requires `alloc-accounting` feature). Approximate to within 64 KiB of un-flushed delta per live thread. | -| thread_NNN_comet_memory_reserved | Memory reserved by Comet's DataFusion memory pool (summed across all contexts on the thread). NNN is the Rust thread ID. | -| thread_NNN_comet_jvm_shuffle | Off-heap memory allocated by Comet for columnar shuffle. NNN is the Rust thread ID. | +| Label | Meaning | +| -------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| jvm_heap_used | JVM heap memory usage of live objects for the executor process | +| jemalloc_allocated | Native memory usage for the executor process (requires `jemalloc` feature) | +| jvm_arrow_allocated | Bytes charged to Comet's Arrow allocator tree on the JVM, including buffers imported from native over the Arrow C Data Interface | +| jvm_arrow_imported | Bytes charged to the Arrow C Data Interface import allocator, a subset of `jvm_arrow_allocated`. An allocator charge, not a measure of where the bytes were allocated; see above. | +| native_allocated | Bytes handed out by the Rust global allocator, process-wide (requires `alloc-accounting` feature). Approximate to within 64 KiB of un-flushed delta per live thread. | +| thread_NNN_comet_memory_reserved | Memory reserved by Comet's DataFusion memory pool (summed across all contexts on the thread). NNN is the Rust thread ID. | +| thread_NNN_comet_jvm_shuffle | Off-heap memory allocated by Comet for columnar shuffle. NNN is the Rust thread ID. | diff --git a/spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java b/spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java index d8dea73135..d2881bfe44 100644 --- a/spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java +++ b/spark/src/main/java/org/apache/comet/udf/CometUdfBridge.java @@ -27,6 +27,7 @@ import org.apache.arrow.c.Data; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.vector.FieldVector; import org.apache.arrow.vector.ValueVector; +import org.apache.arrow.vector.util.TransferPair; import org.apache.spark.TaskContext; import org.apache.spark.comet.CometTaskContextShim; import org.apache.spark.util.TaskCompletionListener; @@ -210,17 +211,26 @@ public class CometUdfBridge { assert udf != null : "reflective instantiation returned null for " + udfClassName; BufferAllocator allocator = org.apache.comet.package$.MODULE$.CometArrowAllocator(); + // See CometArrowImportAllocator: inputs are imported against that allocator, so that tracing + // can report the import path's charges apart from the rest of Comet's Arrow memory. + BufferAllocator importAllocator = org.apache.comet.package$.MODULE$.CometArrowImportAllocator(); ValueVector[] inputs = new ValueVector[inputArrayPtrs.length]; ValueVector result = null; + // Whether the UDF handed back one of the vectors it was given. Such a result is closed by the + // input loop below, so the result branch there must leave it alone. + boolean resultIsInput = false; + ValueVector transferred = null; try { for (int i = 0; i < inputArrayPtrs.length; i++) { ArrowArray inArr = ArrowArray.wrap(inputArrayPtrs[i]); ArrowSchema inSch = ArrowSchema.wrap(inputSchemaPtrs[i]); - inputs[i] = Data.importVector(allocator, inArr, inSch, null); + inputs[i] = Data.importVector(importAllocator, inArr, inSch, null); } result = udf.evaluate(inputs, numRows); + // Recorded before the checks below, so the invariant holds however this exits. + resultIsInput = isOneOf(result, inputs); if (!(result instanceof FieldVector)) { throw new RuntimeException( "CometUDF.evaluate() must return a FieldVector, got: " + result.getClass().getName()); @@ -232,9 +242,27 @@ public class CometUdfBridge { + " rows, expected " + numRows); } + // The UDF may allocate its result from the allocator it found on its inputs, which is + // the import allocator. Data.exportVector does not re-own the buffers, so the result + // would stay charged there for as long as the export holds it and be reported as + // imported memory. TransferPair moves ownership without copying the payload. + // + // A result that *is* one of the inputs is left alone: those buffers were imported, so the + // import allocator is the right place for them, and transferring would move a foreign + // charge onto the root. The check is reference identity, so a result that merely shares + // buffers with an input (a slice, say) is still transferred. That is a limit of allocator + // accounting rather than something this can close; see CometArrowImportAllocator. + FieldVector toExport = (FieldVector) result; + if (!resultIsInput && result.getAllocator() != allocator) { + TransferPair transferPair = result.getTransferPair(allocator); + transferPair.transfer(); + transferred = transferPair.getTo(); + toExport = (FieldVector) transferred; + } + ArrowArray outArr = ArrowArray.wrap(outArrayPtr); ArrowSchema outSch = ArrowSchema.wrap(outSchemaPtr); - Data.exportVector(allocator, (FieldVector) result, null, outArr, outSch); + Data.exportVector(allocator, toExport, null, outArr, outSch); } finally { for (ValueVector v : inputs) { if (v != null) { @@ -245,13 +273,30 @@ public class CometUdfBridge { } } } - if (result != null) { + if (result != null && !resultIsInput) { try { result.close(); } catch (RuntimeException ignored) { // do not mask the original throwable } } + if (transferred != null) { + try { + transferred.close(); + } catch (RuntimeException ignored) { + // do not mask the original throwable + } + } + } + } + + /** Whether the UDF handed back one of the vectors it was given, rather than a new one. */ + private static boolean isOneOf(ValueVector result, ValueVector[] inputs) { + for (ValueVector input : inputs) { + if (result == input) { + return true; + } } + return false; } } diff --git a/spark/src/main/scala/org/apache/comet/CometExecIterator.scala b/spark/src/main/scala/org/apache/comet/CometExecIterator.scala index 4da95af18d..4de949b9de 100644 --- a/spark/src/main/scala/org/apache/comet/CometExecIterator.scala +++ b/spark/src/main/scala/org/apache/comet/CometExecIterator.scala @@ -357,6 +357,7 @@ class CometExecIterator( private def traceMemoryUsage(): Unit = { nativeLib.logMemoryUsage("jvm_heap_used", memoryMXBean.getHeapMemoryUsage.getUsed) + Tracing.logArrowMemory() } } diff --git a/spark/src/main/scala/org/apache/comet/Tracing.scala b/spark/src/main/scala/org/apache/comet/Tracing.scala index 13f44ce270..a3cb6fd8b5 100644 --- a/spark/src/main/scala/org/apache/comet/Tracing.scala +++ b/spark/src/main/scala/org/apache/comet/Tracing.scala @@ -23,6 +23,18 @@ object Tracing { private val nativeLib = new Native + /** + * Emits the Arrow memory counters for the JVM side: what Comet's Arrow allocator tree is + * charged for in total, and how much of that the C Data Interface import path is charged for. + * + * Both are allocator charges rather than allocation origin, so their difference is not a bound + * on the Arrow memory the JVM allocated itself. See [[CometArrowImportAllocator]]. + */ + def logArrowMemory(): Unit = { + nativeLib.logMemoryUsage("jvm_arrow_allocated", CometArrowAllocator.getAllocatedMemory) + nativeLib.logMemoryUsage("jvm_arrow_imported", CometArrowImportAllocator.getAllocatedMemory) + } + def withTrace[T](label: String, tracingEnabled: Boolean, fun: => T): T = { try { if (tracingEnabled) { diff --git a/spark/src/main/scala/org/apache/comet/package.scala b/spark/src/main/scala/org/apache/comet/package.scala index 0eb65c9ba6..702359e335 100644 --- a/spark/src/main/scala/org/apache/comet/package.scala +++ b/spark/src/main/scala/org/apache/comet/package.scala @@ -21,7 +21,7 @@ package org.apache import java.util.Properties -import org.apache.arrow.memory.RootAllocator +import org.apache.arrow.memory.{BufferAllocator, RootAllocator} import org.apache.spark.internal.Logging package object comet { @@ -35,6 +35,30 @@ package object comet { */ val CometArrowAllocator = new RootAllocator(Long.MaxValue) + /** + * The allocator that the Arrow C Data Interface import path allocates from. + * + * Arrow charges a buffer to whichever allocator owns it, so imports taken directly against + * [[CometArrowAllocator]] are indistinguishable from buffers the JVM allocated itself. Giving + * the import path its own child keeps the two separable for tracing. The child reserves + * nothing, so every byte still escalates to the parent and the root keeps reporting the total. + * Like the root, it is never closed: imported buffers are reference counted and routinely + * outlive the task that imported them. + * + * What this counts is what the import path is charged for, not where the bytes were allocated. + * Ownership and allocation come apart in both directions. Bytes the JVM allocated land here: + * Arrow's importer allocates the owning `ArrowArray` struct from this allocator, and + * `BitVectorHelper.loadValidityBuffer` allocates a validity bitmap here when an imported vector + * is all-valid or all-null and carries no validity buffer. Imported bytes land elsewhere: an + * ownership transfer re-parents a charge without moving the payload, so a vector that shares + * buffers with an import can leave the root accountable for memory the producer allocated. + * + * So read this and the root's total as allocator charges. Their difference is not a bound on + * the Arrow memory the JVM allocated itself, and neither is a count of unique physical bytes. + */ + val CometArrowImportAllocator: BufferAllocator = + CometArrowAllocator.newChildAllocator("comet-ffi-imports", 0, Long.MaxValue) + /** * Provides access to build information about the Comet libraries. This will be used by the * benchmarking software to provide the source revision and repository. In addition, the build diff --git a/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala b/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala index 173086d2fd..5ef2432173 100644 --- a/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala +++ b/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala @@ -31,7 +31,7 @@ import org.apache.spark.sql.comet.util.Utils import org.apache.spark.sql.execution.vectorized.ConstantColumnVector import org.apache.spark.sql.vectorized.ColumnarBatch -import org.apache.comet.CometArrowAllocator +import org.apache.comet.{CometArrowAllocator, CometArrowImportAllocator} /** * Provides functionality for importing Arrow vectors from native code and wrapping them as @@ -51,7 +51,7 @@ class NativeUtil extends AutoCloseable { private val allocator = CometArrowAllocator /** ArrowImporter does not hold any state and does not need to be closed */ - private val importer = new ArrowImporter(allocator) + private val importer = new ArrowImporter(CometArrowImportAllocator) /** * Dictionary provider to use for the lifetime of this instance of NativeUtil. The dictionary diff --git a/spark/src/test/scala/org/apache/comet/udf/CometUdfBridgeSuite.scala b/spark/src/test/scala/org/apache/comet/udf/CometUdfBridgeSuite.scala new file mode 100644 index 0000000000..7efc5f12cb --- /dev/null +++ b/spark/src/test/scala/org/apache/comet/udf/CometUdfBridgeSuite.scala @@ -0,0 +1,207 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.udf + +import org.scalatest.funsuite.AnyFunSuite + +import org.apache.arrow.memory.BufferAllocator +import org.apache.arrow.vector.{IntVector, ValueVector} +import org.apache.spark.sql.execution.vectorized.ConstantColumnVector +import org.apache.spark.sql.types.IntegerType +import org.apache.spark.sql.vectorized.{ColumnarBatch, ColumnVector} + +import org.apache.comet.{CometArrowAllocator, CometArrowImportAllocator} +import org.apache.comet.vector.NativeUtil + +/** + * A UDF that reports which allocator owns its inputs. + * + * The bridge closes the imported inputs before `evaluate` returns, so this can only be observed + * from inside the call. + */ +class ImportAllocatorProbeUdf extends CometUDF { + override def evaluate(inputs: Array[ValueVector], numRows: Int): ValueVector = { + ImportAllocatorProbeUdf.inputAllocator = Some(inputs.head.getAllocator) + + // Allocated from the root: a UDF result is memory the JVM allocated. + val out = new IntVector("out", CometArrowAllocator) + out.allocateNew(numRows) + (0 until numRows).foreach(row => out.setSafe(row, 7)) + out.setValueCount(numRows) + out + } +} + +object ImportAllocatorProbeUdf { + @volatile var inputAllocator: Option[BufferAllocator] = None +} + +/** + * A UDF that allocates its output from the allocator that owns its inputs. + * + * The `CometUDF` interface permits this, and nothing about it is wrong from the UDF's side. It + * matters here because that allocator is the FFI import allocator, so without intervention the + * output would be charged to it and counted as imported memory. + */ +class InputAllocatorOutputUdf extends CometUDF { + override def evaluate(inputs: Array[ValueVector], numRows: Int): ValueVector = { + val out = new IntVector("out", inputs.head.getAllocator) + out.allocateNew(numRows) + (0 until numRows).foreach(row => out.setSafe(row, 7)) + out.setValueCount(numRows) + out + } +} + +/** + * A UDF that hands back the vector it was given. + * + * The `CometUDF` interface permits this. It matters here because the result's buffers were + * imported, so re-parenting them to the root before export would charge the root for memory the + * producer allocated. + */ +class IdentityUdf extends CometUDF { + override def evaluate(inputs: Array[ValueVector], numRows: Int): ValueVector = inputs.head +} + +class CometUdfBridgeSuite extends AnyFunSuite { + + /** Exports a single-column batch and invokes the bridge against it. */ + private def withBridgeCall(udfClassName: String, numRows: Int)(check: () => Unit): Unit = { + val col = new ConstantColumnVector(numRows, IntegerType) + col.setInt(42) + val batch = new ColumnarBatch(Array[ColumnVector](col), numRows) + + val nativeUtil = new NativeUtil + try { + val (inputArrayAddrs, inputSchemaAddrs, _) = nativeUtil.exportBatchToAddresses(batch) + val (outArrays, outSchemas) = nativeUtil.allocateArrowStructs(1) + + try { + CometUdfBridge.evaluate( + udfClassName, + inputArrayAddrs, + inputSchemaAddrs, + outArrays(0).memoryAddress(), + outSchemas(0).memoryAddress(), + numRows, + null, + null) + + // Checked before releasing the exported structs: the export keeps the result's buffers + // alive, so whichever allocator owns them is still charged at this point. + check() + } finally { + outArrays(0).release() + outArrays(0).close() + outSchemas(0).release() + outSchemas(0).close() + } + } finally { + nativeUtil.close() + } + } + + test("evaluate imports its input vectors against the FFI import allocator") { + withBridgeCall(classOf[ImportAllocatorProbeUdf].getName, 4) { () => + assert( + ImportAllocatorProbeUdf.inputAllocator.contains(CometArrowImportAllocator), + "the UDF's inputs were imported against " + + s"${ImportAllocatorProbeUdf.inputAllocator.map(_.getName)}, so their bytes are not " + + "reported as imported memory") + } + } + + test("a UDF output allocated from the input allocator is not left charged as imported") { + // The CometUDF interface lets a UDF allocate its result from inputs.head.getAllocator, which + // is the import allocator. Data.exportVector does not re-own the buffers, so without a + // transfer the output stays charged to the import allocator for as long as the export holds + // it, and jvm_arrow_imported counts JVM-created bytes as imported native memory. + val numRows = 4096 + val before = CometArrowImportAllocator.getAllocatedMemory + + withBridgeCall(classOf[InputAllocatorOutputUdf].getName, numRows) { () => + val during = CometArrowImportAllocator.getAllocatedMemory + assert( + during - before < numRows.toLong * 4, + s"the UDF's output is still charged to the import allocator ($before -> $during bytes, " + + s"output is ${numRows * 4} bytes), so it would be reported as imported native memory") + } + } + + test("a UDF that returns one of its inputs leaves the charge on the import allocator") { + // The transfer above exists for a UDF that allocates its output from the import allocator. A + // UDF that returns an input is the opposite case: those buffers really were imported, so + // transferring them would move a charge for the producer's memory onto the root and inflate + // the JVM-side reading. Reference identity does not cover a result that merely shares buffers + // with an input, which is why the counters are documented as allocator charges rather than as + // a measure of where the bytes were allocated. + val numRows = 4096 + val importBefore = CometArrowImportAllocator.getAllocatedMemory + + val col = new ConstantColumnVector(numRows, IntegerType) + col.setInt(42) + val batch = new ColumnarBatch(Array[ColumnVector](col), numRows) + + val nativeUtil = new NativeUtil + try { + val (inputArrayAddrs, inputSchemaAddrs, _) = nativeUtil.exportBatchToAddresses(batch) + val (outArrays, outSchemas) = nativeUtil.allocateArrowStructs(1) + + CometUdfBridge.evaluate( + classOf[IdentityUdf].getName, + inputArrayAddrs, + inputSchemaAddrs, + outArrays(0).memoryAddress(), + outSchemas(0).memoryAddress(), + numRows, + null, + null) + + // Read before the export is consumed. The bridge has dropped its own reference to the + // inputs, so what remains charged is the export's, and it must still be on the import + // allocator. + val during = CometArrowImportAllocator.getAllocatedMemory - importBefore + assert( + during >= numRows.toLong * 4, + s"the returned input was re-parented off the import allocator ($during bytes charged, " + + s"its data buffer alone is ${numRows * 4}), so the root is now charged for memory the " + + "producer allocated") + + // Importing the export back shows the bridge's cleanup left the payload alive and readable. + val result = nativeUtil.importVector(outArrays, outSchemas) + try { + val vector = result.head + assert(vector.getValueVector.getValueCount == numRows) + assert( + (0 until numRows).forall(i => vector.getInt(i) == 42), + "the exported payload did not survive the bridge's cleanup") + } finally { + result.foreach(_.close()) + } + } finally { + nativeUtil.close() + } + + assert( + CometArrowImportAllocator.getAllocatedMemory == importBefore, + "the import allocator did not return to its prior charge once the export was released") + } +} diff --git a/spark/src/test/scala/org/apache/comet/vector/NativeUtilSuite.scala b/spark/src/test/scala/org/apache/comet/vector/NativeUtilSuite.scala index ec9dde945c..fb22dd0abc 100644 --- a/spark/src/test/scala/org/apache/comet/vector/NativeUtilSuite.scala +++ b/spark/src/test/scala/org/apache/comet/vector/NativeUtilSuite.scala @@ -22,6 +22,7 @@ package org.apache.comet.vector import java.io.IOException import java.nio.charset.StandardCharsets +import scala.jdk.CollectionConverters._ import scala.util.Using import org.apache.arrow.c.{ArrowArray, ArrowSchema, Data} @@ -38,7 +39,7 @@ import org.apache.spark.sql.execution.vectorized.ConstantColumnVector import org.apache.spark.sql.types.{IntegerType, StringType, StructField, StructType} import org.apache.spark.sql.vectorized.{ColumnarBatch, ColumnVector} -import org.apache.comet.CometConf +import org.apache.comet.{CometArrowAllocator, CometArrowImportAllocator, CometConf} import org.apache.comet.serde.{OperatorOuterClass, QueryPlanSerde} class NativeUtilSuite extends CometTestBase { @@ -371,6 +372,85 @@ class NativeUtilSuite extends CometTestBase { } } + test("imports are charged to the import allocator and roll up into the root") { + // Arrow charges a buffer to whichever allocator owns it, so the tracing counters can only + // report the import path apart from the rest of Comet's Arrow memory if imports go to their + // own allocator. The child must still roll up into the root, because jvm_arrow_imported is + // reported as a subset of jvm_arrow_allocated. + val numRows = 4 + val col = new ConstantColumnVector(numRows, IntegerType) + col.setInt(42) + val batch = new ColumnarBatch(Array[ColumnVector](col), numRows) + + val nativeUtil = new NativeUtil + var imported: ColumnarBatch = null + try { + val (arrayAddrs, schemaAddrs, _) = nativeUtil.exportBatchToAddresses(batch) + val vectors = + nativeUtil.importVector( + arrayAddrs.map(ArrowArray.wrap), + schemaAddrs.map(ArrowSchema.wrap)) + imported = new ColumnarBatch(vectors.toArray, numRows) + + assert( + CometArrowAllocator.getChildAllocators.asScala.exists(_ eq CometArrowImportAllocator), + "the FFI import allocator must be a child of the root, so the root keeps reporting the " + + "total across both") + val importedBytes = CometArrowImportAllocator.getAllocatedMemory + assert( + importedBytes > 0, + "imported buffers were charged somewhere other than the FFI import allocator") + assert( + CometArrowAllocator.getAllocatedMemory >= importedBytes, + "the root's total must include imported bytes, otherwise the subtraction is meaningless") + } finally { + if (imported != null) { + imported.close() + } + nativeUtil.close() + } + } + + test("the import allocator is also charged for the JVM-side cost of importing") { + // Characterization, not an aspiration: Arrow's importer allocates the owning ArrowArray + // struct from the import allocator (ArrayImporter calls ArrowArray.allocateNew(allocator)), + // and loadValidityBuffer allocates a validity bitmap there when an imported vector is + // all-valid and carries no validity buffer. So the import allocator is charged for more than + // the imported buffers, which is one half of why jvm_arrow_imported is an allocator charge + // rather than a measure of where the bytes were allocated. Compared against every imported + // buffer, not just the data one, so the excess measured here is JVM-allocated rather than the + // foreign validity buffer. + val numRows = 4096 + val col = new ConstantColumnVector(numRows, IntegerType) + col.setInt(42) + val batch = new ColumnarBatch(Array[ColumnVector](col), numRows) + + val nativeUtil = new NativeUtil + var imported: ColumnarBatch = null + val before = CometArrowImportAllocator.getAllocatedMemory + try { + val (arrayAddrs, schemaAddrs, _) = nativeUtil.exportBatchToAddresses(batch) + val vectors = + nativeUtil.importVector( + arrayAddrs.map(ArrowArray.wrap), + schemaAddrs.map(ArrowSchema.wrap)) + imported = new ColumnarBatch(vectors.toArray, numRows) + + val charged = CometArrowImportAllocator.getAllocatedMemory - before + val foreign = vectors.head.getValueVector.getBuffers(false).map(_.capacity()).sum + assert( + charged > foreign, + s"expected the import allocator to hold more than the $foreign bytes of imported " + + "buffers, since the importer allocates its own ArrowArray struct there, but it held " + + s"$charged bytes") + } finally { + if (imported != null) { + imported.close() + } + nativeUtil.close() + } + } + test("Variant schema identity round-trips through native Arrow FFI") { val variantType = Utils.variantType.getOrElse { cancel("VariantType requires Spark 4.0+") --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
