From fb9e28807e375632743b6e10876c7b95699e928d Mon Sep 17 00:00:00 2001 From: lisen Date: Wed, 16 Sep 2026 09:46:13 +0800 Subject: [PATCH] [client][server] Close remaining resources if one close fails Keep lake scanners, KV shutdown, limited batch collection, and log segments from leaking when an earlier close or poll throws. Co-authored-by: Cursor --- .../table/scanner/batch/BatchScanUtils.java | 12 +++---- .../batch/LakeSnapshotAndLogSplitScanner.java | 15 ++------- .../batch/CompositeBatchScannerTest.java | 31 +++++++++++++++++++ .../org/apache/fluss/server/kv/KvManager.java | 12 +++---- .../apache/fluss/server/log/LogSegments.java | 3 +- 5 files changed, 45 insertions(+), 28 deletions(-) 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 aaf1b9711e8..0c1bbf47f6f 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 d8292d11c8c..393be460f31 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 129fc028b13..0d254bc3ba2 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 14a14714bc5..bb073bc2980 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 e908752a28d..326bdc6e461 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); } }