Skip to content

Commit df27bbc

Browse files
authored
Fix table delete routing to region leaders (#18779)
* Fix table delete routing to region leaders * spotless
1 parent e163f3b commit df27bbc

3 files changed

Lines changed: 179 additions & 1 deletion

File tree

Lines changed: 152 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,152 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
package org.apache.iotdb.relational.it.db.it;
21+
22+
import org.apache.iotdb.consensus.ConsensusFactory;
23+
import org.apache.iotdb.isession.SessionConfig;
24+
import org.apache.iotdb.it.env.EnvFactory;
25+
import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper;
26+
import org.apache.iotdb.it.framework.IoTDBTestRunner;
27+
import org.apache.iotdb.itbase.category.TableClusterIT;
28+
29+
import org.awaitility.Awaitility;
30+
import org.junit.AfterClass;
31+
import org.junit.BeforeClass;
32+
import org.junit.Test;
33+
import org.junit.experimental.categories.Category;
34+
import org.junit.runner.RunWith;
35+
36+
import java.sql.Connection;
37+
import java.sql.ResultSet;
38+
import java.sql.SQLException;
39+
import java.sql.Statement;
40+
import java.util.concurrent.TimeUnit;
41+
42+
import static org.apache.iotdb.itbase.env.BaseEnv.TABLE_SQL_DIALECT;
43+
import static org.junit.Assert.assertEquals;
44+
import static org.junit.Assert.assertFalse;
45+
import static org.junit.Assert.assertTrue;
46+
47+
@RunWith(IoTDBTestRunner.class)
48+
@Category(TableClusterIT.class)
49+
public class IoTDBDeletionTableClusterIT {
50+
51+
private static final String DATABASE = "sc3";
52+
53+
@BeforeClass
54+
public static void setUpClass() throws SQLException {
55+
EnvFactory.getEnv()
56+
.getConfig()
57+
.getCommonConfig()
58+
.setConfigNodeConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
59+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
60+
.setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
61+
.setSchemaReplicationFactor(2)
62+
.setDataReplicationFactor(2);
63+
EnvFactory.getEnv().initClusterEnvironment(3, 3);
64+
try (Connection connection = EnvFactory.getEnv().getConnection(TABLE_SQL_DIALECT);
65+
Statement statement = connection.createStatement()) {
66+
statement.execute("CREATE DATABASE " + DATABASE);
67+
}
68+
}
69+
70+
@AfterClass
71+
public static void tearDownClass() {
72+
EnvFactory.getEnv().cleanClusterEnvironment();
73+
}
74+
75+
/**
76+
* deleting time=2 from three flushed devices must leave exactly (1,d1,1) and (3,d3,3), through
77+
* every DataNode coordinator, both before and after another flush.
78+
*/
79+
@Test
80+
public void testDeleteByTimeAfterFlush() throws SQLException {
81+
checkDeletionOnEveryDataNode("time_delete_", " WHERE time = 2", 1, 3);
82+
}
83+
84+
/**
85+
* an unconditional DELETE must remove all three flushed rows through every DataNode coordinator,
86+
* and another flush must not make the deleted data visible again.
87+
*/
88+
@Test
89+
public void testDeleteAllAfterFlush() throws SQLException {
90+
checkDeletionOnEveryDataNode("full_delete_", "");
91+
}
92+
93+
private void checkDeletionOnEveryDataNode(
94+
String tablePrefix, String predicate, int... remainingTimes) throws SQLException {
95+
for (int i = 0; i < EnvFactory.getEnv().getDataNodeWrapperList().size(); i++) {
96+
String table = DATABASE + "." + tablePrefix + i;
97+
// Pin each DELETE to a different coordinator and use fresh data for every attempt.
98+
try (Connection connection =
99+
EnvFactory.getEnv()
100+
.getWriteOnlyConnectionWithSpecifiedDataNode(
101+
EnvFactory.getEnv().getDataNodeWrapper(i), TABLE_SQL_DIALECT);
102+
Statement statement = connection.createStatement()) {
103+
statement.execute("CREATE TABLE " + table + "(device_id STRING TAG, s1 INT32 FIELD)");
104+
statement.execute(
105+
"INSERT INTO "
106+
+ table
107+
+ "(time, device_id, s1) VALUES (1,'d1',1),(2,'d2',2),(3,'d3',3)");
108+
statement.execute("FLUSH");
109+
assertRowsOnEveryDataNode(table, 1, 2, 3);
110+
111+
statement.execute("DELETE FROM " + table + predicate);
112+
// Successful execution alone cannot detect the silent no-op reported in TDB-449.
113+
// Allow replica propagation before concluding that the deletion did not take effect.
114+
Awaitility.await()
115+
.atMost(10, TimeUnit.SECONDS)
116+
.untilAsserted(() -> assertRowsOnEveryDataNode(table, remainingTimes));
117+
statement.execute("FLUSH");
118+
assertRowsOnEveryDataNode(table, remainingTimes);
119+
}
120+
}
121+
}
122+
123+
private void assertRowsOnEveryDataNode(String table, int... expectedTimes) throws SQLException {
124+
for (DataNodeWrapper dataNode : EnvFactory.getEnv().getDataNodeWrapperList()) {
125+
String context = table + " on " + dataNode.getIpAndPortString();
126+
try (Connection connection =
127+
EnvFactory.getEnv()
128+
.getConnection(
129+
dataNode,
130+
SessionConfig.DEFAULT_USER,
131+
SessionConfig.DEFAULT_PASSWORD,
132+
TABLE_SQL_DIALECT);
133+
Statement statement = connection.createStatement()) {
134+
try (ResultSet resultSet = statement.executeQuery("SELECT count(*) FROM " + table)) {
135+
assertTrue(context, resultSet.next());
136+
assertEquals(context, expectedTimes.length, resultSet.getLong(1));
137+
assertFalse(context, resultSet.next());
138+
}
139+
try (ResultSet resultSet =
140+
statement.executeQuery("SELECT time, device_id, s1 FROM " + table + " ORDER BY time")) {
141+
for (int time : expectedTimes) {
142+
assertTrue(context, resultSet.next());
143+
assertEquals(context, time, resultSet.getLong("time"));
144+
assertEquals(context, "d" + time, resultSet.getString("device_id"));
145+
assertEquals(context, time, resultSet.getInt("s1"));
146+
}
147+
assertFalse(context, resultSet.next());
148+
}
149+
}
150+
}
151+
}
152+
}

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/response/partition/GetRegionGroupsByTimeResp.java‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,12 +19,16 @@
1919

2020
package org.apache.iotdb.confignode.consensus.response.partition;
2121

22+
import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
2223
import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
2324
import org.apache.iotdb.common.rpc.thrift.TSStatus;
2425
import org.apache.iotdb.confignode.rpc.thrift.TGetRegionGroupsByTimeResp;
2526
import org.apache.iotdb.consensus.common.DataSet;
2627
import org.apache.iotdb.rpc.TSStatusCode;
2728

29+
import java.util.Collections;
30+
import java.util.HashSet;
31+
import java.util.Map;
2832
import java.util.Set;
2933

3034
public class GetRegionGroupsByTimeResp implements DataSet {
@@ -43,6 +47,27 @@ public TSStatus getStatus() {
4347
return status;
4448
}
4549

50+
/** Return a response whose replica lists put the current leader first. */
51+
public GetRegionGroupsByTimeResp reorderByLeader(
52+
final Map<TConsensusGroupId, Integer> regionLeaderMap) {
53+
if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
54+
return this;
55+
}
56+
final Set<TRegionReplicaSet> reordered = new HashSet<>();
57+
for (final TRegionReplicaSet replicaSet : regionReplicaSets) {
58+
final TRegionReplicaSet copy = replicaSet.deepCopy();
59+
final int leaderId = regionLeaderMap.getOrDefault(copy.getRegionId(), -1);
60+
for (int i = 0; i < copy.getDataNodeLocationsSize(); i++) {
61+
if (copy.getDataNodeLocations().get(i).getDataNodeId() == leaderId) {
62+
Collections.swap(copy.getDataNodeLocations(), 0, i);
63+
break;
64+
}
65+
}
66+
reordered.add(copy);
67+
}
68+
return new GetRegionGroupsByTimeResp(status, reordered);
69+
}
70+
4671
public TGetRegionGroupsByTimeResp convertToRpcResp() {
4772
TGetRegionGroupsByTimeResp resp = new TGetRegionGroupsByTimeResp();
4873
resp.setStatus(status);

‎iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1329,7 +1329,8 @@ public GetRegionGroupsByTimeResp getRegionGroupsByTime(final TGetRegionGroupsByT
13291329
final GetRegionGroupsByTimePlan plan =
13301330
new GetRegionGroupsByTimePlan(req.getDatabase(), req.getStartTime(), req.getEndTime());
13311331
try {
1332-
return (GetRegionGroupsByTimeResp) getConsensusManager().read(plan);
1332+
return ((GetRegionGroupsByTimeResp) getConsensusManager().read(plan))
1333+
.reorderByLeader(getLoadManager().getRegionLeaderMap());
13331334
} catch (final ConsensusException e) {
13341335
LOGGER.warn(CONSENSUS_READ_ERROR, e);
13351336
final TSStatus res = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());

0 commit comments

Comments
 (0)