Skip to content

Commit a333099

Browse files
authored
Fix scan termination for finished fragment instances (#18701)
1 parent af6f8f5 commit a333099

5 files changed

Lines changed: 294 additions & 7 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/driver/Driver.java‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import org.apache.iotdb.db.conf.IoTDBDescriptor;
2727
import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
2828
import org.apache.iotdb.db.queryengine.execution.exchange.sink.ISink;
29+
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceFinishedException;
2930
import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext;
3031
import org.apache.iotdb.db.queryengine.execution.schedule.task.DriverTaskId;
3132
import org.apache.iotdb.db.queryengine.metric.QueryMetricsManager;
@@ -248,6 +249,11 @@ private ListenableFuture<?> processInternal() {
248249
}
249250
}
250251
return NOT_BLOCKED;
252+
} catch (FragmentInstanceFinishedException e) {
253+
// The fragment may finish before its notification thread closes this driver. Release the
254+
// driver through the normal cleanup path without aborting other fragments of the query.
255+
state.compareAndSet(State.ALIVE, State.NEED_DESTRUCTION);
256+
return NOT_BLOCKED;
251257
} catch (Throwable t) {
252258
Throwable actualCause = t;
253259
if (actualCause.getCause() instanceof IoTDBRuntimeException) {
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
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.db.queryengine.execution.fragment;
21+
22+
import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
23+
import org.apache.iotdb.db.queryengine.common.FragmentInstanceId;
24+
25+
/**
26+
* Internal control signal for stopping a scan after its fragment has finished. This is unchecked so
27+
* scan operators do not wrap it as an I/O failure before the driver can handle normal termination.
28+
*/
29+
public class FragmentInstanceFinishedException extends RuntimeException {
30+
31+
public FragmentInstanceFinishedException(FragmentInstanceId fragmentInstanceId) {
32+
super(
33+
String.format(
34+
DataNodeQueryMessages.EXCEPTION_FRAGMENT_INSTANCE_ARG_IS_ALREADY_ARG_B44984B4,
35+
fragmentInstanceId,
36+
FragmentInstanceState.FINISHED));
37+
}
38+
}

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesScanUtil.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import org.apache.iotdb.db.exception.CorruptedTsFileException;
2626
import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
2727
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext;
28+
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceFinishedException;
2829
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceState;
2930
import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext;
3031
import org.apache.iotdb.db.queryengine.metric.SeriesScanCostMetricSet;
@@ -1587,6 +1588,9 @@ private void checkFragmentInstanceState() throws IOException {
15871588
return;
15881589
}
15891590
FragmentInstanceState state = context.getStateMachine().getState();
1591+
if (state == FragmentInstanceState.FINISHED) {
1592+
throw new FragmentInstanceFinishedException(context.getId());
1593+
}
15901594
if (state.isDone()) {
15911595
// A scan over many overlapping files may stay in one operator call long after cancellation.
15921596
// Exit on the driver thread so it can release its lock and finish resource cleanup.

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesScanUtilCancellationTest.java‎

Lines changed: 165 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import org.apache.iotdb.calc.execution.operator.Operator;
2323
import org.apache.iotdb.calc.plan.planner.memory.MemoryReservationManager;
2424
import org.apache.iotdb.commons.exception.QueryTimeoutException;
25+
import org.apache.iotdb.commons.exception.SemanticException;
2526
import org.apache.iotdb.commons.path.AlignedFullPath;
2627
import org.apache.iotdb.commons.path.NonAlignedFullPath;
2728
import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -34,6 +35,7 @@
3435
import org.apache.iotdb.db.queryengine.execution.exchange.sink.ISink;
3536
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext;
3637
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceExecution;
38+
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceFinishedException;
3739
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceState;
3840
import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine;
3941
import org.apache.iotdb.db.queryengine.execution.schedule.IDriverScheduler;
@@ -63,6 +65,7 @@
6365
import java.util.concurrent.Executor;
6466
import java.util.concurrent.ExecutorService;
6567
import java.util.concurrent.Executors;
68+
import java.util.concurrent.Future;
6669
import java.util.concurrent.TimeUnit;
6770
import java.util.function.IntConsumer;
6871

@@ -71,9 +74,11 @@
7174
import static org.junit.Assert.assertSame;
7275
import static org.junit.Assert.assertThrows;
7376
import static org.junit.Assert.assertTrue;
77+
import static org.mockito.ArgumentMatchers.any;
7478
import static org.mockito.ArgumentMatchers.anyList;
7579
import static org.mockito.Mockito.doReturn;
7680
import static org.mockito.Mockito.mock;
81+
import static org.mockito.Mockito.never;
7782
import static org.mockito.Mockito.verify;
7883
import static org.mockito.Mockito.when;
7984

@@ -139,8 +144,12 @@ public void testTerminalStateBeforeScanDoesNotLoadMetadata() {
139144
throw new AssertionError(state);
140145
}
141146

142-
IOException exception = assertThrows(IOException.class, scanner::hasNextFile);
143-
assertSame(context.getFailureCause().orElse(null), exception.getCause());
147+
if (state == FragmentInstanceState.FINISHED) {
148+
assertThrows(FragmentInstanceFinishedException.class, scanner::hasNextFile);
149+
} else {
150+
IOException exception = assertThrows(IOException.class, scanner::hasNextFile);
151+
assertSame(context.getFailureCause().orElse(null), exception.getCause());
152+
}
144153
assertEquals(0, loadedFiles);
145154
assertEquals(state, context.getStateMachine().getState());
146155
}
@@ -237,12 +246,22 @@ public void testCompactionContextWithoutStateMachine() throws IOException {
237246

238247
@Test(timeout = 15000)
239248
public void testTimeoutUnblocksDriverResourceCleanup() throws Exception {
249+
assertFailureUnblocksDriverResourceCleanup(new QueryTimeoutException());
250+
}
251+
252+
@Test(timeout = 15000)
253+
public void testSemanticFailureUnblocksDriverResourceCleanup() throws Exception {
254+
assertFailureUnblocksDriverResourceCleanup(
255+
new SemanticException("Scalar sub-query has returned multiple rows."));
256+
}
257+
258+
private void assertFailureUnblocksDriverResourceCleanup(RuntimeException failure)
259+
throws Exception {
240260
ExecutorService notifications = Executors.newSingleThreadExecutor();
241261
FragmentInstanceContext context = newContext(notifications);
242262
context.initializeNumOfDrivers(1);
243263
try {
244264
CountDownLatch closeRequested = new CountDownLatch(1);
245-
QueryTimeoutException timeout = new QueryTimeoutException();
246265
SeriesScanUtil scanner =
247266
newScanner(
248267
context,
@@ -251,7 +270,7 @@ public void testTimeoutUnblocksDriverResourceCleanup() throws Exception {
251270
4,
252271
count -> {
253272
if (count == 2) {
254-
context.failed(timeout);
273+
context.failed(failure);
255274
try {
256275
// The notification thread has requested close while this thread owns the
257276
// driver lock, just as in a query that is still loading metadata.
@@ -296,10 +315,9 @@ public void close() {
296315
exchangeManager);
297316

298317
assertSame(
299-
timeout,
318+
failure,
300319
assertThrows(
301-
QueryTimeoutException.class,
302-
() -> driver.processFor(new Duration(1, TimeUnit.SECONDS))));
320+
RuntimeException.class, () -> driver.processFor(new Duration(1, TimeUnit.SECONDS))));
303321
// This can complete only after the preceding cleanup callback gets past allDriversClosed.
304322
notifications.submit(() -> {}).get(5, TimeUnit.SECONDS);
305323
assertEquals(2, loadedFiles);
@@ -316,6 +334,146 @@ public void close() {
316334
}
317335
}
318336

337+
@Test(timeout = 60000)
338+
public void testFinishedScanUnblocksDriverResourceCleanup() throws Exception {
339+
for (boolean sequence : new boolean[] {true, false}) {
340+
for (boolean waitForClose : new boolean[] {false, true}) {
341+
assertFinishedScanUnblocksDriverResourceCleanup(sequence, waitForClose);
342+
}
343+
}
344+
}
345+
346+
private void assertFinishedScanUnblocksDriverResourceCleanup(
347+
boolean sequence, boolean waitForClose) throws Exception {
348+
ExecutorService notifications = Executors.newSingleThreadExecutor();
349+
ExecutorService worker = Executors.newSingleThreadExecutor();
350+
CountDownLatch allowNotifications = new CountDownLatch(waitForClose ? 0 : 1);
351+
CountDownLatch scanEntered = new CountDownLatch(1);
352+
CountDownLatch resumeScan = new CountDownLatch(1);
353+
CountDownLatch closeRequested = new CountDownLatch(1);
354+
// Cover both orderings: the FI is FINISHED before driver.close(), and close is already pending.
355+
notifications.submit(() -> await(allowNotifications));
356+
FragmentInstanceContext context = newContext(notifications);
357+
context.initializeNumOfDrivers(1);
358+
Operator operator = mock(Operator.class);
359+
doReturn(NOT_BLOCKED).when(operator).isBlocked();
360+
when(operator.hasNextWithTimer()).thenReturn(true);
361+
SeriesScanUtil scanner =
362+
newScanner(
363+
context,
364+
Ordering.ASC,
365+
sequence,
366+
4,
367+
count -> {
368+
if (count == 1) {
369+
scanEntered.countDown();
370+
await(resumeScan);
371+
}
372+
});
373+
when(operator.nextWithTimer())
374+
.thenAnswer(
375+
invocation -> {
376+
scanner.hasNextFile();
377+
return null;
378+
});
379+
ISink sink = mock(ISink.class);
380+
doReturn(NOT_BLOCKED).when(sink).isFull();
381+
DataDriverContext driverContext = new DataDriverContext(context, 0);
382+
driverContext.setSink(sink);
383+
DataDriver driver =
384+
new DataDriver(operator, driverContext, 0) {
385+
@Override
386+
public void close() {
387+
super.close();
388+
closeRequested.countDown();
389+
}
390+
};
391+
MPPDataExchangeManager exchangeManager = mock(MPPDataExchangeManager.class);
392+
IDriverScheduler scheduler = mock(IDriverScheduler.class);
393+
FragmentInstanceExecution.createFragmentInstanceExecution(
394+
scheduler,
395+
context.getId(),
396+
context,
397+
Collections.singletonList(driver),
398+
sink,
399+
context.getStateMachine(),
400+
1000,
401+
false,
402+
exchangeManager);
403+
try {
404+
Future<?> execution =
405+
worker.submit(() -> driver.processFor(new Duration(1, TimeUnit.SECONDS)));
406+
await(scanEntered);
407+
context.finished();
408+
if (waitForClose) {
409+
await(closeRequested);
410+
}
411+
resumeScan.countDown();
412+
execution.get(5, TimeUnit.SECONDS);
413+
assertTrue(driver.isFinished());
414+
assertEquals(1, loadedFiles);
415+
assertEquals(FragmentInstanceState.FINISHED, context.getStateMachine().getState());
416+
assertTrue(context.getStateMachine().getFailureCauses().isEmpty());
417+
allowNotifications.countDown();
418+
notifications.submit(() -> {}).get(5, TimeUnit.SECONDS);
419+
verify(operator).close();
420+
verify(context.getMemoryReservationContext()).releaseAllReservedMemory();
421+
verify(exchangeManager)
422+
.deRegisterFragmentInstanceFromMemoryPool(
423+
context.getId().getQueryId().getId(), context.getId().getFragmentInstanceId(), true);
424+
verify(scheduler, never()).abortFragmentInstance(any(), any());
425+
} finally {
426+
resumeScan.countDown();
427+
allowNotifications.countDown();
428+
worker.shutdownNow();
429+
assertTrue(worker.awaitTermination(5, TimeUnit.SECONDS));
430+
driver.close();
431+
context.decrementNumOfUnClosedDriver();
432+
// Let FI cleanup finish after opening its latch, without interrupting its driver-close wait.
433+
notifications.shutdown();
434+
assertTrue(notifications.awaitTermination(5, TimeUnit.SECONDS));
435+
}
436+
}
437+
438+
@Test
439+
public void testFinishedFragmentDoesNotSuppressIOException() throws Exception {
440+
FragmentInstanceContext context = newContext(Runnable::run);
441+
context.initializeNumOfDrivers(1);
442+
IOException failure = new IOException("Failed to read file metadata");
443+
Operator operator = mock(Operator.class);
444+
doReturn(NOT_BLOCKED).when(operator).isBlocked();
445+
when(operator.hasNextWithTimer()).thenReturn(true);
446+
when(operator.nextWithTimer())
447+
.thenAnswer(
448+
invocation -> {
449+
context.finished();
450+
throw failure;
451+
});
452+
ISink sink = mock(ISink.class);
453+
doReturn(NOT_BLOCKED).when(sink).isFull();
454+
DataDriverContext driverContext = new DataDriverContext(context, 0);
455+
driverContext.setSink(sink);
456+
DataDriver driver = new DataDriver(operator, driverContext, 0);
457+
try {
458+
RuntimeException exception =
459+
assertThrows(
460+
RuntimeException.class, () -> driver.processFor(new Duration(1, TimeUnit.SECONDS)));
461+
assertSame(failure, exception.getCause());
462+
assertSame(failure, context.getFailureCause().get());
463+
} finally {
464+
driver.close();
465+
}
466+
}
467+
468+
private static void await(CountDownLatch latch) {
469+
try {
470+
assertTrue(latch.await(5, TimeUnit.SECONDS));
471+
} catch (InterruptedException e) {
472+
Thread.currentThread().interrupt();
473+
throw new AssertionError(e);
474+
}
475+
}
476+
319477
private FragmentInstanceContext newContext(Executor executor) {
320478
FragmentInstanceId id =
321479
new FragmentInstanceId(new PlanFragmentId(new QueryId("scan_cancellation"), 0), "0");

0 commit comments

Comments
 (0)