Skip to content

Commit 825bf9e

Browse files
authored
Stop data region write retry promptly (#18730)
1 parent 50d86cf commit 825bf9e

2 files changed

Lines changed: 57 additions & 4 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataRegionStateMachine.java‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -265,11 +265,13 @@ protected TSStatus write(PlanNode planNode) {
265265
.PIPE_LOG_WRITE_OPERATION_STILL_FAILED_AFTER_RETRY_TIMES_BECAUSE_15EEA702,
266266
MAX_WRITE_RETRY_TIMES,
267267
result.getCode());
268+
break;
268269
}
269270
try {
270-
Thread.sleep(WRITE_RETRY_WAIT_TIME_IN_MS);
271+
waitBeforeNextWriteRetry();
271272
} catch (InterruptedException e) {
272273
Thread.currentThread().interrupt();
274+
break;
273275
}
274276
} else {
275277
if (TSStatusCode.TABLE_NOT_EXISTS.getStatusCode() == result.getCode()
@@ -282,6 +284,10 @@ protected TSStatus write(PlanNode planNode) {
282284
return result;
283285
}
284286

287+
protected void waitBeforeNextWriteRetry() throws InterruptedException {
288+
Thread.sleep(WRITE_RETRY_WAIT_TIME_IN_MS);
289+
}
290+
285291
@Override
286292
public DataSet read(IConsensusRequest request) {
287293
if (region == null) {

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/consensus/statemachine/dataregion/DataRegionStateMachineTest.java‎

Lines changed: 50 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -141,13 +141,14 @@ public void testPipeGeneratedWriteProcessRejectSkipsStateMachineRetry() {
141141

142142
@Test
143143
public void testNonPipeWriteProcessRejectCanStillRetry() {
144-
final DataRegionStateMachine stateMachine = new DataRegionStateMachine(null);
144+
final TestingDataRegionStateMachine stateMachine = new TestingDataRegionStateMachine(false);
145145
final RetryControlledPlanNode planNode = new RetryControlledPlanNode(false);
146146

147147
final TSStatus status = stateMachine.write(planNode);
148148

149149
Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), status.getCode());
150150
Assert.assertEquals(2, planNode.getAcceptCount());
151+
Assert.assertEquals(1, stateMachine.getWaitCount());
151152
}
152153

153154
@Test
@@ -164,7 +165,7 @@ public void testMetadataLeaseFencedNoRetry() {
164165

165166
@Test
166167
public void testMetadataLeaseFencedRetryRequiredRetriesAndFails() {
167-
final DataRegionStateMachine stateMachine = new DataRegionStateMachine(null);
168+
final TestingDataRegionStateMachine stateMachine = new TestingDataRegionStateMachine(false);
168169
final FixedStatusPlanNode planNode =
169170
new FixedStatusPlanNode(TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode());
170171

@@ -173,11 +174,12 @@ public void testMetadataLeaseFencedRetryRequiredRetriesAndFails() {
173174
Assert.assertEquals(
174175
TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode(), status.getCode());
175176
Assert.assertEquals(5, planNode.getAcceptCount());
177+
Assert.assertEquals(4, stateMachine.getWaitCount());
176178
}
177179

178180
@Test
179181
public void testMetadataLeaseFencedRetryRequiredRetriesAndSucceeds() {
180-
final DataRegionStateMachine stateMachine = new DataRegionStateMachine(null);
182+
final TestingDataRegionStateMachine stateMachine = new TestingDataRegionStateMachine(false);
181183
final SequenceStatusPlanNode planNode =
182184
new SequenceStatusPlanNode(
183185
TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode(),
@@ -188,6 +190,51 @@ public void testMetadataLeaseFencedRetryRequiredRetriesAndSucceeds() {
188190

189191
Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), status.getCode());
190192
Assert.assertEquals(3, planNode.getAcceptCount());
193+
Assert.assertEquals(2, stateMachine.getWaitCount());
194+
}
195+
196+
@Test
197+
public void testInterruptedWriteRetryStopsImmediately() {
198+
Thread.interrupted();
199+
try {
200+
final TestingDataRegionStateMachine stateMachine = new TestingDataRegionStateMachine(true);
201+
final FixedStatusPlanNode planNode =
202+
new FixedStatusPlanNode(
203+
TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode());
204+
205+
final TSStatus status = stateMachine.write(planNode);
206+
207+
Assert.assertEquals(
208+
TSStatusCode.METADATA_LEASE_FENCED_RETRY_REQUIRED.getStatusCode(), status.getCode());
209+
Assert.assertEquals(1, planNode.getAcceptCount());
210+
Assert.assertEquals(1, stateMachine.getWaitCount());
211+
Assert.assertTrue(Thread.currentThread().isInterrupted());
212+
} finally {
213+
Thread.interrupted();
214+
}
215+
}
216+
217+
private static class TestingDataRegionStateMachine extends DataRegionStateMachine {
218+
219+
private final boolean interruptOnWait;
220+
private int waitCount;
221+
222+
private TestingDataRegionStateMachine(final boolean interruptOnWait) {
223+
super(null);
224+
this.interruptOnWait = interruptOnWait;
225+
}
226+
227+
private int getWaitCount() {
228+
return waitCount;
229+
}
230+
231+
@Override
232+
protected void waitBeforeNextWriteRetry() throws InterruptedException {
233+
++waitCount;
234+
if (interruptOnWait) {
235+
throw new InterruptedException();
236+
}
237+
}
191238
}
192239

193240
private static class FixedStatusPlanNode extends PlanNode {

0 commit comments

Comments
 (0)