diff --git a/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java b/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java index 9d9b2a4869d0..24efec2dd606 100644 --- a/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java +++ b/iotdb-core/consensus/src/main/i18n/en/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java @@ -256,6 +256,9 @@ private IoTConsensusMessages() {} "waiting target request timeout. current index: {}, target index: {}"; public static final String CURRENT_WAITING_INTERRUPTED = "current waiting is interrupted. SyncIndex: {}. Exception: "; + public static final String + MESSAGE_SYNC_LOG_REQUEST_WITH_SYNC_INDEX_ARG_WAS_INTERRUPTED_WHILE_WAITING_81B4ABB2 = + "SyncLog request with SyncIndex %d was interrupted while waiting"; // ===================== SyncStatus ===================== diff --git a/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java b/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java index 9fca07f7bde9..810ae0a3761e 100644 --- a/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java +++ b/iotdb-core/consensus/src/main/i18n/zh/org/apache/iotdb/consensus/i18n/IoTConsensusMessages.java @@ -254,6 +254,9 @@ private IoTConsensusMessages() {} "等待目标请求超时。当前 index:{},目标 index:{}"; public static final String CURRENT_WAITING_INTERRUPTED = "当前等待被中断。SyncIndex:{}。异常:"; + public static final String + MESSAGE_SYNC_LOG_REQUEST_WITH_SYNC_INDEX_ARG_WAS_INTERRUPTED_WHILE_WAITING_81B4ABB2 = + "等待期间 SyncLog 请求(SyncIndex %d)被中断"; // ===================== SyncStatus ===================== diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java index 1d55489f7cb8..430507a40489 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java @@ -1435,8 +1435,19 @@ private TSStatus cacheAndInsertLatestNode(DeserializedBatchIndexedConsensusReque } catch (InterruptedException e) { logger.warn( IoTConsensusMessages.CURRENT_WAITING_INTERRUPTED, request.getStartSyncIndex(), e); + requestCache.remove(request); + queueSortCondition.signalAll(); Thread.currentThread().interrupt(); - break; + return new TSStatus() + .setSubStatus( + Collections.nCopies( + request.getInsertNodes().size(), + RpcUtils.getStatus( + TSStatusCode.INTERNAL_SERVER_ERROR, + String.format( + IoTConsensusMessages + .MESSAGE_SYNC_LOG_REQUEST_WITH_SYNC_INDEX_ARG_WAS_INTERRUPTED_WHILE_WAITING_81B4ABB2, + request.getStartSyncIndex())))); } } long sortTime = System.nanoTime(); diff --git a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/ReplicateTest.java b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/ReplicateTest.java index ca96ffec4859..e096f5597734 100644 --- a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/ReplicateTest.java +++ b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/ReplicateTest.java @@ -21,16 +21,23 @@ import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType; import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.consensus.ConsensusGroupId; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.exception.StartupException; +import org.apache.iotdb.commons.request.IConsensusRequest; import org.apache.iotdb.consensus.ConsensusFactory; import org.apache.iotdb.consensus.common.ConsensusGroup; import org.apache.iotdb.consensus.common.Peer; +import org.apache.iotdb.consensus.common.request.DeserializedBatchIndexedConsensusRequest; +import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest; import org.apache.iotdb.consensus.config.ConsensusConfig; +import org.apache.iotdb.consensus.config.IoTConsensusConfig; import org.apache.iotdb.consensus.exception.ConsensusException; +import org.apache.iotdb.consensus.i18n.IoTConsensusMessages; import org.apache.iotdb.consensus.iot.util.TestEntry; import org.apache.iotdb.consensus.iot.util.TestStateMachine; +import org.apache.iotdb.rpc.TSStatusCode; import org.apache.ratis.util.FileUtils; import org.junit.After; @@ -48,6 +55,9 @@ import java.util.Arrays; import java.util.Collections; import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; public class ReplicateTest { @@ -313,6 +323,85 @@ public void replicateUsingWALTest() throws IOException, InterruptedException, Co } } + /** + * Verifies that a SyncLog request interrupted while waiting for its turn is rejected without + * being applied, and that the thread's interrupted status is preserved. + */ + @Test + public void syncLogInterruptedWhileWaitingTest() throws Exception { + servers.get(0).createLocalPeer(group.getGroupId(), group.getPeers()); + IoTConsensusServerImpl server = servers.get(0).getImpl(gid); + DeserializedBatchIndexedConsensusRequest request = + new DeserializedBatchIndexedConsensusRequest(1, 1, 1, peers.get(1).getNodeId(), 1); + request.add( + new IndexedConsensusRequest( + 1, 1, Collections.singletonList(new TestEntry(1, peers.get(1))))); + AtomicReference status = new AtomicReference<>(); + AtomicBoolean interrupted = new AtomicBoolean(); + + Thread syncLogThread = + new Thread( + () -> { + Thread.currentThread().interrupt(); + status.set(server.syncLog(peers.get(1).getNodeId(), request)); + interrupted.set(Thread.currentThread().isInterrupted()); + }); + syncLogThread.start(); + syncLogThread.join(); + + Assert.assertTrue(interrupted.get()); + Assert.assertEquals(1, status.get().getSubStatusSize()); + Assert.assertEquals( + TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode(), + status.get().getSubStatus().get(0).getCode()); + Assert.assertEquals( + String.format( + IoTConsensusMessages + .MESSAGE_SYNC_LOG_REQUEST_WITH_SYNC_INDEX_ARG_WAS_INTERRUPTED_WHILE_WAITING_81B4ABB2, + request.getStartSyncIndex()), + status.get().getSubStatus().get(0).getMessage()); + Assert.assertTrue(stateMachines.get(0).getRequestSet().isEmpty()); + } + + /** + * Verifies that WRITE_PROCESS_ERROR returned by a follower SyncLog write is retried by the leader + * and the request is eventually applied after the transient failure disappears. + */ + @Test + public void syncLogWriteProcessErrorTriggersLeaderRetryTest() throws Exception { + IoTConsensusConfig retryConfig = + IoTConsensusConfig.newBuilder() + .setReplication( + IoTConsensusConfig.Replication.newBuilder() + .setMaxWaitingTimeForWaitBatchInMs(50) + .setBasicRetryWaitTimeMs(10) + .setMaxRetryWaitTimeMs(100) + .build()) + .build(); + ConsensusConfig consensusConfig = + ConsensusConfig.newBuilder().setIoTConsensusConfig(retryConfig).build(); + servers.forEach(server -> server.reloadConsensusConfig(consensusConfig)); + + WriteProcessErrorOnceTestStateMachine failingStateMachine = + new WriteProcessErrorOnceTestStateMachine(); + stateMachines.set(1, failingStateMachine); + for (IoTConsensus server : servers) { + server.createLocalPeer(group.getGroupId(), group.getPeers()); + } + + TestEntry entry = new TestEntry(1, peers.get(0)); + Assert.assertEquals( + TSStatusCode.SUCCESS_STATUS.getStatusCode(), servers.get(0).write(gid, entry).getCode()); + long deadline = System.currentTimeMillis() + 5_000; + while (!failingStateMachine.getData().contains(entry) + && System.currentTimeMillis() < deadline) { + Thread.sleep(20); + } + + Assert.assertTrue(failingStateMachine.getWriteAttempts() >= 2); + Assert.assertTrue(failingStateMachine.getData().contains(entry)); + } + @Test public void parsingAndConstructIDTest() throws Exception { logger.info("Start ParsingAndConstructIDTest"); @@ -347,4 +436,21 @@ private void checkPeerList(IoTConsensusServerImpl iotServerImpl) { peers.stream().map(Peer::getNodeId).collect(Collectors.toSet()), iotServerImpl.getConfiguration().stream().map(Peer::getNodeId).collect(Collectors.toSet())); } + + private static class WriteProcessErrorOnceTestStateMachine extends TestStateMachine { + + private final AtomicInteger writeAttempts = new AtomicInteger(); + + @Override + public TSStatus write(IConsensusRequest request) { + if (writeAttempts.incrementAndGet() == 1) { + return new TSStatus(TSStatusCode.WRITE_PROCESS_ERROR.getStatusCode()); + } + return super.write(request); + } + + private int getWriteAttempts() { + return writeAttempts.get(); + } + } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/RetryUtils.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/RetryUtils.java index b07a95045c73..be84a0b2e3c1 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/RetryUtils.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/utils/RetryUtils.java @@ -40,6 +40,7 @@ public static boolean needRetryForWrite(int statusCode) { return statusCode == TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode() || statusCode == TSStatusCode.SYSTEM_READ_ONLY.getStatusCode() || statusCode == TSStatusCode.WRITE_PROCESS_REJECT.getStatusCode() + || statusCode == TSStatusCode.WRITE_PROCESS_ERROR.getStatusCode() || statusCode == TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode() || statusCode == TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode(); } diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/RetryUtilsTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/RetryUtilsTest.java new file mode 100644 index 000000000000..0c7d1dee9967 --- /dev/null +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/utils/RetryUtilsTest.java @@ -0,0 +1,46 @@ +/* + * 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.commons.utils; + +import org.apache.iotdb.rpc.TSStatusCode; + +import org.junit.Assert; +import org.junit.Test; + +public class RetryUtilsTest { + + /** + * Verifies that transient write failures, including WRITE_PROCESS_ERROR returned for an + * IOException during SyncLog, are retried while successful writes are not. + */ + @Test + public void testNeedRetryForWrite() { + Assert.assertTrue( + RetryUtils.needRetryForWrite(TSStatusCode.INTERNAL_SERVER_ERROR.getStatusCode())); + Assert.assertTrue(RetryUtils.needRetryForWrite(TSStatusCode.SYSTEM_READ_ONLY.getStatusCode())); + Assert.assertTrue( + RetryUtils.needRetryForWrite(TSStatusCode.WRITE_PROCESS_REJECT.getStatusCode())); + Assert.assertTrue( + RetryUtils.needRetryForWrite(TSStatusCode.WRITE_PROCESS_ERROR.getStatusCode())); + Assert.assertTrue( + RetryUtils.needRetryForWrite(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode())); + Assert.assertFalse(RetryUtils.needRetryForWrite(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + } +}