diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvCloseMode.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvCloseMode.java new file mode 100644 index 00000000000..c8e94c450e1 --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvCloseMode.java @@ -0,0 +1,27 @@ +/* + * 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.fluss.server.kv; + +/** Defines whether closing a KV tablet must persist its current local RocksDB state. */ +public enum KvCloseMode { + /** Preserve the existing local-reopen behavior by allowing RocksDB to flush on close. */ + PRESERVE_LOCAL_STATE, + + /** Skip flushing data that recovery will rebuild from a snapshot and changelog. */ + DISCARD_UNPERSISTED_STATE +} 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 6a5c6079c94..2ff8aa23db9 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 @@ -63,6 +63,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; @@ -203,10 +204,17 @@ public void startup() { } public void shutdown() { - LOG.info("Shutting down KvManager"); + shutdown(KvCloseMode.PRESERVE_LOCAL_STATE); + } + + public void shutdown(KvCloseMode closeMode) { + Objects.requireNonNull(closeMode, "closeMode"); + LOG.info("Shutting down KvManager with close mode {}.", closeMode); isShutdown = true; List kvs = new ArrayList<>(currentKvs.values()); - closeTabletsConcurrently(kvs, "kv-tablet-closing", this::closeKvTablet).join(); + closeTabletsConcurrently( + kvs, "kv-tablet-closing", kvTablet -> closeKvTablet(kvTablet, closeMode)) + .join(); arrowBufferAllocator.close(); memorySegmentPool.close(); if (sharedRocksDBRateLimiter != null) { @@ -215,11 +223,15 @@ public void shutdown() { LOG.info("Shut down KvManager complete."); } - private void closeKvTablet(KvTablet kvTablet) { + private void closeKvTablet(KvTablet kvTablet, KvCloseMode closeMode) { try { - kvTablet.close(); + kvTablet.close(closeMode); } catch (Exception e) { - LOG.warn("Exception while closing kv tablet {}.", kvTablet.getTableBucket(), e); + LOG.warn( + "Exception while closing kv tablet {} with mode {}.", + kvTablet.getTableBucket(), + closeMode, + e); } } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java index da6b518266b..4b8c1109a8d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java @@ -842,7 +842,15 @@ public KvBatchWriter createKvBatchWriter() { } public void close() throws Exception { - LOG.debug("close kv tablet {} for table {}.", tableBucket, physicalPath); + close(KvCloseMode.PRESERVE_LOCAL_STATE); + } + + public void close(KvCloseMode closeMode) throws Exception { + LOG.debug( + "Close kv tablet {} for table {} with mode {}.", + tableBucket, + physicalPath, + closeMode); inWriteLock( kvLock, () -> { @@ -852,7 +860,7 @@ public void close() throws Exception { // Note: RocksDB metrics lifecycle is managed by TableMetricGroup // No need to close it here if (rocksDBKv != null) { - rocksDBKv.close(); + rocksDBKv.close(closeMode); } isClosed = true; }); @@ -864,7 +872,7 @@ public void drop() throws Exception { kvLock, () -> { // first close the kv. - close(); + close(KvCloseMode.DISCARD_UNPERSISTED_STATE); // then delete the directory. FileUtils.deleteDirectory(kvTabletDir); }); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java index f3998f4435c..4de4cfbc985 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java @@ -17,10 +17,12 @@ package org.apache.fluss.server.kv.rocksdb; +import org.apache.fluss.annotation.VisibleForTesting; import org.apache.fluss.exception.FlussRuntimeException; import org.apache.fluss.metrics.Counter; import org.apache.fluss.metrics.Histogram; import org.apache.fluss.rocksdb.RocksDBOperationUtils; +import org.apache.fluss.server.kv.KvCloseMode; import org.apache.fluss.server.utils.ResourceGuard; import org.apache.fluss.utils.BytesUtils; import org.apache.fluss.utils.IOUtils; @@ -28,6 +30,7 @@ import org.rocksdb.Cache; import org.rocksdb.ColumnFamilyHandle; import org.rocksdb.ColumnFamilyOptions; +import org.rocksdb.MutableDBOptions; import org.rocksdb.ReadOptions; import org.rocksdb.RocksDB; import org.rocksdb.RocksDBException; @@ -40,6 +43,7 @@ import java.io.IOException; import java.util.ArrayList; import java.util.List; +import java.util.Objects; /** A wrapper for the operation of {@link org.rocksdb.RocksDB}. */ public class RocksDBKv implements AutoCloseable { @@ -178,6 +182,13 @@ public void checkIfRocksDBClosed() { @Override public void close() throws Exception { + close(KvCloseMode.PRESERVE_LOCAL_STATE); + } + + /** Closes this RocksDB KV instance using the provided persistence mode. */ + public void close(KvCloseMode closeMode) throws Exception { + Objects.requireNonNull(closeMode); + if (this.closed) { return; } @@ -188,29 +199,45 @@ public void close() throws Exception { // parallel. rocksDBResourceGuard.close(); - // IMPORTANT: null reference to signal potential async checkpoint workers that the db was - // disposed, as - // working on the disposed object results in SEGFAULTS. - if (db != null) { + RocksDBException avoidFlushFailure = null; + try { + if (closeMode == KvCloseMode.DISCARD_UNPERSISTED_STATE) { + setAvoidFlushDuringShutdown(); + } + } catch (RocksDBException e) { + avoidFlushFailure = e; + } finally { + // IMPORTANT: null reference to signal potential async checkpoint workers that the db + // was disposed, as working on the disposed object results in SEGFAULTS. + if (db != null) { + + // RocksDB's native memory management requires that *all* CFs (including default) + // are closed before the DB is closed. See: + // https://github.com/facebook/rocksdb/wiki/RocksJava-Basics#opening-a-database-with-column-families + // Start with default CF ... + List columnFamilyOptions = new ArrayList<>(); + RocksDBOperationUtils.addColumnFamilyOptionsToCloseLater( + columnFamilyOptions, defaultColumnFamilyHandle); + IOUtils.closeQuietly(defaultColumnFamilyHandle); - // RocksDB's native memory management requires that *all* CFs (including default) are - // closed before the - // DB is closed. See: - // https://github.com/facebook/rocksdb/wiki/RocksJava-Basics#opening-a-database-with-column-families - // Start with default CF ... - List columnFamilyOptions = new ArrayList<>(); - RocksDBOperationUtils.addColumnFamilyOptionsToCloseLater( - columnFamilyOptions, defaultColumnFamilyHandle); - IOUtils.closeQuietly(defaultColumnFamilyHandle); + // ... and finally close the DB instance ... + IOUtils.closeQuietly(db); - // ... and finally close the DB instance ... - IOUtils.closeQuietly(db); + columnFamilyOptions.forEach(IOUtils::closeQuietly); - columnFamilyOptions.forEach(IOUtils::closeQuietly); + IOUtils.closeQuietly(optionsContainer); + } + this.closed = true; + } - IOUtils.closeQuietly(optionsContainer); + if (avoidFlushFailure != null) { + throw avoidFlushFailure; } - this.closed = true; + } + + @VisibleForTesting + void setAvoidFlushDuringShutdown() throws RocksDBException { + db.setDBOptions(MutableDBOptions.builder().setAvoidFlushDuringShutdown(true).build()); } public RocksDB getDb() { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java index 3d1ad4ef701..184a583f7d1 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/tablet/TabletServer.java @@ -39,6 +39,7 @@ import org.apache.fluss.server.authorizer.AuthorizerLoader; import org.apache.fluss.server.coordinator.LakeCatalogDynamicLoader; import org.apache.fluss.server.coordinator.MetadataManager; +import org.apache.fluss.server.kv.KvCloseMode; import org.apache.fluss.server.kv.KvManager; import org.apache.fluss.server.kv.scan.ScannerManager; import org.apache.fluss.server.kv.snapshot.DefaultCompletedKvSnapshotCommitter; @@ -571,9 +572,7 @@ void shutdownReplicaManager() throws InterruptedException { @VisibleForTesting void shutdownTabletManagers() throws IOException { - if (kvManager != null) { - kvManager.shutdown(); - } + shutdownKvManager(KvCloseMode.DISCARD_UNPERSISTED_STATE); if (remoteLogManager != null) { remoteLogManager.close(); @@ -584,6 +583,13 @@ void shutdownTabletManagers() throws IOException { } } + @VisibleForTesting + void shutdownKvManager(KvCloseMode closeMode) { + if (kvManager != null) { + kvManager.shutdown(closeMode); + } + } + private void controlledShutDown() { long startTime = System.currentTimeMillis(); LOG.info("Starting controlled shutdown."); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java index 530e92f4f7b..35e4895edb3 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java @@ -60,6 +60,7 @@ import java.io.File; import java.io.IOException; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; @@ -232,6 +233,29 @@ void testRecoveryAfterKvManagerShutDown(String partitionName) throws Exception { assertThat(kv2.multiGet(kv2Keys)).containsExactlyElementsOf(kv2Values); } + @ParameterizedTest + @MethodSource("partitionProvider") + void testDiscardShutdownDoesNotPersistUnflushedState(String partitionName) throws Exception { + initTableBuckets(partitionName); + KvTablet kv = getOrCreateKv(tablePath1, partitionName, tableBucket1); + byte[] key = "discarded-key".getBytes(StandardCharsets.UTF_8); + put(kv, kvRecordFactory.ofRecord(key, new Object[] {1, "value"})); + + kvManager.shutdown(KvCloseMode.DISCARD_UNPERSISTED_STATE); + kvManager = + KvManager.create( + conf, + zkClient, + logManager, + TestingMetricGroups.TABLET_SERVER_METRICS, + localDiskManager); + kvManager.startup(); + + KvTablet reopened = getOrCreateKv(tablePath1, partitionName, tableBucket1); + assertThat(reopened.multiGet(Collections.singletonList(key))) + .containsExactly((byte[]) null); + } + @ParameterizedTest @MethodSource("partitionProvider") void testRecoveryWithSchemaChange(String partitionName) throws Exception { @@ -317,6 +341,8 @@ void testSameTableNameInDifferentDb(String partitionName) throws Exception { void testDropKv(String partitionName) throws Exception { initTableBuckets(partitionName); KvTablet kv1 = getOrCreateKv(tablePath1, partitionName, tableBucket1); + byte[] key = "dropped-key".getBytes(StandardCharsets.UTF_8); + put(kv1, kvRecordFactory.ofRecord(key, new Object[] {1, "old"})); kvManager.dropKv(kv1.getTableBucket()); assertThat(kv1.getKvTabletDir()).doesNotExist(); @@ -324,6 +350,7 @@ void testDropKv(String partitionName) throws Exception { kv1 = getOrCreateKv(tablePath1, partitionName, tableBucket1); assertThat(kv1.getKvTabletDir()).exists(); + assertThat(kv1.multiGet(Collections.singletonList(key))).containsExactly((byte[]) null); assertThat(kvManager.getKv(tableBucket1)).isPresent(); } @@ -334,6 +361,26 @@ void testGetNonExistentKv() { assertThat(kv).isNotPresent(); } + @Test + void testShutdownRejectsNullCloseModeBeforeClosingTablets() throws Exception { + initTableBuckets(null); + KvTablet kv = getOrCreateKv(tablePath1, null, tableBucket1); + byte[] key = "still-usable-key".getBytes(StandardCharsets.UTF_8); + KvRecord record = kvRecordFactory.ofRecord(key, new Object[] {1, "value"}); + put(kv, record); + + assertThatThrownBy(() -> kvManager.shutdown(null)) + .isInstanceOf(NullPointerException.class) + .hasMessage("closeMode"); + + verifyMultiGet(kv, key, valueOf(record)); + + kvManager.shutdown(); + assertThatThrownBy(kv.getRocksDBKv()::checkIfRocksDBClosed) + .isInstanceOf(FlussRuntimeException.class); + kvManager = null; + } + @Test void testShutdownClosesKvTabletsConcurrently() throws Exception { int maxClosingThreads = 2; @@ -352,7 +399,7 @@ void testShutdownClosesKvTabletsConcurrently() throws Exception { KvManager managerToShutdown = kvManager; kvManager = null; ExecutorService shutdownExecutor = Executors.newSingleThreadExecutor(); - Future shutdownFuture = shutdownExecutor.submit(managerToShutdown::shutdown); + Future shutdownFuture = shutdownExecutor.submit(() -> managerToShutdown.shutdown()); try { waitUntil( () -> diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java index 80d27f8bd16..c6db154c1c5 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java @@ -18,15 +18,21 @@ package org.apache.fluss.server.kv.rocksdb; import org.apache.fluss.config.Configuration; +import org.apache.fluss.exception.FlussRuntimeException; +import org.apache.fluss.server.kv.KvCloseMode; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import org.rocksdb.RocksDBException; import java.io.File; import java.nio.file.Path; import java.util.Arrays; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.spy; /** Test for {@link org.apache.fluss.server.kv.rocksdb.RocksDBKv}. */ class RocksDBKvTest { @@ -34,15 +40,7 @@ class RocksDBKvTest { @Test void testRocksDbKv(@TempDir Path tempDir) throws Exception { File instanceBasePath = tempDir.toFile(); - RocksDBResourceContainer rocksDBResourceContainer = - new RocksDBResourceContainer(new Configuration(), instanceBasePath); - RocksDBKvBuilder rocksDBKvBuilder = - new RocksDBKvBuilder( - instanceBasePath, - rocksDBResourceContainer, - rocksDBResourceContainer.getColumnOptions()); - - try (RocksDBKv rocksDBKv = rocksDBKvBuilder.build()) { + try (RocksDBKv rocksDBKv = buildRocksDBKv(instanceBasePath)) { // put the k/v byte[] key = new byte[] {1, 2, 3}; byte[] val = new byte[] {1, 2}; @@ -64,4 +62,58 @@ void testRocksDbKv(@TempDir Path tempDir) throws Exception { assertThat(rocksDBKv.multiGet(Arrays.asList(key, key2))).containsExactly(null, val2); } } + + @Test + void testClosePreservesUnflushedStateByDefault(@TempDir Path tempDir) throws Exception { + byte[] key = new byte[] {1}; + byte[] value = new byte[] {2}; + + RocksDBKv rocksDBKv = buildRocksDBKv(tempDir.toFile()); + rocksDBKv.put(key, value); + rocksDBKv.close(); + + try (RocksDBKv reopened = buildRocksDBKv(tempDir.toFile())) { + assertThat(reopened.get(key)).isEqualTo(value); + } + } + + @Test + void testDiscardCloseDoesNotPersistUnflushedState(@TempDir Path tempDir) throws Exception { + byte[] key = new byte[] {1}; + + RocksDBKv rocksDBKv = buildRocksDBKv(tempDir.toFile()); + rocksDBKv.put(key, new byte[] {2}); + rocksDBKv.close(KvCloseMode.DISCARD_UNPERSISTED_STATE); + + try (RocksDBKv reopened = buildRocksDBKv(tempDir.toFile())) { + assertThat(reopened.get(key)).isNull(); + } + } + + @Test + void testDiscardOptionFailureStillClosesNativeResources(@TempDir Path tempDir) + throws Exception { + RocksDBKv rocksDBKv = spy(buildRocksDBKv(tempDir.toFile())); + doThrow(new RocksDBException("expected")).when(rocksDBKv).setAvoidFlushDuringShutdown(); + + assertThatThrownBy(() -> rocksDBKv.close(KvCloseMode.DISCARD_UNPERSISTED_STATE)) + .isInstanceOf(RocksDBException.class) + .hasMessage("expected"); + assertThatThrownBy(rocksDBKv::checkIfRocksDBClosed) + .isInstanceOf(FlussRuntimeException.class); + + // Reopening the same path also proves the native DB handle was released. + try (RocksDBKv ignored = buildRocksDBKv(tempDir.toFile())) {} + } + + private RocksDBKv buildRocksDBKv(File instanceBasePath) throws Exception { + RocksDBResourceContainer rocksDBResourceContainer = + new RocksDBResourceContainer(new Configuration(), instanceBasePath); + RocksDBKvBuilder rocksDBKvBuilder = + new RocksDBKvBuilder( + instanceBasePath, + rocksDBResourceContainer, + rocksDBResourceContainer.getColumnOptions()); + return rocksDBKvBuilder.build(); + } } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java index 761cf051997..533a032b7bd 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java @@ -169,6 +169,7 @@ void testConfigurationOptionsFromConfig() throws Exception { assertThat(dbOptions.keepLogFileNum()).isEqualTo(10); assertThat(dbOptions.maxLogFileSize()).isEqualTo(2 * SizeUnit.MB); assertThat(dbOptions.statistics()).isNotNull(); + assertThat(dbOptions.avoidFlushDuringShutdown()).isFalse(); ColumnFamilyOptions columnOptions = optionsContainer.getColumnOptions(); assertThat(columnOptions.compactionStyle()).isEqualTo(CompactionStyle.LEVEL); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownTest.java b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownTest.java index 6cb07f77f2e..4b2375845d0 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownTest.java @@ -20,6 +20,7 @@ import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.exception.FlussRuntimeException; +import org.apache.fluss.server.kv.KvCloseMode; import org.junit.jupiter.api.Test; @@ -60,6 +61,7 @@ void testQuiescesTabletUsersBeforeClosingTabletManagers() throws Exception { assertThat(server.getShutdownEvents()) .containsExactly( "rpc-started", "rpc-completed", "replica-manager", "tablet-managers"); + assertThat(server.getKvCloseMode()).isEqualTo(KvCloseMode.DISCARD_UNPERSISTED_STATE); } finally { server.completeRpcShutdown(); shutdownExecutor.shutdownNow(); @@ -81,6 +83,7 @@ void testClosesTabletManagersWhenReplicaManagerShutdownFails() throws Exception assertThat(server.getShutdownEvents()) .containsExactly( "rpc-started", "rpc-completed", "replica-manager", "tablet-managers"); + assertThat(server.getKvCloseMode()).isEqualTo(KvCloseMode.DISCARD_UNPERSISTED_STATE); } private static final class TestingTabletServer extends TabletServer { @@ -88,6 +91,7 @@ private static final class TestingTabletServer extends TabletServer { private final CountDownLatch rpcShutdownStarted = new CountDownLatch(1); private final List shutdownEvents = new CopyOnWriteArrayList<>(); private boolean failReplicaManagerShutdown; + private KvCloseMode kvCloseMode; private TestingTabletServer(Configuration conf) { super(conf); @@ -110,8 +114,9 @@ void shutdownReplicaManager() { } @Override - void shutdownTabletManagers() { + void shutdownKvManager(KvCloseMode closeMode) { shutdownEvents.add("tablet-managers"); + kvCloseMode = closeMode; } private boolean awaitRpcShutdownStarted() throws InterruptedException { @@ -129,5 +134,9 @@ private void failReplicaManagerShutdown() { private List getShutdownEvents() { return shutdownEvents; } + + private KvCloseMode getKvCloseMode() { + return kvCloseMode; + } } }