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;