diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/BatchScanUtils.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/BatchScanUtils.java index aaf1b9711e..0c1bbf47f6 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/BatchScanUtils.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/BatchScanUtils.java @@ -19,6 +19,7 @@ import org.apache.fluss.row.InternalRow; import org.apache.fluss.utils.CloseableIterator; +import org.apache.fluss.utils.IOUtils; import java.time.Duration; import java.util.ArrayList; @@ -43,6 +44,7 @@ public static List collectRows(BatchScanner scanner) { } rows.addAll(toList(iterator)); } catch (Exception e) { + IOUtils.closeQuietly(scanner); throw new RuntimeException("Failed to collect rows", e); } } @@ -75,6 +77,8 @@ public static List collectLimitedRows(List scanners, scanner.close(); } } catch (Exception e) { + IOUtils.closeQuietly(scanner); + IOUtils.closeAllQuietly(scannerQueue); throw new RuntimeException("Failed to collect rows", e); } if (rows.size() >= limit) { @@ -82,13 +86,7 @@ public static List collectLimitedRows(List scanners, } } // may collect enough rows before drain all scanners, close all scanners in the queue - for (BatchScanner scanner : scannerQueue) { - try { - scanner.close(); - } catch (Exception e) { - throw new RuntimeException("Failed to close scanner", e); - } - } + IOUtils.closeAllQuietly(scannerQueue); return rows.size() > limit ? rows.subList(0, limit) : rows; } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java index d8292d11c8..393be460f3 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java @@ -33,6 +33,7 @@ import org.apache.fluss.row.InternalRow; import org.apache.fluss.row.KeyValueRow; import org.apache.fluss.utils.CloseableIterator; +import org.apache.fluss.utils.IOUtils; import javax.annotation.Nullable; @@ -197,17 +198,7 @@ private void pollLogRecords(Duration timeout) { @Override public void close() throws IOException { - try { - if (logScanner != null) { - logScanner.close(); - } - if (lakeRecordIterators != null) { - for (CloseableIterator iterator : lakeRecordIterators) { - iterator.close(); - } - } - } catch (Exception e) { - throw new IOException("Failed to close resources", e); - } + IOUtils.closeQuietly(logScanner); + IOUtils.closeAllQuietly(lakeRecordIterators); } } diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/CompositeBatchScannerTest.java b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/CompositeBatchScannerTest.java index 129fc028b1..0d254bc3ba 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/CompositeBatchScannerTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/CompositeBatchScannerTest.java @@ -35,6 +35,7 @@ import java.util.Queue; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Test for {@link CompositeBatchScanner}. */ class CompositeBatchScannerTest { @@ -120,6 +121,20 @@ void testCloseOnEmptyQueueDoesNotThrow() throws IOException { composite.close(); // should not throw } + @Test + void testLimitPollClosesRemainingScannersWhenOneFails() { + StubBatchScanner s1 = scanner(1); + FailingStubBatchScanner s2 = new FailingStubBatchScanner(); + StubBatchScanner s3 = scanner(3); + + CompositeBatchScanner composite = new CompositeBatchScanner(Arrays.asList(s1, s2, s3), 10); + + assertThatThrownBy(() -> composite.pollBatch(TIMEOUT)).isInstanceOf(RuntimeException.class); + assertThat(s1.closed).isTrue(); + assertThat(s2.closed).isTrue(); + assertThat(s3.closed).isTrue(); + } + // ------------------------------------------------------------------------- // Helpers // ------------------------------------------------------------------------- @@ -186,4 +201,20 @@ public void close() { closed = true; } } + + /** A stub {@link BatchScanner} that fails on {@link #pollBatch(Duration)}. */ + private static class FailingStubBatchScanner implements BatchScanner { + boolean closed = false; + + @Nullable + @Override + public CloseableIterator pollBatch(Duration timeout) throws IOException { + throw new IOException("injected poll failure"); + } + + @Override + public void close() { + closed = true; + } + } } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java index 14a14714bc..bb073bc298 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java @@ -377,14 +377,10 @@ public void shutdown(KvCloseMode closeMode) { .join(); IOUtils.closeQuietly(sharedWriteBufferManager); IOUtils.closeQuietly(sharedWriteBufferAccountingCache); - arrowBufferAllocator.close(); - memorySegmentPool.close(); - if (sharedRocksDBRateLimiter != null) { - sharedRocksDBRateLimiter.close(); - } - if (sharedBlockCache != null) { - sharedBlockCache.close(); - } + IOUtils.closeQuietly(arrowBufferAllocator); + IOUtils.closeQuietly(memorySegmentPool); + IOUtils.closeQuietly(sharedRocksDBRateLimiter); + IOUtils.closeQuietly(sharedBlockCache); LOG.info("Shut down KvManager complete."); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/log/LogSegments.java b/fluss-server/src/main/java/org/apache/fluss/server/log/LogSegments.java index e908752a28..326bdc6e46 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/log/LogSegments.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/log/LogSegments.java @@ -19,6 +19,7 @@ import org.apache.fluss.annotation.Internal; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.utils.IOUtils; import javax.annotation.concurrent.ThreadSafe; @@ -70,7 +71,7 @@ public void clear() { public void close() { for (LogSegment logSegment : segments.values()) { - logSegment.close(); + IOUtils.closeQuietly(logSegment::close); } }