From 4442bdffb3a468717140c5c16d97ffc97b644c03 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 30 Sep 2026 16:06:13 +0800 Subject: [PATCH 1/2] Fix table delete routing to region leaders --- .../it/db/it/IoTDBDeletionTableClusterIT.java | 152 ++++++++++++++++++ .../partition/GetRegionGroupsByTimeResp.java | 25 +++ .../manager/partition/PartitionManager.java | 3 +- 3 files changed, 179 insertions(+), 1 deletion(-) create mode 100644 integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java diff --git a/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java new file mode 100644 index 0000000000000..34804aa6f2ec0 --- /dev/null +++ b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java @@ -0,0 +1,152 @@ +/* + * 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.relational.it.db.it; + +import org.apache.iotdb.consensus.ConsensusFactory; +import org.apache.iotdb.isession.SessionConfig; +import org.apache.iotdb.it.env.EnvFactory; +import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper; +import org.apache.iotdb.it.framework.IoTDBTestRunner; +import org.apache.iotdb.itbase.category.TableClusterIT; + +import org.awaitility.Awaitility; +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.Test; +import org.junit.experimental.categories.Category; +import org.junit.runner.RunWith; + +import java.sql.Connection; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.util.concurrent.TimeUnit; + +import static org.apache.iotdb.itbase.env.BaseEnv.TABLE_SQL_DIALECT; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +@RunWith(IoTDBTestRunner.class) +@Category(TableClusterIT.class) +public class IoTDBDeletionTableClusterIT { + + private static final String DATABASE = "sc3"; + + @BeforeClass + public static void setUpClass() throws SQLException { + EnvFactory.getEnv() + .getConfig() + .getCommonConfig() + .setConfigNodeConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS) + .setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS) + .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS) + .setSchemaReplicationFactor(2) + .setDataReplicationFactor(2); + EnvFactory.getEnv().initClusterEnvironment(3, 3); + try (Connection connection = EnvFactory.getEnv().getConnection(TABLE_SQL_DIALECT); + Statement statement = connection.createStatement()) { + statement.execute("CREATE DATABASE " + DATABASE); + } + } + + @AfterClass + public static void tearDownClass() { + EnvFactory.getEnv().cleanClusterEnvironment(); + } + + /** + * deleting time=2 from three flushed devices must leave exactly (1,d1,1) and (3,d3,3), + * through every DataNode coordinator, both before and after another flush. + */ + @Test + public void testDeleteByTimeAfterFlush() throws SQLException { + checkDeletionOnEveryDataNode("time_delete_", " WHERE time = 2", 1, 3); + } + + /** + * an unconditional DELETE must remove all three flushed rows through every DataNode + * coordinator, and another flush must not make the deleted data visible again. + */ + @Test + public void testDeleteAllAfterFlush() throws SQLException { + checkDeletionOnEveryDataNode("full_delete_", ""); + } + + private void checkDeletionOnEveryDataNode( + String tablePrefix, String predicate, int... remainingTimes) throws SQLException { + for (int i = 0; i < EnvFactory.getEnv().getDataNodeWrapperList().size(); i++) { + String table = DATABASE + "." + tablePrefix + i; + // Pin each DELETE to a different coordinator and use fresh data for every attempt. + try (Connection connection = + EnvFactory.getEnv() + .getWriteOnlyConnectionWithSpecifiedDataNode( + EnvFactory.getEnv().getDataNodeWrapper(i), TABLE_SQL_DIALECT); + Statement statement = connection.createStatement()) { + statement.execute("CREATE TABLE " + table + "(device_id STRING TAG, s1 INT32 FIELD)"); + statement.execute( + "INSERT INTO " + + table + + "(time, device_id, s1) VALUES (1,'d1',1),(2,'d2',2),(3,'d3',3)"); + statement.execute("FLUSH"); + assertRowsOnEveryDataNode(table, 1, 2, 3); + + statement.execute("DELETE FROM " + table + predicate); + // Successful execution alone cannot detect the silent no-op reported in TDB-449. + // Allow replica propagation before concluding that the deletion did not take effect. + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertRowsOnEveryDataNode(table, remainingTimes)); + statement.execute("FLUSH"); + assertRowsOnEveryDataNode(table, remainingTimes); + } + } + } + + private void assertRowsOnEveryDataNode(String table, int... expectedTimes) throws SQLException { + for (DataNodeWrapper dataNode : EnvFactory.getEnv().getDataNodeWrapperList()) { + String context = table + " on " + dataNode.getIpAndPortString(); + try (Connection connection = + EnvFactory.getEnv() + .getConnection( + dataNode, + SessionConfig.DEFAULT_USER, + SessionConfig.DEFAULT_PASSWORD, + TABLE_SQL_DIALECT); + Statement statement = connection.createStatement()) { + try (ResultSet resultSet = statement.executeQuery("SELECT count(*) FROM " + table)) { + assertTrue(context, resultSet.next()); + assertEquals(context, expectedTimes.length, resultSet.getLong(1)); + assertFalse(context, resultSet.next()); + } + try (ResultSet resultSet = + statement.executeQuery("SELECT time, device_id, s1 FROM " + table + " ORDER BY time")) { + for (int time : expectedTimes) { + assertTrue(context, resultSet.next()); + assertEquals(context, time, resultSet.getLong("time")); + assertEquals(context, "d" + time, resultSet.getString("device_id")); + assertEquals(context, time, resultSet.getInt("s1")); + } + assertFalse(context, resultSet.next()); + } + } + } + } +} diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/partition/GetRegionGroupsByTimeResp.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/partition/GetRegionGroupsByTimeResp.java index 3166ec30df14f..a51a9ef513cf0 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/partition/GetRegionGroupsByTimeResp.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/partition/GetRegionGroupsByTimeResp.java @@ -19,12 +19,16 @@ package org.apache.iotdb.confignode.consensus.response.partition; +import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId; import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.confignode.rpc.thrift.TGetRegionGroupsByTimeResp; import org.apache.iotdb.consensus.common.DataSet; import org.apache.iotdb.rpc.TSStatusCode; +import java.util.Collections; +import java.util.HashSet; +import java.util.Map; import java.util.Set; public class GetRegionGroupsByTimeResp implements DataSet { @@ -43,6 +47,27 @@ public TSStatus getStatus() { return status; } + /** Return a response whose replica lists put the current leader first. */ + public GetRegionGroupsByTimeResp reorderByLeader( + final Map regionLeaderMap) { + if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + return this; + } + final Set reordered = new HashSet<>(); + for (final TRegionReplicaSet replicaSet : regionReplicaSets) { + final TRegionReplicaSet copy = replicaSet.deepCopy(); + final int leaderId = regionLeaderMap.getOrDefault(copy.getRegionId(), -1); + for (int i = 0; i < copy.getDataNodeLocationsSize(); i++) { + if (copy.getDataNodeLocations().get(i).getDataNodeId() == leaderId) { + Collections.swap(copy.getDataNodeLocations(), 0, i); + break; + } + } + reordered.add(copy); + } + return new GetRegionGroupsByTimeResp(status, reordered); + } + public TGetRegionGroupsByTimeResp convertToRpcResp() { TGetRegionGroupsByTimeResp resp = new TGetRegionGroupsByTimeResp(); resp.setStatus(status); diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java index 249a4f9bd0e05..48ebe498a50db 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java @@ -1329,7 +1329,8 @@ public GetRegionGroupsByTimeResp getRegionGroupsByTime(final TGetRegionGroupsByT final GetRegionGroupsByTimePlan plan = new GetRegionGroupsByTimePlan(req.getDatabase(), req.getStartTime(), req.getEndTime()); try { - return (GetRegionGroupsByTimeResp) getConsensusManager().read(plan); + return ((GetRegionGroupsByTimeResp) getConsensusManager().read(plan)) + .reorderByLeader(getLoadManager().getRegionLeaderMap()); } catch (final ConsensusException e) { LOGGER.warn(CONSENSUS_READ_ERROR, e); final TSStatus res = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()); From 73303cd3cb31939f768ccdc58df35ea0e628dba6 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 30 Sep 2026 17:16:19 +0800 Subject: [PATCH 2/2] spotless --- .../relational/it/db/it/IoTDBDeletionTableClusterIT.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java index 34804aa6f2ec0..2ecd502a0a686 100644 --- a/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBDeletionTableClusterIT.java @@ -73,8 +73,8 @@ public static void tearDownClass() { } /** - * deleting time=2 from three flushed devices must leave exactly (1,d1,1) and (3,d3,3), - * through every DataNode coordinator, both before and after another flush. + * deleting time=2 from three flushed devices must leave exactly (1,d1,1) and (3,d3,3), through + * every DataNode coordinator, both before and after another flush. */ @Test public void testDeleteByTimeAfterFlush() throws SQLException { @@ -82,8 +82,8 @@ public void testDeleteByTimeAfterFlush() throws SQLException { } /** - * an unconditional DELETE must remove all three flushed rows through every DataNode - * coordinator, and another flush must not make the deleted data visible again. + * an unconditional DELETE must remove all three flushed rows through every DataNode coordinator, + * and another flush must not make the deleted data visible again. */ @Test public void testDeleteAllAfterFlush() throws SQLException {