Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@
* A memory consumer of {@link TaskMemoryManager} that supports spilling.
*
* Note: this only supports allocation / spilling of Tungsten memory.
*
* {@link TaskMemoryManager} tracks consumers by identity, so distinct consumer instances are
* always tracked, spilled and reported separately, regardless of how they implement
* {@code equals} and {@code hashCode}. Code that keeps consumers in a collection should likewise
* key them by identity rather than relying on {@code equals}/{@code hashCode}.
*/
public abstract class MemoryConsumer {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,10 +113,12 @@ public class TaskMemoryManager {
final MemoryMode tungstenMemoryMode;

/**
* Tracks spillable memory consumers.
* Tracks spillable memory consumers. Consumers are tracked by identity rather than by

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, added a paragraph to the MemoryConsumer class Javadoc.

* {@code equals}, so that distinct consumers which happen to be equal (for example, Scala case
* classes) are all offered for spilling and reported in the memory usage breakdown.
*/
@GuardedBy("this")
private final HashSet<MemoryConsumer> consumers;
private final Set<MemoryConsumer> consumers;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.


/**
* The amount of memory that is acquired but not used.
Expand Down Expand Up @@ -181,7 +183,9 @@ public TaskMemoryManager(MemoryManager memoryManager, long taskAttemptId) {
this.memoryManager = memoryManager;
this.tungstenMemoryAllocator = tungstenMemoryAllocator;
this.taskAttemptId = taskAttemptId;
this.consumers = new HashSet<>();
// A task rarely has more than a few consumers; the default expected size of IdentityHashMap
// would eagerly allocate a much larger table for every task.
this.consumers = Collections.newSetFromMap(new IdentityHashMap<>(4));
}

/**
Expand Down
117 changes: 117 additions & 0 deletions core/src/test/java/org/apache/spark/memory/TaskMemoryManagerSuite.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.spark.memory;

import java.io.IOException;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;

import org.junit.jupiter.api.Assertions;
Expand Down Expand Up @@ -187,6 +188,33 @@ public long spill(long size, MemoryConsumer trigger) {
}
}

/**
* A consumer with value-based equality and string representation: two instances on the same
* manager and memory mode are equal, share a hash code and print the same, even though they
* track memory independently.
*/
private static final class ValueEqualConsumer extends TestMemoryConsumer {
ValueEqualConsumer(TaskMemoryManager memoryManager) {
super(memoryManager);
}

@Override
public boolean equals(Object other) {
return other instanceof ValueEqualConsumer that &&
taskMemoryManager == that.taskMemoryManager && getMode() == that.getMode();
}

@Override
public int hashCode() {
return Objects.hash(taskMemoryManager, getMode());
}

@Override
public String toString() {
return "ValueEqualConsumer";
}
}

@Test
public void leakedPageMemoryIsDetected() {
final TaskMemoryManager manager = new TaskMemoryManager(
Expand Down Expand Up @@ -861,6 +889,95 @@ public void shouldNotForceSpillingInDifferentModes() {
Assertions.assertEquals(80, c1.getUsed()); // not spilled
}

@Test
public void equalButDistinctConsumersAreAllSpilled() {
final TestMemoryManager memoryManager = new TestMemoryManager(new SparkConf());
memoryManager.limit(100);
final TaskMemoryManager manager = new TaskMemoryManager(memoryManager, 0);

ValueEqualConsumer c1 = new ValueEqualConsumer(manager);
ValueEqualConsumer c2 = new ValueEqualConsumer(manager);
Assertions.assertEquals(c1, c2);
TestMemoryConsumer c3 = new TestMemoryConsumer(manager);
c1.use(50);
c2.use(50);

// Both equal consumers must be offered for spilling to satisfy the request.
c3.use(100);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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. spillConsumersForPageAllocation builds its own candidate map from consumers.
  • Self-spill, where the requester is the equal-but-distinct consumer. It was never in the set, so it never got the key 0 and was never asked to spill itself. For example, c1 acquires first and then frees its memory, c2.use(100) fills the limit, and then c2.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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Assertions.assertEquals(0, c1.getUsed());
Assertions.assertEquals(0, c2.getUsed());
Assertions.assertEquals(100, c3.getUsed());

c3.free(100);
Assertions.assertEquals(0, manager.cleanUpAllAllocatedMemory());
}

@Test
public void equalButDistinctConsumerCanSelfSpill() {
final TestMemoryManager memoryManager = new TestMemoryManager(new SparkConf());
memoryManager.limit(100);
final TaskMemoryManager manager = new TaskMemoryManager(memoryManager, 0);

ValueEqualConsumer c1 = new ValueEqualConsumer(manager);
ValueEqualConsumer c2 = new ValueEqualConsumer(manager);
c1.use(50);
c1.free(50);
c2.use(100);

// The requesting consumer is spilled last, but it must still be spilled when no other
// consumer holds memory, even if an equal consumer was registered first.
c2.use(50);
Assertions.assertEquals(0, c1.getUsed());
Assertions.assertEquals(50, c2.getUsed());

c2.free(50);
Assertions.assertEquals(0, manager.cleanUpAllAllocatedMemory());
}

@Test
public void equalButDistinctConsumersAreAllSpilledForPageAllocation() {
final TestMemoryManager memoryManager = new TestMemoryManager(new SparkConf());
memoryManager.limit(5120);
final TestAllocator allocator = new TestAllocator(1);
final TaskMemoryManager manager = new TaskMemoryManager(memoryManager, 0, allocator);
ValueEqualConsumer c1 = new ValueEqualConsumer(manager);
ValueEqualConsumer c2 = new ValueEqualConsumer(manager);
final PageAllocatingConsumer requestingConsumer = new PageAllocatingConsumer(manager, 4096);
c1.use(512);
c2.use(512);

// The grant succeeds without spilling, so only the allocator-failure recovery path spills.
final MemoryBlock page = requestingConsumer.allocate(4096);
Assertions.assertNotNull(page);
Assertions.assertEquals(0, c1.getUsed());
Assertions.assertEquals(0, c2.getUsed());

requestingConsumer.freeAllocatedPage(page);
Assertions.assertEquals(0, manager.cleanUpAllAllocatedMemory());
}

@Test
public void equalButDistinctConsumersAreAllInMemoryConsumptionBreakdown() {
final TestMemoryManager memoryManager = new TestMemoryManager(new SparkConf());
memoryManager.limit(100);
final TaskMemoryManager manager = new TaskMemoryManager(memoryManager, 0);

ValueEqualConsumer c1 = new ValueEqualConsumer(manager);
ValueEqualConsumer c2 = new ValueEqualConsumer(manager);
c1.use(20);
c2.use(30);

String breakdown = manager.getMemoryConsumptionBreakdown();
Assertions.assertTrue(breakdown.contains("ValueEqualConsumer: 20.0 B"), breakdown);
Assertions.assertTrue(breakdown.contains("ValueEqualConsumer: 30.0 B"), breakdown);
Assertions.assertFalse(
breakdown.contains("(not attributed to a specific consumer)"), breakdown);

c1.free(20);
c2.free(30);
Assertions.assertEquals(0, manager.cleanUpAllAllocatedMemory());
}

@Test
public void memoryConsumptionBreakdownIsEmptyWhenThereIsNoMemoryToReport() {
final TestMemoryManager memoryManager = new TestMemoryManager(new SparkConf());
Expand Down