andygrove opened a new pull request, #6269:
URL: https://github.com/apache/datafusion-comet/pull/6269
## Which issue does this PR close?
Closes #6257.
## Rationale for this change
`CometTaskMemoryManager.acquireMemory` logs a warning every time Spark
grants less memory than a native pool asked for. It then calls
`TaskMemoryManager.showMemoryUsage()`, which logs at least three more lines at
INFO. A partial grant is not an error, though. `try_grow` hands it back and
refuses the reservation, which is how a native operator knows to spill. `grow`
carries the shortfall as overcommit. So a query that spills logs a warning and
a memory dump for every reservation Spark refuses.
The memory sweep suite behind the issue runs five spilling or failing native
plans with 96 MB of off-heap memory at `local[4]`. Together they logged 257 of
these warnings and 1,028 dump lines in a 35-second run, about a third of the
log. The issue attributed all 257 to the aggregate, but they came from all five
tests. The aggregate alone logged 73 in about 3 seconds.
Nothing is lost when a reservation really fails. The error the pool returns
already says how much Spark granted, and `TrackConsumersPool` lists the top
consumers. Here is the one from that aggregate:
```
CometNativeException: Additional allocation failed for
FinalHashAggregateStream[0] with top memory consumers (across reservations) as:
FinalHashAggregateStream[0]#29(can spill: true) consumed 23.1 MB, peak
23.1 MB.
Error: Failed to acquire 1828080 bytes plus 0 bytes overcommitted, only got
917424 bytes. Reserved: 24248400 bytes
```
The dump is also unsafe, which is why this PR removes it rather than moving
it to DEBUG. sunchao found this cycle while reviewing #5613:
1. One native thread gets a short grant. The grant is still charged to the
task until the thread returns to native code and hands it back.
2. Another acquire of the same task waits inside Spark's
`ExecutionMemoryPool`, below its 1/2N share. It waits in `lock.wait()`, which
releases the memory manager's monitor, but it still holds the
`TaskMemoryManager` monitor, which `acquireExecutionMemory` takes around its
whole body.
3. The first thread calls `showMemoryUsage()`, which takes `synchronized
(this)` on the same `TaskMemoryManager`. It blocks, so it never hands back the
bytes the waiting acquire needs.
The task then hangs until some other task frees memory. On main the default
`fair_unified` pool holds its lock across the call into Spark, which prevents
this between two native threads. `greedy_unified` takes no lock, so it can
reach the cycle today. By my reading of the code, `fair_unified` could reach it
only when the waiting acquire comes from a JVM off-heap consumer of the same
task. #5613 drops the call for the same reason.
Spark takes the same approach with its own partial grants.
`TaskMemoryManager` logs them at DEBUG, and the only caller of
`showMemoryUsage()` is `MemoryConsumer.throwOom`, just before it throws.
## What changes are included in this PR?
- `acquireMemory` logs a partial grant at DEBUG and no longer calls
`showMemoryUsage()`. The DEBUG line keeps this manager's total and
`getMemoryConsumptionForThisTask()`. That call takes only the memory manager's
monitor, which a waiting acquire gives up. The `isDebugEnabled()` check means
it isn't called at all unless DEBUG is on.
- A comment in `acquireMemory` says why the method must not take the
`TaskMemoryManager` monitor.
- The debugging guide explains why these lines are at DEBUG and how to turn
them on.
The issue also offered logging once per task at INFO. I went with DEBUG
because a `CometTaskMemoryManager` exists per native plan rather than per task,
so once per task would need state keyed by task. Spill metrics already show
when an operator spilled, and `spark.comet.debug.memory` logs every refused
`try_grow`.
## How are these changes tested?
- Two new tests in `CometTaskMemoryManagerSuite`:
- Partial and zero grants log nothing at INFO or above, from either
`CometTaskMemoryManager` or Spark's `TaskMemoryManager`. On main this fails
with two warnings and eight dump lines.
- At DEBUG, a partial grant logs the request and the grant, and there is
still no memory dump.
- Each of these changes fails at least one of the tests: going back to
main's code, keeping the dump behind DEBUG, or removing the DEBUG line.
- The suite now extends `SparkFunSuite` so that it can use
`withLogAppender`. That helper leaves a config behind for each logger that had
none. The config copies the root config's additivity, which is off in Comet's
test `log4j2.properties`, so the logger's events would miss `unit-tests.log`
for the rest of the JVM. The test removes the configs it caused. A probe suite
run after it in the same JVM confirmed that both loggers reach the file again.
- A throwaway probe, not committed, ran the interleaving above against a
real `UnifiedMemoryManager`, adapted from #5613's regression test to main by
using a 1-byte acquire in place of that PR's anchor:
| `acquireMemory` | at INFO | at DEBUG |
| --- | --- | --- |
| main (WARN and dump) | hangs in `showMemoryUsage` | hangs |
| dump kept behind DEBUG | completes | hangs |
| this PR | completes | completes |
#5613 carries a committed version of that test, built on its anchor. I
left it out here so the two PRs don't add duplicate helpers to the same suite.
- The sweep's `MemSweepPressureSuite`, which is not committed, logs 257
warnings and 1,028 dump lines before this change and none after. Its aggregate
fails the same way in both runs, which is #6254.
- The suite passes on Spark 4.1 and on Spark 3.4 with Scala 2.12. scalafix
(on 3.5), scalastyle, spotless and prettier are clean.
--
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]