This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new cd72eac711 [#13399] fix(core): make SegmentedLock.withGlobalLock
exclusive against in-flight segment ops (#13400)
cd72eac711 is described below
commit cd72eac7116dcafb56a660476e64bddabf4e5477
Author: YangJie <[email protected]>
AuthorDate: Tue Sep 22 08:07:33 2026 -0400
[#13399] fix(core): make SegmentedLock.withGlobalLock exclusive against
in-flight segment ops (#13400)
### What changes were proposed in this pull request?
Replace the entry latch in `SegmentedLock` with a
`ReentrantReadWriteLock` gate. Every `withLock` variant now holds the
read lock across its full critical section (segment-lock acquire plus
the action), and `withGlobalLock` takes the write lock, so the global
action waits for all in-flight segment operations to finish before it
runs. The concurrent-global `IllegalStateException` and `isClearing()`
semantics are preserved.
### Why are the changes needed?
`withGlobalLock` promised exclusive access to all segments but acquired
no segment lock, so a segment operation already inside its critical
section ran concurrently with the global action. In
`CaffeineEntityCache`, a `doPut` racing `clear()` could write its index
entry into the tree `clear()` was swapping out, leaving an entity in
`cacheData` but absent from the active `cacheIndex` and served stale
until TTL.
Fix: #13399
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Added `TestSegmentedLock.testGlobalClearingWaitsForInFlightOperations`,
a deterministic test that pins that a global action does not proceed
until an in-flight segment operation has finished. It fails against the
pre-fix code.
---
.../org/apache/gravitino/cache/SegmentedLock.java | 120 ++++++++++++---------
.../apache/gravitino/cache/TestSegmentedLock.java | 60 +++++++++++
2 files changed, 129 insertions(+), 51 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
b/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
index cfbb8d649e..b6b2d5a8dc 100644
--- a/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
+++ b/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
@@ -21,9 +21,9 @@ package org.apache.gravitino.cache;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.util.concurrent.Striped;
-import java.util.concurrent.CountDownLatch;
-import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
/**
* Segmented lock for improved concurrency. Divides locks into segments to
reduce contention.
@@ -34,8 +34,18 @@ public class SegmentedLock {
private final Striped<Lock> stripedLocks;
- /** CountDownLatch for global operations - null when no operation is in
progress */
- private final AtomicReference<CountDownLatch> globalOperationLatch = new
AtomicReference<>();
+ /**
+ * Gates segment operations against global operations: segment operations
hold the read lock
+ * across their whole critical section, global operations hold the write
lock, so a global
+ * operation excludes every segment operation, including ones already in
flight when it starts.
+ * The lock is fair so that a steady stream of segment operations cannot
indefinitely starve a
+ * waiting global operation; reentrant read reacquisition is still permitted
while a writer is
+ * queued, so the nested cache paths that reacquire the read lock do not
deadlock.
+ */
+ private final ReentrantReadWriteLock globalGate = new
ReentrantReadWriteLock(true);
+
+ /** True while a global operation is in progress, used to reject concurrent
global operations. */
+ private final AtomicBoolean clearing = new AtomicBoolean(false);
/**
* Creates a SegmentedLock with the specified number of segments. Guava's
Striped automatically
@@ -82,14 +92,19 @@ public class SegmentedLock {
* @throws RuntimeException if interrupted
*/
public void withLock(Object key, Runnable action) {
- waitForGlobalComplete();
- Lock lock = getSegmentLock(key);
+ Lock readLock = globalGate.readLock();
try {
- lock.lockInterruptibly();
+ readLock.lockInterruptibly();
try {
- action.run();
+ Lock lock = getSegmentLock(key);
+ lock.lockInterruptibly();
+ try {
+ action.run();
+ } finally {
+ lock.unlock();
+ }
} finally {
- lock.unlock();
+ readLock.unlock();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -108,14 +123,19 @@ public class SegmentedLock {
* @throws RuntimeException if interrupted
*/
public <T> T withLock(Object key, java.util.function.Supplier<T> action) {
- waitForGlobalComplete();
- Lock lock = getSegmentLock(key);
+ Lock readLock = globalGate.readLock();
try {
- lock.lockInterruptibly();
+ readLock.lockInterruptibly();
try {
- return action.get();
+ Lock lock = getSegmentLock(key);
+ lock.lockInterruptibly();
+ try {
+ return action.get();
+ } finally {
+ lock.unlock();
+ }
} finally {
- lock.unlock();
+ readLock.unlock();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -136,14 +156,19 @@ public class SegmentedLock {
*/
public <T, E extends Exception> T withLockAndThrow(
Object key, EntityCache.ThrowingSupplier<T, E> action) throws E {
- waitForGlobalComplete();
- Lock lock = getSegmentLock(key);
+ Lock readLock = globalGate.readLock();
try {
- lock.lockInterruptibly();
+ readLock.lockInterruptibly();
try {
- return action.get();
+ Lock lock = getSegmentLock(key);
+ lock.lockInterruptibly();
+ try {
+ return action.get();
+ } finally {
+ lock.unlock();
+ }
} finally {
- lock.unlock();
+ readLock.unlock();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -162,14 +187,19 @@ public class SegmentedLock {
*/
public <E extends Exception> void withLockAndThrow(
Object key, EntityCache.ThrowingRunnable<E> action) throws E {
- waitForGlobalComplete();
- Lock lock = getSegmentLock(key);
+ Lock readLock = globalGate.readLock();
try {
- lock.lockInterruptibly();
+ readLock.lockInterruptibly();
try {
- action.run();
+ Lock lock = getSegmentLock(key);
+ lock.lockInterruptibly();
+ try {
+ action.run();
+ } finally {
+ lock.unlock();
+ }
} finally {
- lock.unlock();
+ readLock.unlock();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -180,7 +210,7 @@ public class SegmentedLock {
/** Checks if a global operation is currently in progress. */
@VisibleForTesting
public boolean isClearing() {
- return globalOperationLatch.get() != null;
+ return clearing.get();
}
/**
@@ -196,42 +226,30 @@ public class SegmentedLock {
* Executes a global clearing operation with exclusive access to all
segments. This method sets
* the clearing flag and ensures no other operations can proceed until the
clearing is complete.
*
+ * <p>Exclusivity holds against every segment operation, including ones that
were already in
+ * flight when this method is called: segment operations hold the {@code
globalGate} read lock
+ * across their whole critical section, and the write lock acquired here
waits for all of them.
+ *
+ * <p>Must not be called from inside a {@code withLock} action on the same
instance: the read lock
+ * cannot upgrade to the write lock, so such a call would deadlock.
+ *
* @param action The clearing action to execute
*/
public void withGlobalLock(Runnable action) {
- // Create a new CountDownLatch for this operation
- CountDownLatch latch = new CountDownLatch(1);
-
- // Atomically set the latch, fail if another operation is already in
progress
- if (!globalOperationLatch.compareAndSet(null, latch)) {
+ // Mark the global operation in progress, fail if another one is already
running
+ if (!clearing.compareAndSet(false, true)) {
throw new IllegalStateException("Global operation already in progress");
}
try {
- synchronized (this) {
+ globalGate.writeLock().lock();
+ try {
action.run();
+ } finally {
+ globalGate.writeLock().unlock();
}
} finally {
- // Clear state first, then signal completion
- globalOperationLatch.set(null);
- latch.countDown();
- }
- }
-
- /**
- * Waits for any ongoing global operation to complete. This method is called
by regular operations
- * to ensure they don't interfere with global operations.
- */
- private void waitForGlobalComplete() {
- CountDownLatch latch = globalOperationLatch.get();
- if (latch != null) {
- try {
- latch.await();
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- throw new RuntimeException(
- "Thread was interrupted while waiting for global operation to
complete", e);
- }
+ clearing.set(false);
}
}
}
diff --git
a/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
b/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
index e377c7751f..edce98096a 100644
--- a/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
+++ b/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
@@ -29,6 +29,7 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
@@ -380,4 +381,63 @@ public class TestSegmentedLock {
});
});
}
+
+ @Test
+ @Timeout(30)
+ void testGlobalClearingWaitsForInFlightOperations() throws
InterruptedException {
+ SegmentedLock lock = new SegmentedLock(4);
+ CountDownLatch insideSegmentOp = new CountDownLatch(1);
+ CountDownLatch releaseSegmentOp = new CountDownLatch(1);
+ CountDownLatch globalActionDone = new CountDownLatch(1);
+ AtomicBoolean overlapped = new AtomicBoolean(false);
+ AtomicBoolean segmentOpActive = new AtomicBoolean(false);
+
+ Thread segmentOpThread =
+ new Thread(
+ () ->
+ lock.withLock(
+ "key1",
+ () -> {
+ segmentOpActive.set(true);
+ insideSegmentOp.countDown();
+ try {
+ releaseSegmentOp.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ segmentOpActive.set(false);
+ }));
+ segmentOpThread.start();
+ assertTrue(insideSegmentOp.await(5, TimeUnit.SECONDS), "segment operation
never started");
+
+ Thread globalThread =
+ new Thread(
+ () ->
+ lock.withGlobalLock(
+ () -> {
+ // The global action must never overlap the segment
operation's critical
+ // section. segmentOpActive stays true from the first to
the last statement of
+ // the segment action, and the write lock is granted
only after the read lock
+ // is released, so a correct gate always observes it
false here.
+ if (segmentOpActive.get()) {
+ overlapped.set(true);
+ }
+ globalActionDone.countDown();
+ }));
+ globalThread.start();
+
+ // While the segment operation is in flight, the global action must not
complete:
+ // withGlobalLock promises exclusive access to all segments, so it must
block
+ // until the in-flight operation releases its critical section. A
completion
+ // here means the global action ran concurrently with it.
+ assertFalse(
+ globalActionDone.await(1, TimeUnit.SECONDS),
+ "global action ran concurrently with an in-flight segment operation");
+ releaseSegmentOp.countDown();
+ assertTrue(globalActionDone.await(5, TimeUnit.SECONDS), "global action
never completed");
+ assertFalse(overlapped.get(), "global action overlapped with the in-flight
segment operation");
+
+ segmentOpThread.join(5000);
+ globalThread.join(5000);
+ }
}