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

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


The following commit(s) were added to refs/heads/master by this push:
     new 5750e9d4014 [Subscription] Fix consumer close deadlock (#18704)
5750e9d4014 is described below

commit 5750e9d40144f982a4f772d2349d472f4c82772e
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 17:09:47 2026 +0800

    [Subscription] Fix consumer close deadlock (#18704)
---
 .../consensus/ConsensusLogToTabletConverter.java   |  4 +
 .../consensus/ConsensusPrefetchingQueue.java       | 12 +--
 .../ConsensusLogToTabletConverterTest.java         | 10 +++
 .../consensus/ConsensusPrefetchingQueueTest.java   | 97 ++++++++++++++++++++++
 4 files changed, 118 insertions(+), 5 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverter.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverter.java
index 30767c82baa..abd4c5cc381 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverter.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverter.java
@@ -104,6 +104,10 @@ public class ConsensusLogToTabletConverter {
     return databaseName;
   }
 
+  boolean isTableModel() {
+    return Objects.nonNull(tablePattern);
+  }
+
   static String safeDeviceIdForLog(final InsertNode node) {
     try {
       final Object deviceId = node.getDeviceID();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index d105600cfbc..f68d37cc6a2 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -135,6 +135,10 @@ public class ConsensusPrefetchingQueue {
 
   private final ConsensusLogToTabletConverter converter;
 
+  // Consumer-group meta pushes hold the meta write lock while waiting for 
prefetch shutdown, so
+  // the prefetch path must use this immutable value instead of acquiring the 
meta read lock.
+  private final boolean tableModel;
+
   private final ConsensusSubscriptionCommitManager commitManager;
 
   private SubscriptionMemoryManager subscriptionMemoryManager;
@@ -540,6 +544,7 @@ public class ConsensusPrefetchingQueue {
     this.consensusReqReader = serverImpl.getConsensusReqReader();
     this.retentionPolicy = retentionPolicy;
     this.converter = converter;
+    this.tableModel = converter.isTableModel();
     this.commitManager = commitManager;
     this.subscriptionMemoryManager = 
SubscriptionDataNodeResourceManager.memory();
     this.fallbackCommittedRegionProgress = fallbackCommittedRegionProgress;
@@ -2170,8 +2175,7 @@ public class ConsensusPrefetchingQueue {
             payload,
             commitContext,
             SubscriptionAgent.broker()
-                .getColumnFilterMatcher(
-                    topicName, 
SubscriptionAgent.consumer().isTableModel(consumerGroupId))
+                .getColumnFilterMatcher(topicName, tableModel)
                 .isTimeSelected(),
             getTimeSelectedByTable(converter.getDatabaseName(), tablets));
 
@@ -2198,9 +2202,7 @@ public class ConsensusPrefetchingQueue {
       return Collections.emptyMap();
     }
     final ColumnFilterMatcher matcher =
-        SubscriptionAgent.broker()
-            .getColumnFilterMatcher(
-                topicName, 
SubscriptionAgent.consumer().isTableModel(consumerGroupId));
+        SubscriptionAgent.broker().getColumnFilterMatcher(topicName, 
tableModel);
     final Map<String, Boolean> tableMap = new HashMap<>();
     for (final Tablet tablet : tablets) {
       if (Objects.nonNull(tablet) && Objects.nonNull(tablet.getTableName())) {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverterTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverterTest.java
index 1b8ad968365..b06ccf453fe 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverterTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusLogToTabletConverterTest.java
@@ -55,6 +55,16 @@ public class ConsensusLogToTabletConverterTest {
 
   private static final String DATABASE_NAME = "db";
 
+  @Test
+  public void testDataModelIsDerivedFromImmutableConversionPattern() {
+    final ConsensusLogToTabletConverter tableConverter = 
createConverter("id1");
+    final ConsensusLogToTabletConverter treeConverter =
+        new ConsensusLogToTabletConverter(new IoTDBTreePattern("root.**"), 
null, null, null);
+
+    Assert.assertTrue(tableConverter.isTableModel());
+    Assert.assertFalse(treeConverter.isTableModel());
+  }
+
   @Test
   public void testConvertRelationalInsertRowNodeWithSingleMatchedColumn() {
     final ConsensusLogToTabletConverter converter = createConverter("id1");
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index 3143f8e08a6..b2127159672 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -39,6 +39,8 @@ import 
org.apache.iotdb.db.storageengine.dataregion.wal.io.WALWriter;
 import org.apache.iotdb.db.storageengine.dataregion.wal.node.WALNode;
 import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileStatus;
 import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileUtils;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
+import org.apache.iotdb.db.subscription.agent.SubscriptionConsumerAgent;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import org.apache.iotdb.db.subscription.resource.SubscriptionMemoryManager;
 import org.apache.iotdb.rpc.subscription.config.TopicConstant;
@@ -67,6 +69,9 @@ import java.util.Iterator;
 import java.util.List;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
@@ -107,6 +112,69 @@ public class ConsensusPrefetchingQueueTest {
                 .getModifiers()));
   }
 
+  @Test
+  public void testPrefetchDoesNotDependOnConsumerMetaReadLock() throws 
Exception {
+    final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+    final File systemDir = 
temporaryFolder.newFolder("prefetch-with-consumer-meta-write-lock");
+    final ExecutorService executor = Executors.newSingleThreadExecutor();
+    ConsensusPrefetchingQueue queue = null;
+    boolean consumerMetaWriteLockAcquired = false;
+    try {
+      final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+      reader.currentSearchIndex = 1L;
+      final IoTConsensusServerImpl serverImpl = 
mock(IoTConsensusServerImpl.class);
+      when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+      when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new 
WriterSafeFrontierTracker());
+
+      final ConsensusLogToTabletConverter converter = 
mock(ConsensusLogToTabletConverter.class);
+      
when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet()));
+      when(converter.getDatabaseName()).thenReturn("db");
+      when(converter.isTableModel()).thenReturn(true);
+
+      queue =
+          new ConsensusPrefetchingQueue(
+              "consumerGroup",
+              "topic",
+              TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+              new DataRegionId(1),
+              serverImpl,
+              new SubscriptionWalRetentionPolicy(
+                  "topic",
+                  SubscriptionWalRetentionPolicy.UNBOUNDED,
+                  SubscriptionWalRetentionPolicy.UNBOUNDED),
+              converter,
+              newCommitManager(systemDir),
+              new RegionProgress(Collections.emptyMap()),
+              1L,
+              1L,
+              true);
+      // Consumer-group metadata pushes hold this write lock while closing 
consensus queues and
+      // waiting for an in-flight prefetch round to finish. The prefetch round 
must therefore not
+      // try to reacquire the metadata read lock when it materializes an event.
+      acquireConsumerMetaWriteLock();
+      consumerMetaWriteLockAcquired = true;
+
+      final ConsensusPrefetchingQueue queueToMaterialize = queue;
+      final Future<Boolean> materialized =
+          executor.submit(
+              () ->
+                  invokeCreateAndEnqueueEvent(
+                      queueToMaterialize, 
Collections.singletonList(createTablet())));
+      assertTrue(materialized.get(5, TimeUnit.SECONDS));
+      assertEquals(1, queue.getPrefetchedEventCount());
+    } finally {
+      if (consumerMetaWriteLockAcquired) {
+        releaseConsumerMetaWriteLock();
+      }
+      executor.shutdownNow();
+      executor.awaitTermination(5, TimeUnit.SECONDS);
+      if (queue != null) {
+        queue.close();
+      }
+      
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+    }
+  }
+
   @Test
   public void testWalFileCommitRequirementUsesPerWriterMaximum() {
     final String regionId = "DataRegion[1]";
@@ -2293,6 +2361,35 @@ public class ConsensusPrefetchingQueueTest {
     return new ConsensusSubscriptionCommitManager(progressFetcher);
   }
 
+  private static void acquireConsumerMetaWriteLock() throws Exception {
+    invokeConsumerMetaLockMethod("acquireWriteLock");
+  }
+
+  private static void releaseConsumerMetaWriteLock() throws Exception {
+    invokeConsumerMetaLockMethod("releaseWriteLock");
+  }
+
+  private static void invokeConsumerMetaLockMethod(final String methodName) 
throws Exception {
+    final Method method = 
SubscriptionConsumerAgent.class.getDeclaredMethod(methodName);
+    method.setAccessible(true);
+    method.invoke(SubscriptionAgent.consumer());
+  }
+
+  private static boolean invokeCreateAndEnqueueEvent(
+      final ConsensusPrefetchingQueue queue, final List<Tablet> tablets) 
throws Exception {
+    final Method method =
+        ConsensusPrefetchingQueue.class.getDeclaredMethod(
+            "createAndEnqueueEvent",
+            List.class,
+            long.class,
+            long.class,
+            long.class,
+            long.class,
+            long.class);
+    method.setAccessible(true);
+    return (Boolean) method.invoke(queue, tablets, 1L, 1L, 1L, 0L, 0L);
+  }
+
   private static final class FakeConsensusReqReader implements 
ConsensusReqReader {
 
     private long currentSearchIndex;

Reply via email to