This is an automated email from the ASF dual-hosted git repository.

FrankChen021 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new 29c6d9207e9 Fix record supplier lock release (#19811)
29c6d9207e9 is described below

commit 29c6d9207e979114254ca30503a8f0059a7ec54a
Author: Frank Chen <[email protected]>
AuthorDate: Tue Aug 4 10:24:08 2026 +0800

    Fix record supplier lock release (#19811)
---
 .../supervisor/SeekableStreamSupervisor.java       |  8 +++----
 .../SeekableStreamSupervisorStateTest.java         | 28 ++++++++++++++++++++++
 2 files changed, 32 insertions(+), 4 deletions(-)

diff --git 
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
 
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
index 58f92ddaab7..c6130303ec9 100644
--- 
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
+++ 
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java
@@ -5406,11 +5406,11 @@ public abstract class 
SeekableStreamSupervisor<PartitionIdType, SequenceOffsetTy
     StreamPartition<PartitionIdType> streamPartition = 
StreamPartition.of(ioConfig.getStream(), partition);
     OrderedSequenceNumber<SequenceOffsetType> sequenceNumber = 
makeSequenceNumber(offsetFromMetadata);
     recordSupplierLock.lock();
-    if (!recordSupplier.getAssignment().contains(streamPartition)) {
-      // this shouldn't happen, but in case it does...
-      throw new IllegalStateException("Record supplier does not match current 
known partitions");
-    }
     try {
+      if (!recordSupplier.getAssignment().contains(streamPartition)) {
+        // this shouldn't happen, but in case it does...
+        throw new IllegalStateException("Record supplier does not match 
current known partitions");
+      }
       return recordSupplier.isOffsetAvailable(streamPartition, sequenceNumber);
     }
     finally {
diff --git 
a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java
 
b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java
index 695a916383e..1bdc228b7e1 100644
--- 
a/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java
+++ 
b/indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java
@@ -112,6 +112,8 @@ import org.junit.Test;
 import javax.annotation.Nullable;
 import java.io.File;
 import java.io.IOException;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
 import java.math.BigInteger;
 import java.util.ArrayList;
 import java.util.Collection;
@@ -248,6 +250,32 @@ public class SeekableStreamSupervisorStateTest extends 
EasyMockSupport
     verifyAll();
   }
 
+  @Test
+  public void 
testCheckOffsetAvailabilityReleasesLockWhenAssignmentDoesNotMatch() throws 
Exception
+  {
+    EasyMock.expect(spec.isSuspended()).andReturn(false);
+    EasyMock.reset(recordSupplier);
+    
EasyMock.expect(recordSupplier.getAssignment()).andReturn(ImmutableSet.of());
+    replayAll();
+
+    final TestSeekableStreamSupervisor supervisor = new 
TestSeekableStreamSupervisor();
+    supervisor.recordSupplier = recordSupplier;
+    final Method checkOffsetAvailability = 
SeekableStreamSupervisor.class.getDeclaredMethod(
+        "checkOffsetAvailability",
+        Object.class,
+        Object.class
+    );
+    checkOffsetAvailability.setAccessible(true);
+
+    final InvocationTargetException exception = Assert.assertThrows(
+        InvocationTargetException.class,
+        () -> checkOffsetAvailability.invoke(supervisor, SHARD_ID, "1")
+    );
+    Assert.assertEquals(IllegalStateException.class, 
exception.getCause().getClass());
+    Assert.assertFalse(supervisor.getRecordSupplierLock().isLocked());
+    verifyAll();
+  }
+
   @Test
   public void testRunningStreamGetSequenceNumberReturnsNull()
   {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to