comphead commented on code in PR #6813:
URL: https://github.com/apache/datafusion-comet/pull/6813#discussion_r4232154086
##########
spark/src/test/scala/org/apache/spark/CometTaskMemoryManagerSuite.scala:
##########
@@ -134,52 +134,138 @@ class CometTaskMemoryManagerSuite extends SparkFunSuite {
}
}
+ test(
+ "an acquire waiting in Spark survives another native plan releasing the
task's last bytes") {
+ checkWaitingAcquireSurvives { _ =>
+ val otherPlan = new CometTaskMemoryManager(2L, 0L)
+ (otherPlan.acquireMemory(_), otherPlan.releaseMemory(_))
+ }
+ }
+
+ test("an acquire waiting in Spark survives a JVM consumer releasing the
task's last bytes") {
Review Comment:
These two race tests reach Spark through the same call. From reading Spark
3.4.3 and 4.1.3, `CometTaskMemoryManager.releaseMemory` and
`MemoryConsumer.freeMemory` both end in
`TaskMemoryManager.releaseExecutionMemory` against the task id, so I expect the
releaser's type not to change what Spark or `acquireMemory` does. Could we keep
one of the two race tests and the refusal test?
##########
docs/source/contributor-guide/memory_management.md:
##########
@@ -296,6 +297,12 @@ JNI, which goes through Spark's ordinary
`TaskMemoryManager`. That means:
granted. While overcommit is outstanding, `try_grow` asks Spark for it on
top of the request and
is refused unless Spark can cover both, so operators spill until the debt is
repaid. Both pools'
`Display` output and their `try_grow` errors report the current overcommit.
+- A request that waits in Spark can lose the task's entry in Spark's execution
pool when another
Review Comment:
A note for whichever of this PR and #6310 lands second. `git merge-tree`
reports conflicts only in `CometTaskMemoryManager.java` and
`CometTaskMemoryManagerSuite.scala`. `memory_management.md` merges cleanly and
keeps both #6310's "A parked acquire can wake up to a missing task entry"
paragraph and this bullet, which describe the same retry. That paragraph also
says the `allocatePage` gap is tracked in #6304, which #6403 has since closed.
Could the rebase keep only one of the two?
##########
spark/src/test/scala/org/apache/spark/CometTaskMemoryManagerSuite.scala:
##########
@@ -134,52 +134,138 @@ class CometTaskMemoryManagerSuite extends SparkFunSuite {
}
}
+ test(
+ "an acquire waiting in Spark survives another native plan releasing the
task's last bytes") {
+ checkWaitingAcquireSurvives { _ =>
+ val otherPlan = new CometTaskMemoryManager(2L, 0L)
+ (otherPlan.acquireMemory(_), otherPlan.releaseMemory(_))
+ }
+ }
+
+ test("an acquire waiting in Spark survives a JVM consumer releasing the
task's last bytes") {
+ checkWaitingAcquireSurvives { taskMemoryManager =>
+ val consumer = new OffHeapConsumer(taskMemoryManager)
+ (consumer.acquireMemory(_), consumer.freeMemory(_))
+ }
+ }
+
+ test("an acquire rethrows a NoSuchElementException other than Spark's
missing task entry") {
+ val otherMessage = noSuchElement("no entry for task 0",
fromExecutionMemoryPool = true)
+ val notFromPool = noSuchElement("key not found: 0",
fromExecutionMemoryPool = false)
+ for (error <- Seq(otherMessage, notFromPool)) {
+ val taskMemoryManager = new FailingTaskMemoryManager(() => error,
failures = Int.MaxValue)
+ withTaskContext(taskMemoryManager) { _ =>
+ val manager = new CometTaskMemoryManager(1L, 0L)
+ val events = logEvents(Level.INFO) {
+ val thrown =
intercept[NoSuchElementException](manager.acquireMemory(20L))
+ assert(thrown eq error)
+ }
+ assert(events.isEmpty, messages(events))
+ assert(taskMemoryManager.calls.get == 1)
+ assert(manager.getUsed == 0L)
+ }
+ }
+ }
+
+ test("an acquire that keeps losing its task entry is refused") {
+ // A new exception on every call, so the test can tell which one is logged.
+ val taskMemoryManager = new FailingTaskMemoryManager(
+ () => noSuchElement("key not found: 0", fromExecutionMemoryPool = true),
+ failures = Int.MaxValue)
+ withTaskContext(taskMemoryManager) { _ =>
+ val manager = new CometTaskMemoryManager(1L, 0L)
+ val consumer = nativeMemoryConsumer(manager)
+ val events = logEvents(Level.INFO) {
+ // A refusal rather than an exception, so native spills as for any
refused reservation.
+ assert(manager.acquireMemory(20L) == 0L)
+ }
+ assert(events.map(_.getLevel) == Seq(Level.INFO, Level.INFO,
Level.WARN), messages(events))
+ assert(events.last.getThrown eq taskMemoryManager.errors.last)
+ assert(taskMemoryManager.calls.get == MaxAcquireAttempts)
+ assert(taskMemoryManager.errors.distinct.size == MaxAcquireAttempts)
+ assert(manager.getUsed == 0L)
+ assert(consumer.getUsed == 0L)
+ assert(taskMemoryManager.getMemoryConsumptionForThisTask == 0L)
+ }
+ }
+
+ test("an acquire that loses its task entry twice is granted on the last
attempt") {
+ val error = noSuchElement("key not found: 0", fromExecutionMemoryPool =
true)
+ val taskMemoryManager =
+ new FailingTaskMemoryManager(() => error, failures = MaxAcquireAttempts
- 1)
+ withTaskContext(taskMemoryManager) { _ =>
+ val manager = new CometTaskMemoryManager(1L, 0L)
+ val events = logEvents(Level.INFO) {
+ assert(manager.acquireMemory(20L) == 20L)
+ }
+ assert(events.map(_.getLevel) == Seq(Level.INFO, Level.INFO),
messages(events))
+ assert(taskMemoryManager.calls.get == MaxAcquireAttempts)
+ assert(manager.getUsed == 20L)
+ assert(taskMemoryManager.getMemoryConsumptionForThisTask == 20L)
+ manager.releaseMemory(20L)
+ assert(taskMemoryManager.getMemoryConsumptionForThisTask == 0L)
+ }
+ }
+
+ /**
+ * The task holds 10 bytes of a 100 byte pool through another consumer, made
by `holder`, and
+ * another task holds 90. A native plan asks for 20 bytes on another thread
and waits in Spark
+ * below the task's minimum share of 25 while the holder releases the task's
last 10 bytes,
+ * which removes the task's entry from Spark's pool. Once the other task
frees its memory the
+ * native request must be granted in full.
+ */
+ private def checkWaitingAcquireSurvives(
Review Comment:
`checkWaitingAcquireSurvives` is close to a copy of
`checkWaitingAllocationSurvives` in `CometUnifiedShuffleMemoryAllocatorSuite`.
The pool setup, the other task's 90 bytes, the waiting thread,
`awaitWaitingInSpark`, the cleanup in `finally` and the single INFO event check
are the same. Only the holder, the request and the final `used` checks differ.
Since this PR adds `TaskMemoryTestUtils` to share these fixtures, could the
harness move there too, taking the holder's acquire and release and the request
as functions?
--
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]