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 76005c94103 Avoid holding ProcedureManager lock while waiting for
table tasks (#18673)
76005c94103 is described below
commit 76005c941037ee2666fd8af1923e9296b524abe3
Author: Caideyipi <[email protected]>
AuthorDate: Mon Sep 21 10:47:13 2026 +0800
Avoid holding ProcedureManager lock while waiting for table tasks (#18673)
---
.../iotdb/confignode/manager/ProcedureManager.java | 23 +--
.../manager/ProcedureManagerTableTaskTest.java | 225 +++++++++++++++++++++
2 files changed, 235 insertions(+), 13 deletions(-)
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index 28b54fa61d5..dfbb440d9d6 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -2254,10 +2254,6 @@ public class ProcedureManager {
return waitingProcedureFinished(procedure);
}
- private TSStatus waitingProcedureFinished(final long procedureId) {
- return waitingProcedureFinished(executor.getProcedures().get(procedureId));
- }
-
protected TSStatus waitingProcedureFinished(final Procedure<?> procedure) {
return waitingProcedureFinished(procedure, PROCEDURE_WAIT_TIME_OUT);
}
@@ -2520,9 +2516,8 @@ public class ProcedureManager {
public TDeleteTableDeviceResp deleteDevices(
final TDeleteTableDeviceReq req, final boolean isGeneratedByPipe) {
- long procedureId;
DeleteDevicesProcedure procedure = null;
- final TSStatus status;
+ final Procedure<?> procedureToWait;
synchronized (this) {
final Pair<Long, Boolean> procedureIdDuplicatePair =
checkDuplicateTableTask(
@@ -2532,7 +2527,7 @@ public class ProcedureManager {
null,
req.queryId,
ProcedureType.DELETE_DEVICES_PROCEDURE);
- procedureId = procedureIdDuplicatePair.getLeft();
+ final long procedureId = procedureIdDuplicatePair.getLeft();
if (procedureId == -1) {
if (Boolean.TRUE.equals(procedureIdDuplicatePair.getRight())) {
@@ -2551,11 +2546,12 @@ public class ProcedureManager {
req.getModInfo(),
isGeneratedByPipe);
this.executor.submitProcedure(procedure);
- status = waitingProcedureFinished(procedure);
+ procedureToWait = procedure;
} else {
- status = waitingProcedureFinished(procedureId);
+ procedureToWait = executor.getProcedures().get(procedureId);
}
}
+ final TSStatus status = waitingProcedureFinished(procedureToWait);
if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return new TDeleteTableDeviceResp(StatusUtils.OK)
.setDeletedNum(
@@ -2618,11 +2614,11 @@ public class ProcedureManager {
final String queryId,
final ProcedureType thisType,
final Procedure<ConfigNodeProcedureEnv> procedure) {
- final long procedureId;
+ final Procedure<?> procedureToWait;
synchronized (this) {
final Pair<Long, Boolean> procedureIdDuplicatePair =
checkDuplicateTableTask(database, table, tableName, newName,
queryId, thisType);
- procedureId = procedureIdDuplicatePair.getLeft();
+ final long procedureId = procedureIdDuplicatePair.getLeft();
if (procedureId == -1) {
if (Boolean.TRUE.equals(procedureIdDuplicatePair.getRight())) {
@@ -2631,11 +2627,12 @@ public class ProcedureManager {
"Some other task is operating table with same name.");
}
this.executor.submitProcedure(procedure);
+ procedureToWait = procedure;
} else {
- return waitingProcedureFinished(procedureId);
+ procedureToWait = executor.getProcedures().get(procedureId);
}
}
- return waitingProcedureFinished(procedure);
+ return waitingProcedureFinished(procedureToWait);
}
public Pair<Long, Boolean> checkDuplicateTableTask(
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTableTaskTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTableTaskTest.java
new file mode 100644
index 00000000000..45b2e7aa230
--- /dev/null
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTableTaskTest.java
@@ -0,0 +1,225 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.confignode.manager;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.schema.table.TsTable;
+import org.apache.iotdb.commons.utils.StatusUtils;
+import org.apache.iotdb.confignode.persistence.ProcedureInfo;
+import org.apache.iotdb.confignode.procedure.Procedure;
+import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
+import
org.apache.iotdb.confignode.procedure.impl.schema.table.CreateTableProcedure;
+import
org.apache.iotdb.confignode.procedure.impl.schema.table.DeleteDevicesProcedure;
+import org.apache.iotdb.confignode.procedure.store.ProcedureType;
+import org.apache.iotdb.confignode.rpc.thrift.TDeleteTableDeviceReq;
+import org.apache.iotdb.confignode.rpc.thrift.TDeleteTableDeviceResp;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Test;
+
+import java.nio.ByteBuffer;
+import java.util.concurrent.ConcurrentHashMap;
+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 static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class ProcedureManagerTableTaskTest {
+
+ @Test
+ public void testWaitingForDuplicateTableTaskDoesNotBlockOtherTableTask()
throws Exception {
+ final ProcedureExecutor<ConfigNodeProcedureEnv> procedureExecutor =
+ mock(ProcedureExecutor.class);
+ final ConcurrentHashMap<Long, Procedure<ConfigNodeProcedureEnv>>
procedures =
+ new ConcurrentHashMap<>();
+ when(procedureExecutor.getProcedures()).thenReturn(procedures);
+
+ final TestProcedureManager procedureManager =
+ new TestProcedureManager(mock(ConfigManager.class),
mock(ProcedureInfo.class));
+ procedureManager.setExecutor(procedureExecutor);
+
+ final String database = "database";
+ final TsTable firstTable = new TsTable("first");
+ final CreateTableProcedure runningProcedure =
+ new CreateTableProcedure(database, firstTable, false);
+ runningProcedure.setProcId(1);
+ procedures.put(runningProcedure.getProcId(), runningProcedure);
+ procedureManager.blockWhenWaitingFor(runningProcedure);
+
+ final CreateTableProcedure duplicateProcedure =
+ new CreateTableProcedure(database, firstTable, false);
+ final TsTable secondTable = new TsTable("second");
+ final CreateTableProcedure independentProcedure =
+ new CreateTableProcedure(database, secondTable, false);
+ final ExecutorService requestExecutor = Executors.newFixedThreadPool(2);
+
+ try {
+ final Future<TSStatus> duplicateRequest =
+ requestExecutor.submit(
+ () ->
+ procedureManager.executeWithoutDuplicate(
+ database,
+ firstTable,
+ firstTable.getTableName(),
+ null,
+ ProcedureType.CREATE_TABLE_PROCEDURE,
+ duplicateProcedure));
+ assertTrue(procedureManager.awaitBlockedWait());
+
+ final Future<TSStatus> independentRequest =
+ requestExecutor.submit(
+ () ->
+ procedureManager.executeWithoutDuplicate(
+ database,
+ secondTable,
+ secondTable.getTableName(),
+ null,
+ ProcedureType.CREATE_TABLE_PROCEDURE,
+ independentProcedure));
+
+ assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ independentRequest.get(5, TimeUnit.SECONDS).getCode());
+ verify(procedureExecutor).submitProcedure(independentProcedure);
+ verify(procedureExecutor, never()).submitProcedure(duplicateProcedure);
+
+ procedureManager.releaseBlockedWait();
+ assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ duplicateRequest.get(5, TimeUnit.SECONDS).getCode());
+ } finally {
+ procedureManager.releaseBlockedWait();
+ requestExecutor.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testWaitingForDuplicateDeleteDevicesDoesNotBlockOtherTableTask()
throws Exception {
+ final ProcedureExecutor<ConfigNodeProcedureEnv> procedureExecutor =
+ mock(ProcedureExecutor.class);
+ final ConcurrentHashMap<Long, Procedure<ConfigNodeProcedureEnv>>
procedures =
+ new ConcurrentHashMap<>();
+ when(procedureExecutor.getProcedures()).thenReturn(procedures);
+
+ final TestProcedureManager procedureManager =
+ new TestProcedureManager(mock(ConfigManager.class),
mock(ProcedureInfo.class));
+ procedureManager.setExecutor(procedureExecutor);
+
+ final String database = "database";
+ final String tableName = "first";
+ final String queryId = "query";
+ final byte[] emptyBytes = new byte[0];
+ final DeleteDevicesProcedure runningProcedure =
+ new DeleteDevicesProcedure(
+ database, tableName, queryId, emptyBytes, emptyBytes, emptyBytes,
false);
+ runningProcedure.setProcId(1);
+ procedures.put(runningProcedure.getProcId(), runningProcedure);
+ procedureManager.blockWhenWaitingFor(runningProcedure);
+
+ final TDeleteTableDeviceReq duplicateRequest =
+ new TDeleteTableDeviceReq(
+ database,
+ tableName,
+ queryId,
+ ByteBuffer.wrap(emptyBytes),
+ ByteBuffer.wrap(emptyBytes),
+ ByteBuffer.wrap(emptyBytes));
+ final TsTable secondTable = new TsTable("second");
+ final CreateTableProcedure independentProcedure =
+ new CreateTableProcedure(database, secondTable, false);
+ final ExecutorService requestExecutor = Executors.newFixedThreadPool(2);
+
+ try {
+ final Future<TDeleteTableDeviceResp> deleteDevicesRequest =
+ requestExecutor.submit(() ->
procedureManager.deleteDevices(duplicateRequest, false));
+ assertTrue(procedureManager.awaitBlockedWait());
+
+ final Future<TSStatus> independentRequest =
+ requestExecutor.submit(
+ () ->
+ procedureManager.executeWithoutDuplicate(
+ database,
+ secondTable,
+ secondTable.getTableName(),
+ null,
+ ProcedureType.CREATE_TABLE_PROCEDURE,
+ independentProcedure));
+
+ assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ independentRequest.get(5, TimeUnit.SECONDS).getCode());
+ verify(procedureExecutor).submitProcedure(independentProcedure);
+
+ procedureManager.releaseBlockedWait();
+ assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ deleteDevicesRequest.get(5, TimeUnit.SECONDS).getStatus().getCode());
+ } finally {
+ procedureManager.releaseBlockedWait();
+ requestExecutor.shutdownNow();
+ }
+ }
+
+ private static class TestProcedureManager extends ProcedureManager {
+
+ private final CountDownLatch waitStarted = new CountDownLatch(1);
+ private final CountDownLatch waitRelease = new CountDownLatch(1);
+ private Procedure<?> blockedProcedure;
+
+ private TestProcedureManager(
+ final ConfigManager configManager, final ProcedureInfo procedureInfo) {
+ super(configManager, procedureInfo);
+ }
+
+ private void blockWhenWaitingFor(final Procedure<?> procedure) {
+ blockedProcedure = procedure;
+ }
+
+ private boolean awaitBlockedWait() throws InterruptedException {
+ return waitStarted.await(5, TimeUnit.SECONDS);
+ }
+
+ private void releaseBlockedWait() {
+ waitRelease.countDown();
+ }
+
+ @Override
+ protected TSStatus waitingProcedureFinished(final Procedure<?> procedure) {
+ if (procedure == blockedProcedure) {
+ waitStarted.countDown();
+ try {
+ waitRelease.await();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ return StatusUtils.OK;
+ }
+ }
+}