Repository navigation
Conversation
…anager TaskMemoryManager tracked spillable consumers in a HashSet, so a consumer with value-based equals/hashCode (for example a Scala case class such as HybridRowQueue) was deduplicated against an equal but distinct consumer. The dropped instance was never offered for spilling, and its memory was reported as not attributed to any consumer in the OOM breakdown. Track consumers in an identity-based set instead, so every consumer instance is considered regardless of how it defines equality. Co-authored-by: Claude Code <noreply@anthropic.com>
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for making this a general fix, @viirya. Switching consumers to an identity set looks correct to me. consumers is only used for add, iteration and clear, and the code that compares consumers already uses reference equality: c == requestingConsumer in both spill loops, request.consumer == consumer in allocatePage, and trigger == this in HybridQueue.spill. I also traced both new tests against the old HashSet, and both fail there as described. I left 6 inline comments, mostly about test coverage. In summary:
-
TaskMemoryManagerSuite.javaL889: The breakdown test relies on distincttoStringnamesEqual consumers whose
toStringis also value-based (HybridRowQueueon the current master) are still listed as indistinguishable lines. -
TaskMemoryManager.javaL121: Equal consumers are now retained until the task endsNothing removes an entry from
consumersbeforecleanUpAllAllocatedMemory(), so every equal instance stays reachable and is scanned on every spill. -
TaskMemoryManagerSuite.javaL908: Two affected paths are not testedThe page allocation recovery path (
spillConsumersForPageAllocation) and self-spill of a consumer that used to be deduplicated. -
TaskMemoryManagerSuite.javaL869: nit: Placement of the helper classThe other helper classes are grouped at the top of the suite.
-
TaskMemoryManager.javaL116: The identity contract is only documented on a private fieldStating it on
MemoryConsumerwould keep a future hash collection keyed by consumers from bringing the bug back. -
TaskMemoryManager.javaL186: nit:IdentityHashMapallocates its table eagerlyThe no-arg constructor allocates a 64-slot table for every task.
| } | ||
|
|
||
| @Override | ||
| public String toString() { |
There was a problem hiding this comment.
ValueEqualConsumer returns a distinct name from toString, so this test passes even though the breakdown cannot tell equal consumers apart when their toString is also value-based. That is the case for HybridRowQueue on the current master: two equal queues are both listed with the same case-class toString, e.g. HybridRowQueue(org.apache.spark.memory.TaskMemoryManager@...,<dir>,1,...,false): 30.0 B. #59282 adds an identity suffix to HybridRowQueue.toString, but only for that class.
Could we make toString identical for both instances here and assert that both amounts are listed (: 20.0 B and : 30.0 B)? Alternatively, renderConsumerBreakdown and logMemoryUsage could append System.identityHashCode to each consumer, which keeps every line distinguishable regardless of how a consumer implements toString.
There was a problem hiding this comment.
Good point. ValueEqualConsumer now prints the same for both instances, and the test asserts that both ValueEqualConsumer: 20.0 B and ValueEqualConsumer: 30.0 B are listed. I'd rather not change the breakdown format for every consumer in this PR. Most consumers already print a distinct Object.toString, and #59282 covers HybridRowQueue. If we want an identity suffix in general, I can open a separate JIRA.
| */ | ||
| @GuardedBy("this") | ||
| private final HashSet<MemoryConsumer> consumers; | ||
| private final Set<MemoryConsumer> consumers; |
There was a problem hiding this comment.
Nothing removes an entry from consumers before cleanUpAllAllocatedMemory(). So every equal-but-distinct instance that the old HashSet collapsed into one now stays reachable until the task ends, together with everything it references, and both spill loops iterate over it on every spill.
This already holds for any other distinct consumer, so it is not a new kind of problem, and a task usually has only a few consumers. Still, it may be worth a separate JIRA to drop a consumer from the set once its usage reaches zero (or when it is closed), instead of only ever adding to it.
There was a problem hiding this comment.
Agreed. This is the same lifetime that any other distinct consumer already has, so I'll keep this PR focused. Filed SPARK-60134 to drop a consumer from the set once it no longer holds memory.
| c2.use(50); | ||
|
|
||
| // Both equal consumers must be offered for spilling to satisfy the request. | ||
| c3.use(100); |
There was a problem hiding this comment.
These tests cover another consumer spilling through acquireExecutionMemory. Two other paths also missed the deduplicated consumer before this change, and they are not covered:
- The page allocation recovery path.
spillConsumersForPageAllocationbuilds its own candidate map fromconsumers. - Self-spill, where the requester is the equal-but-distinct consumer. It was never in the set, so it never got the key
0and was never asked to spill itself. For example,c1acquires first and then frees its memory,c2.use(100)fills the limit, and thenc2.use(50)gets nothing without this fix.
A small case for each would keep a later change to either loop from bringing the bug back.
There was a problem hiding this comment.
Added equalButDistinctConsumerCanSelfSpill (your scenario) and equalButDistinctConsumersAreAllSpilledForPageAllocation. In the second test the grant succeeds without spilling, so only spillConsumersForPageAllocation spills. Both fail with the old HashSet.
| * A consumer with value-based equality: two instances on the same manager and memory mode are | ||
| * equal (and share a hash code) even though they track memory independently. | ||
| */ | ||
| private static final class ValueEqualConsumer extends TestMemoryConsumer { |
There was a problem hiding this comment.
nit: The other helper classes in this suite (TestAllocator through NonSpillingAllocatingConsumer) are grouped at the top of the class. Could we move ValueEqualConsumer next to them?
There was a problem hiding this comment.
Done, moved next to the other helpers.
|
|
||
| /** | ||
| * Tracks spillable memory consumers. | ||
| * Tracks spillable memory consumers. Consumers are tracked by identity rather than by |
There was a problem hiding this comment.
This Javadoc is the only place that says consumers are tracked by identity, and it is on a private field. Since the goal is to cover all current and future consumers, it may be worth stating the same in the MemoryConsumer class Javadoc as well, e.g. that consumers are tracked by identity and equals/hashCode must not be relied on for that. Then someone who later adds a hash-based collection keyed by MemoryConsumer (for example, for per-consumer statistics) would see the requirement where they look.
There was a problem hiding this comment.
Done, added a paragraph to the MemoryConsumer class Javadoc.
| this.tungstenMemoryAllocator = tungstenMemoryAllocator; | ||
| this.taskAttemptId = taskAttemptId; | ||
| this.consumers = new HashSet<>(); | ||
| this.consumers = Collections.newSetFromMap(new IdentityHashMap<>()); |
There was a problem hiding this comment.
nit: The no-arg IdentityHashMap constructor eagerly allocates a 64-slot table (for the default expected maximum size of 21) for every TaskMemoryManager, i.e. once per task, while the previous HashSet allocated its table lazily on the first add. Since a task rarely has more than a few consumers, a smaller expected size such as new IdentityHashMap<>(4) would avoid most of that. It is minor either way.
There was a problem hiding this comment.
Done, it now uses new IdentityHashMap<>(4).
- Make ValueEqualConsumer print the same for equal instances and assert both amounts are listed in the breakdown; move it next to the other helpers. - Add tests for self-spill and the page allocation recovery path. - Document identity tracking in the MemoryConsumer Javadoc. - Size the IdentityHashMap for a few consumers instead of the default. Co-authored-by: Claude Code <noreply@anthropic.com>
What changes were proposed in this pull request?
TaskMemoryManagertracks spillable consumers in aHashSet<MemoryConsumer>, which deduplicates byequals/hashCode. This PR switches it to an identity-based set (Collections.newSetFromMap(new IdentityHashMap<>())), so every consumer instance is tracked regardless of how it defines equality.The set is only used for
add, iteration andclear(nocontains/remove), and the spill heuristic already compares the requesting consumer by reference (c == requestingConsumer), so the switch does not affect any other logic. This is the only hash-based collection keyed byMemoryConsumerin the codebase.Why are the changes needed?
With value-based equality, an equal but distinct consumer is silently dropped from the set. That instance:
UNABLE_TO_ACQUIRE_MEMORYerrors, where its memory shows up as "not attributed to a specific consumer".This affects master today:
HybridRowQueueis a Scala case class extendingMemoryConsumer, so two queues created with the same arguments in the same task are deduplicated. #59282 fixes this forHybridRowQueuespecifically; as suggested by @dongjoon-hyun in that review, this PR fixes it in one place for all current and future consumers. The two PRs are independent and can be merged in either order.Does this PR introduce any user-facing change?
Yes, as a bug fix: distinct memory consumers that compare equal (currently
HybridRowQueueinstances with identical arguments in the same task) are now all eligible for spilling and are all listed in the OOM memory breakdown. There is no API or configuration change.How was this patch tested?
Added two tests to
TaskMemoryManagerSuiteusing a test consumer with value-basedequals/hashCode:equalButDistinctConsumersAreAllSpilled: before this change, the second equal consumer was not spilled (50 B left).equalButDistinctConsumersAreAllInMemoryConsumptionBreakdown: before this change, the second consumer's memory was reported as not attributed to a specific consumer.Both fail without the fix and pass with it. Also ran the
org.apache.spark.memorysuites,UnsafeExternalSorterSuite,ShuffleExternalSorterSuite,BytesToBytesMap{On,Off}HeapSuite,ExternalSorterSpillSuiteandExternalAppendOnlyMapSuite, pluscore/checkstyleandcore/Test/checkstyle.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5.5)
This pull request and its description were written by Isaac.