diff --git a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java index e497ea0703..d88615673d 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java @@ -2681,9 +2681,22 @@ public class ConfigOptions { + "so that clients begin throttling before the storage engine is forced to throttle itself. " + "The gap between this value and the storage trigger forms the throttle ramp-up window."); - // ------------------------------------------------------------------------ + // ConfigOptions for KV lazy open + + public static final ConfigOption KV_LAZY_OPEN_ENABLED = + key("kv.lazy-open.enabled") + .booleanType() + .defaultValue(false) + .withDescription("Whether to enable KvTablet lazy open."); + + public static final ConfigOption KV_LAZY_OPEN_IDLE_TIMEOUT = + key("kv.lazy-open.idle-timeout") + .durationType() + .defaultValue(Duration.ofHours(24)) + .withDescription( + "Idle time before an open KvTablet is eligible for release back to lazy state."); + // ConfigOptions for metrics - // ------------------------------------------------------------------------ public static final ConfigOption> METRICS_REPORTERS = key("metrics.reporters") .stringType() @@ -2704,9 +2717,7 @@ public class ConfigOptions { + "the CoordinatorServer) it is advisable to use a port range " + "like 9250-9260."); - // ------------------------------------------------------------------------ // ConfigOptions for prometheus push gateway reporter - // ------------------------------------------------------------------------ public static final ConfigOption METRICS_REPORTER_PROMETHEUS_PUSHGATEWAY_HOST_URL = key("metrics.reporter.prometheus-push.host-url") .stringType() @@ -2771,9 +2782,7 @@ public class ConfigOptions { + "The value is automatically redacted when the configuration " + "is logged or displayed."); - // ------------------------------------------------------------------------ // ConfigOptions for jmx reporter - // ------------------------------------------------------------------------ public static final ConfigOption METRICS_REPORTER_JMX_HOST = key("metrics.reporter.jmx.port") .stringType() @@ -2786,9 +2795,7 @@ public class ConfigOptions { + "the CoordinatorServer) it is advisable to use a port range " + "like 9990-9999."); - // ------------------------------------------------------------------------ // ConfigOptions for influxdb reporter - // ------------------------------------------------------------------------ public static final ConfigOption METRICS_REPORTER_INFLUXDB_VERSION = key("metrics.reporter.influxdb.version") .stringType() @@ -2830,9 +2837,7 @@ public class ConfigOptions { .defaultValue(Duration.ofSeconds(10)) .withDescription("The interval of reporting metrics to InfluxDB."); - // ------------------------------------------------------------------------ // ConfigOptions for lakehouse storage - // ------------------------------------------------------------------------ public static final ConfigOption DATALAKE_ENABLED = key("datalake.enabled") .booleanType() @@ -2853,9 +2858,7 @@ public class ConfigOptions { "The datalake format used by Fluss as lakehouse storage. Currently, supported formats are Paimon, Iceberg, Hudi, and Lance. " + "In the future, more kinds of data lake format will be supported, such as DeltaLake."); - // ------------------------------------------------------------------------ // ConfigOptions for tiering service - // ------------------------------------------------------------------------ public static final ConfigOption LAKE_TIERING_AUTO_EXPIRE_SNAPSHOT = key("lake.tiering.auto-expire-snapshot") @@ -2877,9 +2880,7 @@ public class ConfigOptions { + "If not configured and the tiering service runs in a Flink job, Fluss uses " + "Flink's IO temporary directories with a 'fluss' child directory."); - // ------------------------------------------------------------------------ // ConfigOptions for fluss kafka - // ------------------------------------------------------------------------ public static final ConfigOption KAFKA_ENABLED = key("kafka.enabled") .booleanType() diff --git a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java index d1ab0ae6c9..a9a2f3d8ff 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java +++ b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java @@ -175,6 +175,9 @@ public class MetricNames { "preWriteBufferTruncateAsDuplicatedPerSecond"; public static final String KV_PRE_WRITE_BUFFER_TRUNCATE_AS_ERROR_RATE = "preWriteBufferTruncateAsErrorPerSecond"; + public static final String KV_TABLET_OPEN_COUNT = "kvTabletOpenCount"; + public static final String KV_TABLET_LAZY_COUNT = "kvTabletLazyCount"; + public static final String KV_TABLET_FAILED_COUNT = "kvTabletFailedCount"; // -------------------------------------------------------------------------------------------- // RocksDB metrics diff --git a/fluss-filesystems/fluss-fs-oss/pom.xml b/fluss-filesystems/fluss-fs-oss/pom.xml index 7ebd7c4dd1..4cd62233b5 100644 --- a/fluss-filesystems/fluss-fs-oss/pom.xml +++ b/fluss-filesystems/fluss-fs-oss/pom.xml @@ -124,6 +124,7 @@ jar + test-jar @@ -237,18 +238,6 @@ - - - org.apache.maven.plugins - maven-jar-plugin - - - - test-jar - - - - diff --git a/fluss-flink/fluss-flink-common/pom.xml b/fluss-flink/fluss-flink-common/pom.xml index dcb81c741b..7c0918dce6 100644 --- a/fluss-flink/fluss-flink-common/pom.xml +++ b/fluss-flink/fluss-flink-common/pom.xml @@ -117,13 +117,6 @@ test - - org.apache.flink - flink-connector-base - ${flink.minor.version} - test - - org.apache.flink flink-connector-base diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java index 4376405934..7df75a2427 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java @@ -214,7 +214,10 @@ void testBoundedPkTableEmitsKvBatchSplits() throws Throwable { @Test void testBoundedPkTableEmitsSnapshotSplitsByDefault() throws Throwable { - createTable(DEFAULT_TABLE_PATH, DEFAULT_PK_TABLE_DESCRIPTOR); + long tableId = createTable(DEFAULT_TABLE_PATH, DEFAULT_PK_TABLE_DESCRIPTOR); + for (int bucketId = 0; bucketId < DEFAULT_BUCKET_NUM; bucketId++) { + FLUSS_CLUSTER_EXTENSION.waitUntilAllReplicaReady(new TableBucket(tableId, bucketId)); + } int numSubtasks = DEFAULT_BUCKET_NUM; try (MockSplitEnumeratorContext context = new MockSplitEnumeratorContext<>(numSubtasks)) { 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 29eedc5eb6..53b894a51d 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 @@ -83,6 +83,7 @@ import java.util.Optional; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Semaphore; import java.util.function.Predicate; import static org.apache.fluss.utils.Preconditions.checkState; @@ -129,6 +130,10 @@ public static RateLimiter getDefaultRateLimiter() { return DEFAULT_RATE_LIMITER; } + private static final int LAZY_OPEN_MAX_CONCURRENT_OPENS = 10; + + static final long RELEASE_DRAIN_TIMEOUT_MS = 5_000; + private final LogManager logManager; private final LocalDiskManager localDiskManager; @@ -174,6 +179,9 @@ public static RateLimiter getDefaultRateLimiter() { private final KvFlushScheduler kvFlushScheduler; + private final boolean lazyOpenEnabled; + private final @Nullable Semaphore openSemaphore; + private volatile boolean isShutdown = false; private KvManager( @@ -194,6 +202,8 @@ private KvManager( this.zkClient = zkClient; this.clock = clock; this.remoteKvDir = FlussPaths.remoteKvDir(conf); + this.lazyOpenEnabled = conf.get(ConfigOptions.KV_LAZY_OPEN_ENABLED); + this.openSemaphore = lazyOpenEnabled ? new Semaphore(LAZY_OPEN_MAX_CONCURRENT_OPENS) : null; this.remoteFileSystem = remoteKvDir.getFileSystem(); this.serverMetricGroup = tabletServerMetricGroup; this.sharedRocksDBRateLimiter = createSharedRateLimiter(conf); @@ -347,6 +357,37 @@ public static KvManager create( clock); } + /** Returns whether KV tablets open their RocksDB resources on first access. */ + public boolean isLazyOpenEnabled() { + return lazyOpenEnabled; + } + + /** Releases idle tablets without changing their registry ownership or deleting local SSTs. */ + public void releaseIdleTablets(long idleTimeoutMs, long nowMs) { + for (KvTablet tablet : currentKvs.values()) { + try { + if (tablet.canRelease(idleTimeoutMs, nowMs)) { + tablet.releaseKv(); + } + } catch (Exception e) { + LOG.warn("Failed to release idle KV tablet {}", tablet.getTableBucket(), e); + } + } + } + + /** Creates an unregistered tablet that opens RocksDB on first access. */ + public KvTablet createLazyTablet(PhysicalTablePath path, TableBucket bucket, LogTablet log) { + KvTablet tablet = + new KvTablet( + path, + bucket, + log, + getTabletDir(log.getDataDir(), path, bucket), + serverMetricGroup); + tablet.getLifecycle().configureTiming(clock, openSemaphore, RELEASE_DRAIN_TIMEOUT_MS); + return tablet; + } + /** * Returns the shared block cache usage in bytes, or 0 if shared cache is disabled. * @@ -514,11 +555,8 @@ private void closeKvTablet(KvTablet kvTablet, KvCloseMode closeMode) { try { kvTablet.close(closeMode); } catch (Exception e) { - LOG.warn( - "Exception while closing kv tablet {} with mode {}.", - kvTablet.getTableBucket(), - closeMode, - e); + throw new KvStorageException( + "Failed to close KV tablet " + kvTablet.getTableBucket(), e); } } @@ -546,10 +584,56 @@ public KvTablet getOrCreateKv( ArrowCompressionInfo arrowCompressionInfo, @Nullable Runnable flushCompleteListener) throws Exception { + return createKv( + tablePath, + tableBucket, + logTablet, + kvFormat, + schemaGetter, + tableConfig, + arrowCompressionInfo, + flushCompleteListener, + true); + } + + /** Creates a tablet without replacing the registered lazy tablet. */ + public KvTablet createKvTabletUnregistered( + PhysicalTablePath tablePath, + TableBucket tableBucket, + LogTablet logTablet, + KvFormat kvFormat, + SchemaGetter schemaGetter, + TableConfig tableConfig, + ArrowCompressionInfo arrowCompressionInfo, + @Nullable Runnable flushCompleteListener) + throws Exception { + return createKv( + tablePath, + tableBucket, + logTablet, + kvFormat, + schemaGetter, + tableConfig, + arrowCompressionInfo, + flushCompleteListener, + false); + } + + private KvTablet createKv( + PhysicalTablePath tablePath, + TableBucket tableBucket, + LogTablet logTablet, + KvFormat kvFormat, + SchemaGetter schemaGetter, + TableConfig tableConfig, + ArrowCompressionInfo arrowCompressionInfo, + @Nullable Runnable flushCompleteListener, + boolean register) + throws Exception { return inKvLock( tableBucket, () -> { - if (currentKvs.containsKey(tableBucket)) { + if (register && currentKvs.containsKey(tableBucket)) { return currentKvs.get(tableBucket); } @@ -586,7 +670,9 @@ public KvTablet getOrCreateKv( autoIncrementManager, clock, tableConfig); - currentKvs.put(tableBucket, tablet); + if (register) { + currentKvs.put(tableBucket, tablet); + } LOG.info( "Created kv tablet for bucket {} in dir {}.", @@ -621,20 +707,28 @@ public Optional getKv(TableBucket tableBucket) { } public void dropKv(TableBucket tableBucket) { - inKvLock( - tableBucket, - () -> { - doDropKv(tableBucket); - return null; - }); - } - - private void doDropKv(TableBucket tableBucket) { - KvTablet dropKvTablet = currentKvs.remove(tableBucket); + KvTablet lazyTablet = + inKvLock( + tableBucket, + () -> { + KvTablet current = currentKvs.get(tableBucket); + if (current != null && !current.isLazyMode()) { + doDropKv(tableBucket, current); + return null; + } + return current; + }); + // Lazy open callbacks also need the bucket lock. Retain registry ownership while + // waiting, but do not hold that lock across lazy lifecycle cleanup. + doDropKv(tableBucket, lazyTablet); + } + + private void doDropKv(TableBucket tableBucket, @Nullable KvTablet dropKvTablet) { if (dropKvTablet != null) { TablePath tablePath = dropKvTablet.getTablePath(); try { dropKvTablet.drop(); + inKvLock(tableBucket, () -> currentKvs.remove(tableBucket, dropKvTablet)); if (dropKvTablet.getPartitionName() == null) { LOG.info( "Deleted kv bucket {} for table {} in file path {}.", @@ -675,7 +769,27 @@ public KvTablet loadKv( physicalTablePath, tableBucket, schemaGetter, - flushCompleteListener)); + flushCompleteListener, + true)); + } + + /** Loads local state without replacing the registered lazy tablet. */ + public KvTablet loadKvUnregistered( + File tabletDir, SchemaGetter schemaGetter, @Nullable Runnable flushCompleteListener) + throws Exception { + Tuple2 pathAndBucket = FlussPaths.parseTabletDir(tabletDir); + PhysicalTablePath physicalTablePath = pathAndBucket.f0; + TableBucket tableBucket = pathAndBucket.f1; + return inKvLock( + tableBucket, + () -> + doLoadKv( + tabletDir, + physicalTablePath, + tableBucket, + schemaGetter, + flushCompleteListener, + false)); } private KvTablet doLoadKv( @@ -683,10 +797,11 @@ private KvTablet doLoadKv( PhysicalTablePath physicalTablePath, TableBucket tableBucket, SchemaGetter schemaGetter, - @Nullable Runnable flushCompleteListener) + @Nullable Runnable flushCompleteListener, + boolean register) throws Exception { KvTablet currentKv = currentKvs.get(tableBucket); - if (currentKv != null) { + if (register && currentKv != null) { throw new IllegalStateException( String.format( "Duplicate kv tablet directories for bucket %s are found in both %s and %s. " @@ -748,11 +863,27 @@ private KvTablet doLoadKv( autoIncrementManager, clock, tableConfig); - currentKvs.put(tableBucket, kvTablet); + if (register) { + currentKvs.put(tableBucket, kvTablet); + } return kvTablet; } + /** Register a KvTablet (e.g. a lazy sentinel) into the {@code currentKvs} registry. */ + public void registerKv(TableBucket tableBucket, KvTablet kvTablet) { + inKvLock( + tableBucket, + () -> { + if (currentKvs.containsKey(tableBucket)) { + throw new IllegalStateException( + "KvTablet already registered for " + tableBucket); + } + currentKvs.put(tableBucket, kvTablet); + return null; + }); + } + public void deleteRemoteKvSnapshot( PhysicalTablePath physicalTablePath, TableBucket tableBucket) { FsPath remoteKvTabletDir = 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 e6bea68f7a..af8352b422 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 @@ -124,7 +124,7 @@ public final class KvTablet { private static final int MAX_FLUSH_RETRY_BACKOFF_SHIFT = 64 - Long.numberOfLeadingZeros(MAX_FLUSH_RETRY_DELAY_MS / MIN_FLUSH_RETRY_DELAY_MS); - private static final long ROW_COUNT_DISABLED = -1; + public static final long ROW_COUNT_DISABLED = -1; // Retain recent historical KV state within this WAL offset distance of lake progress. // TODO: Consider time-based retention after lake coverage is confirmed. @@ -137,11 +137,38 @@ public final class KvTablet { */ private static final int TARGET_ENTRIES_PER_NATIVE_WRITE = 500; + /** Pins an opened lazy tablet until closed; eager tablets use a no-op guard. */ + public static final class Guard implements AutoCloseable { + private final KvTablet tablet; + private final @Nullable KvTabletLazyLifecycle lifecycle; + private boolean released; + + Guard(KvTablet tablet, @Nullable KvTabletLazyLifecycle lifecycle) { + this.tablet = tablet; + this.lifecycle = lifecycle; + } + + /** Returns the opened tablet, valid until this guard is closed. */ + public KvTablet getTablet() { + return tablet; + } + + @Override + public void close() { + if (!released && lifecycle != null) { + released = true; + lifecycle.releasePin(); + } + } + } + + private final Guard eagerGuard = new Guard(this, null); + private final PhysicalTablePath physicalPath; private final TableBucket tableBucket; private final boolean historicalPartition; - private final LogTablet logTablet; + final LogTablet logTablet; private final File kvTabletDir; private final long writeBatchSize; @@ -149,7 +176,7 @@ public final class KvTablet { private final KvPreWriteBuffer kvPreWriteBuffer; private final KvStateAccessor kvStateAccessor; private final KvWriteProcessor kvWriteProcessor; - private final TabletServerMetricGroup serverMetricGroup; + final TabletServerMetricGroup serverMetricGroup; private final KvFlushScheduler kvFlushScheduler; private final boolean closeFlushScheduler; @@ -192,6 +219,12 @@ public final class KvTablet { @GuardedBy("kvLock") private volatile boolean isClosed = false; + @Nullable private final KvTabletLazyLifecycle lifecycle; + + // Keep the opened tablet intact: queued flushes and snapshot callbacks retain its locks + // and state until close has drained all RocksDB resource leases. + @Nullable private volatile KvTablet openedTablet; + private KvTablet( PhysicalTablePath physicalPath, TableBucket tableBucket, @@ -219,6 +252,7 @@ private KvTablet( @Nullable RowTtlTimestampProvider rowTtlTimestampProvider, Clock clock, boolean rowTtlEnabled) { + this.lifecycle = null; this.physicalPath = physicalPath; this.tableBucket = tableBucket; this.historicalPartition = @@ -269,6 +303,89 @@ private KvTablet( : 0L; } + KvTablet( + PhysicalTablePath physicalPath, + TableBucket tableBucket, + LogTablet logTablet, + File kvTabletDir, + @Nullable TabletServerMetricGroup serverMetricGroup) { + this.physicalPath = physicalPath; + this.tableBucket = tableBucket; + this.logTablet = logTablet; + this.serverMetricGroup = serverMetricGroup; + this.historicalPartition = + HISTORICAL_PARTITION_VALUE.equals(physicalPath.getPartitionName()); + this.kvTabletDir = kvTabletDir; + this.writeBatchSize = 0; + this.rocksDBKv = null; + this.kvPreWriteBuffer = null; + this.kvStateAccessor = null; + this.kvWriteProcessor = null; + this.kvFlushScheduler = null; + this.closeFlushScheduler = false; + this.kvValueLayout = null; + this.stateValueEncoder = null; + this.historicalCleanupOffset = new AtomicLong(); + this.rowTtlTimestampProvider = null; + this.rowTtlEnabled = false; + this.autoIncrementManager = null; + this.rocksDBStatistics = null; + this.lifecycle = new KvTabletLazyLifecycle(this); + } + + /** Returns the lazy lifecycle, or null for an eagerly opened tablet. */ + @Nullable + public KvTabletLazyLifecycle getLifecycle() { + return lifecycle; + } + + void installRocksDB(KvTablet source) { + checkState(source.lifecycle == null, "Only an eager tablet can be installed."); + synchronized (historicalCleanupOffset) { + if (historicalPartition) { + source.advanceHistoricalCleanupOffset(historicalCleanupOffset.get()); + } + this.openedTablet = source; + } + } + + void detachRocksDB(KvCloseMode closeMode) { + KvTablet opened = openedTablet; + if (opened != null) { + try { + opened.close(closeMode); + } catch (Exception e) { + throw new KvStorageException("Failed to close KV tablet " + tableBucket, e); + } + openedTablet = null; + } + } + + KvTablet requireOpenedTablet() { + return checkNotNull(openedTablet, "RocksDB is not open for %s", tableBucket); + } + + long currentRowCount() { + KvTablet opened = openedTablet; + return lifecycle == null || opened == null ? rowCount : opened.rowCount; + } + + boolean hasActiveResourceLeases() { + KvTablet opened = openedTablet; + return opened != null && opened.rocksDBKv.getResourceGuard().getLeaseCount() > 0; + } + + void deleteLocalDirectory() { + File dir = getKvTabletDir(); + if (dir != null) { + try { + FileUtils.deleteDirectory(dir); + } catch (IOException e) { + throw new KvStorageException("Failed to delete KV directory " + dir, e); + } + } + } + /** * Creates a kv tablet with a dedicated {@link KvFlushScheduler} that is closed together with * the tablet. Production code must use {@link #create(PhysicalTablePath, TableBucket, @@ -630,6 +747,11 @@ public TablePath getTablePath() { } public long getAutoIncrementCacheSize() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? 0 : opened.getAutoIncrementCacheSize(); + } + return autoIncrementManager.getAutoIncrementCacheSize(); } @@ -648,6 +770,11 @@ public File getKvTabletDir() { /** Returns the total size in bytes of the live RocksDB SST files. */ public long liveSstFilesSize() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? 0L : opened.liveSstFilesSize(); + } + return rocksDBKv.liveSstFilesSize(); } @@ -658,6 +785,11 @@ public long liveSstFilesSize() { */ @Nullable public RocksDBStatistics getRocksDBStatistics() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? null : opened.getRocksDBStatistics(); + } + return rocksDBStatistics; } @@ -673,6 +805,13 @@ void setRowCount(long rowCount) { // row_count is volatile, so it's safe to read without lock public long getRowCount() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + if (opened != null) { + return opened.getRowCount(); + } + } + if (rowCount == ROW_COUNT_DISABLED) { if (rowTtlEnabled) { throw new InvalidTableException( @@ -701,6 +840,10 @@ public long getRowCount() { */ @GuardedBy("kvLock") public TabletState getTabletState() { + if (lifecycle != null) { + return requireOpenedTablet().getTabletState(); + } + return new TabletState( flushedLogOffset, rowCount == ROW_COUNT_DISABLED ? null : rowCount, @@ -858,6 +1001,16 @@ long localLogEndOffset() { } public void requestFlush(long exclusiveUpToLogOffset, FatalErrorHandler fatalErrorHandler) { + if (lifecycle != null) { + KvTablet opened = openedTablet; + if (opened != null) { + // Enqueue flushes even while release drains an existing writer. The opened + // tablet serializes requests with close and ignores requests after closing. + opened.requestFlush(exclusiveUpToLogOffset, fatalErrorHandler); + } + return; + } + asyncFatalErrorHandler = fatalErrorHandler; inWriteLock(kvLock, () -> requestFlushInternal(exclusiveUpToLogOffset)); } @@ -869,11 +1022,28 @@ public void setFlushCompleteListener(@Nullable Runnable flushCompleteListener) { } public long getFlushedLogOffset() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? flushedLogOffset : opened.getFlushedLogOffset(); + } + return flushedLogOffset; } /** Advances the exclusive historical cleanup offset without allowing it to move backwards. */ public boolean advanceHistoricalCleanupOffset(long cleanupOffset) { + if (lifecycle != null) { + checkState(historicalPartition, "%s is not a historical KV tablet", tableBucket); + checkArgument(cleanupOffset >= 0L, "Historical cleanup offset must be non-negative."); + synchronized (historicalCleanupOffset) { + long previous = historicalCleanupOffset.getAndAccumulate(cleanupOffset, Math::max); + KvTablet opened = openedTablet; + return opened == null + ? cleanupOffset > previous + : opened.advanceHistoricalCleanupOffset(cleanupOffset); + } + } + checkState(historicalPartition, "%s is not a historical KV tablet", tableBucket); checkArgument(cleanupOffset >= 0L, "Historical cleanup offset must be non-negative."); long previousCleanupOffset = @@ -883,6 +1053,13 @@ public boolean advanceHistoricalCleanupOffset(long cleanupOffset) { /** Returns the current exclusive cleanup offset for a historical overlay. */ public long getHistoricalCleanupOffset() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null + ? historicalCleanupOffset.get() + : opened.getHistoricalCleanupOffset(); + } + checkState(historicalPartition, "%s is not a historical KV tablet", tableBucket); return historicalCleanupOffset.get(); } @@ -1198,6 +1375,20 @@ void putToPreWriteBuffer( * tablet. */ public Executor getGuardedExecutor() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return runnable -> { + Guard guard = lifecycle.tryAcquireExistingGuard(); + if (guard != null) { + try (Guard ignored = guard) { + if (opened != null && openedTablet == opened) { + opened.getGuardedExecutor().execute(runnable); + } + } + } + }; + } + return runnable -> inWriteLock(kvLock, runnable::run); } @@ -1371,6 +1562,10 @@ public void close() throws Exception { } public void close(KvCloseMode closeMode) throws Exception { + if (lifecycle != null) { + lifecycle.close(closeMode, false); + return; + } LOG.info( "Close kv tablet {} for table {} with mode {}.", tableBucket, @@ -1401,6 +1596,11 @@ public void close(KvCloseMode closeMode) throws Exception { /** Completely delete the kv directory and all contents form the file system with no delay. */ public void drop() throws Exception { + if (lifecycle != null) { + lifecycle.close(KvCloseMode.DISCARD_UNPERSISTED_STATE, true); + return; + } + inWriteLock( kvLock, () -> { @@ -1417,6 +1617,15 @@ public RocksIncrementalSnapshot createIncrementalSnapshot( KvSnapshotDataUploader kvSnapshotDataUploader, long lastCompletedSnapshotId, Counter remoteKvCopyBytes) { + if (lifecycle != null) { + return requireOpenedTablet() + .createIncrementalSnapshot( + uploadedSstFiles, + kvSnapshotDataUploader, + lastCompletedSnapshotId, + remoteKvCopyBytes); + } + return new RocksIncrementalSnapshot( uploadedSstFiles, rocksDBKv.getDb(), @@ -1430,17 +1639,32 @@ public RocksIncrementalSnapshot createIncrementalSnapshot( // only for testing. @VisibleForTesting KvPreWriteBuffer getKvPreWriteBuffer() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? null : opened.getKvPreWriteBuffer(); + } + return kvPreWriteBuffer; } // only for testing. @VisibleForTesting public RocksDBKv getRocksDBKv() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? null : opened.getRocksDBKv(); + } + return rocksDBKv; } /** Returns the recent normalized backpressure pressure in {@code [0, 1)}. */ public float currentPressure() { + if (lifecycle != null) { + KvTablet opened = openedTablet; + return opened == null ? 0F : opened.currentPressure(); + } + return rocksDBKv.currentPressure(); } @@ -1520,4 +1744,30 @@ enum FlushState { RUNNING, STORAGE_BLOCKED } + + /** Returns true if this tablet is in lazy mode (as opposed to eager/traditional mode). */ + public boolean isLazyMode() { + return lifecycle != null; + } + + /** Returns true if this tablet is in lazy mode and currently OPEN (RocksDB loaded). */ + public boolean isLazyOpen() { + return lifecycle != null && lifecycle.isOpen(); + } + + /** + * Opens and pins this tablet. Call outside replica locks to avoid blocking leadership changes. + */ + public Guard acquireGuard() { + return lifecycle != null ? lifecycle.acquireGuard() : eagerGuard; + } + + /** Pre-check for idle release eligibility. */ + public boolean canRelease(long closeIdleIntervalMs, long nowMs) { + return lifecycle != null && lifecycle.canRelease(closeIdleIntervalMs, nowMs); + } + + public boolean releaseKv() { + return lifecycle != null && lifecycle.releaseKv(); + } } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTabletLazyLifecycle.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTabletLazyLifecycle.java new file mode 100644 index 0000000000..6bcad6ebeb --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTabletLazyLifecycle.java @@ -0,0 +1,549 @@ +/* + * 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; + +import org.apache.fluss.exception.KvStorageException; +import org.apache.fluss.server.metrics.group.TabletServerMetricGroup; +import org.apache.fluss.utils.ExponentialBackoff; +import org.apache.fluss.utils.clock.Clock; +import org.apache.fluss.utils.clock.SystemClock; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; +import javax.annotation.concurrent.GuardedBy; + +import java.util.concurrent.CancellationException; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.Condition; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Consumer; +import java.util.function.IntSupplier; + +/** + * Opens KV resources on demand and releases idle tablets. Close/drop fence new requests and wait + * for pins, callbacks and candidate cleanup before relinquishing native resources or local paths. + */ +public final class KvTabletLazyLifecycle { + private static final Logger LOG = LoggerFactory.getLogger(KvTabletLazyLifecycle.class); + + /** Coarsen access timestamp updates to reduce volatile writes on hot path. */ + private static final long ACCESS_TIMESTAMP_GRANULARITY_MS = 1000; + + private static final long OPEN_TIMEOUT_MS = 300_000; + + /** Internal RocksDB lazy lifecycle states. */ + enum LazyState { + LAZY, + + OPENING, + + OPEN, + + RELEASING, + + FAILED, + /** Closing: new requests are fenced while outstanding work drains. */ + CLOSING, + /** Terminal state; local data is deleted only on drop. */ + CLOSED + } + + /** Creates an eager tablet, recovering local SST/WAL or a snapshot as needed. */ + @FunctionalInterface + public interface OpenCallback { + /** Performs the open and returns the opened RocksDB KvTablet instance. */ + KvTablet doOpen(boolean hasLocalData) throws Exception; + } + + private final KvTablet tablet; + + private volatile LazyState lazyState = LazyState.LAZY; + + private final ReentrantLock lazyStateLock = new ReentrantLock(); + + private final Condition lazyStateChanged = lazyStateLock.newCondition(); + + /** Active pin count — prevents release while operations are in-flight. */ + private final AtomicInteger activePins = new AtomicInteger(0); + + /** Semaphore for throttling concurrent open operations across all tablets. */ + private @Nullable Semaphore openSemaphore; + + private long releaseDrainTimeoutMs; + private final ExponentialBackoff failedBackoff = new ExponentialBackoff(5_000, 2, 300_000, 0); + private Clock clock = SystemClock.getInstance(); + + private long failedTimestamp; + + private int failureCount; + private @Nullable Throwable lastFailureCause; + + /** Whether local RocksDB data directory exists from a previous open. */ + @GuardedBy("lazyStateLock") + private boolean hasLocalData; + + /** Access timestamp for idle release decisions (coarsened to reduce volatile writes). */ + private volatile long lastAccessTimestamp; + + private @Nullable IntSupplier leaderEpochSupplier; + + private @Nullable IntSupplier bucketEpochSupplier; + + private @Nullable OpenCallback openCallback; + + private @Nullable Consumer commitCallback; + private @Nullable Consumer releaseCallback; + + /** + * Includes open callbacks, candidate cleanup and idle release cleanup. Guarded by the state + * lock. + */ + private boolean operationInFlight; + + /** Retains ownership if cleaning up an uncommitted open fails. */ + private @Nullable KvTablet failedOpenTablet; + + /** Guarded by the termination monitor; prevents repeated drops from deleting a reused path. */ + private boolean localDirectoryDeleted; + + KvTabletLazyLifecycle(KvTablet tablet) { + this.tablet = tablet; + updateStateGauge(LazyState.LAZY, 1); + } + + /** Configures recovery before this sentinel is registered or accessed. */ + public void configure( + IntSupplier leaderEpoch, + IntSupplier bucketEpoch, + OpenCallback open, + Consumer committed, + Consumer release) { + leaderEpochSupplier = leaderEpoch; + bucketEpochSupplier = bucketEpoch; + openCallback = open; + commitCallback = committed; + releaseCallback = release; + } + + void configureTiming(Clock clock, Semaphore semaphore, long drainTimeout) { + this.clock = clock; + this.openSemaphore = semaphore; + this.releaseDrainTimeoutMs = drainTimeout; + } + + /** Initializes row count from snapshot metadata without opening RocksDB. */ + public void initCachedRowCount(long rowCount) { + tablet.setRowCount(rowCount); + } + + boolean isOpen() { + return lazyState == LazyState.OPEN; + } + + LazyState getLazyState() { + return lazyState; + } + + /** Returns the row count retained across idle release. */ + public long getCachedRowCount() { + return tablet.currentRowCount(); + } + + long getLastAccessTimestamp() { + return lastAccessTimestamp; + } + + int getActivePins() { + return activePins.get(); + } + + /** + * Opens and pins the tablet. Call outside replica locks to avoid blocking leadership changes. + */ + KvTablet.Guard acquireGuard() { + if (openCallback == null) { + throw new IllegalStateException("Lazy tablet recovery is not configured."); + } + + for (int attempt = 0; ; attempt++) { + KvTablet.Guard guard = tryAcquireExistingGuard(); + if (guard != null) { + touchAccessTimestamp(); + return guard; + } + if (attempt >= 3) { + throw new KvStorageException("Failed to pin KV tablet " + tablet.getTableBucket()); + } + ensureOpen(); + } + } + + /** + * Pins maintenance without refreshing idle time. Recheck after incrementing so release cannot + * miss a newly admitted request. + */ + @Nullable + KvTablet.Guard tryAcquireExistingGuard() { + if (lazyState == LazyState.OPEN) { + activePins.incrementAndGet(); + if (lazyState == LazyState.OPEN) { + return new KvTablet.Guard(tablet.requireOpenedTablet(), this); + } + releasePin(); + } + return null; + } + + void releasePin() { + if (activePins.decrementAndGet() == 0) { + if (lazyState != LazyState.OPEN) { + lazyStateLock.lock(); + try { + lazyStateChanged.signalAll(); + } finally { + lazyStateLock.unlock(); + } + } + } + } + + private void touchAccessTimestamp() { + long now = clock.milliseconds(); + if (now - lastAccessTimestamp > ACCESS_TIMESTAMP_GRANULARITY_MS) { + lastAccessTimestamp = now; + } + } + + /** Counts stable states only; transitional states have no gauge. */ + private void updateStateGauge(LazyState state, int delta) { + TabletServerMetricGroup metrics = tablet.serverMetricGroup; + if (metrics == null) { + return; + } + if (state == LazyState.LAZY) { + metrics.kvTabletLazyCount().addAndGet(delta); + } else if (state == LazyState.OPEN) { + metrics.kvTabletOpenCount().addAndGet(delta); + } else if (state == LazyState.FAILED) { + metrics.kvTabletFailedCount().addAndGet(delta); + } + } + + /** Waits for ongoing work or starts a single open after failure backoff. */ + private void ensureOpen() { + lazyStateLock.lock(); + try { + while (true) { + switch (lazyState) { + case OPEN: + return; + + case OPENING: + case RELEASING: + if (!lazyStateChanged.await(OPEN_TIMEOUT_MS, TimeUnit.MILLISECONDS)) { + throw new KvStorageException( + "Timed out opening KV tablet " + tablet.getTableBucket()); + } + continue; + + case FAILED: + long elapsed = clock.milliseconds() - failedTimestamp; + long cooldown = failedBackoff.backoff(failureCount); + if (elapsed < cooldown) { + throw new KvStorageException( + "KvTablet open failed for " + + tablet.getTableBucket() + + ", remaining cooldown: " + + (cooldown - elapsed) + + " ms", + lastFailureCause); + } + // Fall through to LAZY — cooldown expired, retry open. + + case LAZY: + transitionTo(LazyState.OPENING); + operationInFlight = true; + + int myLeaderEpoch = leaderEpochSupplier.getAsInt(); + int myBucketEpoch = bucketEpochSupplier.getAsInt(); + boolean localDataExists = hasLocalData; + lazyStateLock.unlock(); + try { + doSlowOpen(myLeaderEpoch, myBucketEpoch, localDataExists); + } catch (Throwable t) { + // doSlowOpen completes candidate cleanup before publishing failure. + if (t instanceof RuntimeException) { + throw (RuntimeException) t; + } + throw new KvStorageException( + "KvTablet open failed for " + tablet.getTableBucket(), t); + } + return; + + case CLOSING: + case CLOSED: + throw new KvStorageException( + "KvTablet is closed for " + tablet.getTableBucket()); + + default: + throw new IllegalStateException("Unexpected state: " + lazyState); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new KvStorageException( + "Interrupted waiting for KvTablet open: " + tablet.getTableBucket(), e); + } finally { + if (lazyStateLock.isHeldByCurrentThread()) { + lazyStateLock.unlock(); + } + } + } + + private void doSlowOpen(int myLeaderEpoch, int myBucketEpoch, boolean localDataExists) + throws Exception { + KvTablet candidate = null; + boolean committed = false; + boolean semaphoreAcquired = false; + Throwable failure = null; + boolean cleanupFailed = false; + try { + if (openSemaphore != null) { + semaphoreAcquired = + openSemaphore.tryAcquire(OPEN_TIMEOUT_MS, TimeUnit.MILLISECONDS); + if (!semaphoreAcquired) { + throw new KvStorageException( + "KvTablet open semaphore timeout for " + tablet.getTableBucket()); + } + } + candidate = openCallback.doOpen(localDataExists); + lazyStateLock.lock(); + try { + int leaderEpoch = leaderEpochSupplier.getAsInt(); + int bucketEpoch = bucketEpochSupplier.getAsInt(); + if (lazyState != LazyState.OPENING + || leaderEpoch != myLeaderEpoch + || bucketEpoch != myBucketEpoch) { + throw new CancellationException( + "open result fenced by concurrent epoch or lifecycle change"); + } + tablet.installRocksDB(candidate); + committed = true; + transitionTo(LazyState.OPEN); + failureCount = 0; + lastAccessTimestamp = clock.milliseconds(); + + lazyStateChanged.signalAll(); + } finally { + lazyStateLock.unlock(); + } + try { + commitCallback.accept(tablet); + } catch (Exception e) { + LOG.warn("Post-open commit callback failed for {}", tablet.getTableBucket(), e); + } + } catch (Exception | Error e) { + failure = e; + throw e; + } finally { + try { + if (!committed && candidate != null) { + candidate.close(); + candidate.deleteLocalDirectory(); + } + } catch (Exception | Error e) { + cleanupFailed = true; + if (failure != null) { + failure.addSuppressed(e); + } else { + failure = e; + throw e; + } + } finally { + if (semaphoreAcquired) { + openSemaphore.release(); + } + lazyStateLock.lock(); + try { + if (cleanupFailed) { + failedOpenTablet = candidate; + + transitionTo(LazyState.CLOSING); + } else if (!committed && lazyState != LazyState.CLOSING) { + transitionTo(LazyState.FAILED); + failedTimestamp = clock.milliseconds(); + failureCount = + failure instanceof CancellationException ? 0 : failureCount + 1; + lastFailureCause = failure; + } + operationInFlight = false; + lazyStateChanged.signalAll(); + } finally { + lazyStateLock.unlock(); + } + } + } + } + + /** Must be called while holding {@code lazyStateLock}. */ + private void transitionTo(LazyState state) { + updateStateGauge(lazyState, -1); + updateStateGauge(state, 1); + lazyState = state; + } + + boolean canRelease(long closeIdleIntervalMs, long nowMs) { + return lazyState == LazyState.OPEN + && activePins.get() == 0 + && nowMs - lastAccessTimestamp >= closeIdleIntervalMs; + } + + /** + * Releases idle native resources; returns false if data, leases or callbacks prevent release. + */ + boolean releaseKv() { + lazyStateLock.lock(); + try { + if (lazyState != LazyState.OPEN || operationInFlight) { + return false; + } + + transitionTo(LazyState.RELEASING); + operationInFlight = true; + if (!drainPins() + || lazyState == LazyState.CLOSING + || tablet.getFlushedLogOffset() < tablet.logTablet.localLogEndOffset() + || tablet.hasActiveResourceLeases()) { + if (lazyState == LazyState.RELEASING) { + transitionTo(LazyState.OPEN); + } + operationInFlight = false; + lazyStateChanged.signalAll(); + return false; + } + tablet.setRowCount(tablet.currentRowCount()); + tablet.setFlushedLogOffset(tablet.getFlushedLogOffset()); + } finally { + lazyStateLock.unlock(); + } + + boolean releaseSucceeded = false; + boolean nativeCloseStarted = false; + try { + if (releaseCallback != null) { + releaseCallback.accept(tablet); + } + nativeCloseStarted = true; + tablet.detachRocksDB(KvCloseMode.PRESERVE_LOCAL_STATE); + releaseSucceeded = true; + } catch (Exception e) { + LOG.warn("Failed to release KV tablet {}", tablet.getTableBucket(), e); + } finally { + lazyStateLock.lock(); + try { + if (lazyState == LazyState.RELEASING) { + if (releaseSucceeded) { + transitionTo(LazyState.LAZY); + hasLocalData = true; + + } else if (!nativeCloseStarted) { + transitionTo(LazyState.OPEN); + } else { + transitionTo(LazyState.CLOSING); + } + } + operationInFlight = false; + lazyStateChanged.signalAll(); + } finally { + lazyStateLock.unlock(); + } + } + return releaseSucceeded; + } + + /** Serializes close callers without holding the state lock during callbacks or native I/O. */ + synchronized void close(KvCloseMode closeMode, boolean deleteDirectory) { + if (lazyState == LazyState.CLOSED) { + if (deleteDirectory && !localDirectoryDeleted) { + tablet.deleteLocalDirectory(); + localDirectoryDeleted = true; + } + return; + } + lazyStateLock.lock(); + try { + + transitionTo(LazyState.CLOSING); + lazyStateChanged.signalAll(); + // A request timeout cannot transfer ownership of native resources or the local path. + while (operationInFlight || activePins.get() > 0) { + lazyStateChanged.awaitUninterruptibly(); + } + } finally { + lazyStateLock.unlock(); + } + + if (releaseCallback != null) { + releaseCallback.accept(tablet); + } + if (failedOpenTablet != null) { + try { + failedOpenTablet.close(); + failedOpenTablet.deleteLocalDirectory(); + failedOpenTablet = null; + } catch (Exception e) { + throw new KvStorageException( + "Failed to clean up fenced open for " + tablet.getTableBucket(), e); + } + } + tablet.detachRocksDB(closeMode); + if (deleteDirectory) { + tablet.deleteLocalDirectory(); + localDirectoryDeleted = true; + } + lazyStateLock.lock(); + try { + transitionTo(LazyState.CLOSED); + lastFailureCause = null; + lazyStateChanged.signalAll(); + } finally { + lazyStateLock.unlock(); + } + } + + private boolean drainPins() { + long deadline = clock.milliseconds() + releaseDrainTimeoutMs; + while (activePins.get() > 0) { + long remaining = deadline - clock.milliseconds(); + if (remaining <= 0) { + return false; + } + try { + lazyStateChanged.await(remaining, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } + } + return true; + } +} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java index dac5562c10..ad50eac314 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java @@ -34,6 +34,7 @@ import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.LongSupplier; import static org.apache.fluss.utils.Preconditions.checkNotNull; @@ -72,6 +73,9 @@ public class TabletServerMetricGroup extends AbstractMetricGroup { private final Histogram kvFlushLatencyHistogram; private final Counter kvTruncateAsDuplicatedCount; private final Counter kvTruncateAsErrorCount; + private final AtomicInteger kvTabletOpenCount = new AtomicInteger(0); + private final AtomicInteger kvTabletLazyCount = new AtomicInteger(0); + private final AtomicInteger kvTabletFailedCount = new AtomicInteger(0); // aggregated replica metrics private final Counter isrShrinks; @@ -140,6 +144,9 @@ public TabletServerMetricGroup( meter( MetricNames.KV_PRE_WRITE_BUFFER_TRUNCATE_AS_ERROR_RATE, new MeterView(kvTruncateAsErrorCount)); + gauge(MetricNames.KV_TABLET_OPEN_COUNT, kvTabletOpenCount::get); + gauge(MetricNames.KV_TABLET_LAZY_COUNT, kvTabletLazyCount::get); + gauge(MetricNames.KV_TABLET_FAILED_COUNT, kvTabletFailedCount::get); // replica metrics isrExpands = new SimpleCounter(); @@ -301,6 +308,21 @@ public Counter kvTruncateAsErrorCount() { return kvTruncateAsErrorCount; } + /** Returns the gauge-backing counter for currently open lazy KvTablets. */ + public AtomicInteger kvTabletOpenCount() { + return kvTabletOpenCount; + } + + /** Returns the gauge-backing counter for KvTablets in LAZY (released) state. */ + public AtomicInteger kvTabletLazyCount() { + return kvTabletLazyCount; + } + + /** Returns the gauge-backing counter for KvTablets in FAILED state. */ + public AtomicInteger kvTabletFailedCount() { + return kvTabletFailedCount; + } + public Counter isrShrinks() { return isrShrinks; } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java index e3eda80a82..826186216b 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/Replica.java @@ -63,6 +63,7 @@ import org.apache.fluss.server.kv.KvRecoverHelper; import org.apache.fluss.server.kv.KvStateLookupResult; import org.apache.fluss.server.kv.KvTablet; +import org.apache.fluss.server.kv.KvTabletLazyLifecycle; import org.apache.fluss.server.kv.RemoteLogFetcher; import org.apache.fluss.server.kv.autoinc.AutoIncIDRange; import org.apache.fluss.server.kv.historical.HistoricalValueLookup; @@ -112,6 +113,7 @@ import org.apache.fluss.types.RowType; import org.apache.fluss.utils.ByteArraySlice; import org.apache.fluss.utils.CloseableRegistry; +import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.FlussPaths; import org.apache.fluss.utils.IOUtils; import org.apache.fluss.utils.clock.Clock; @@ -229,7 +231,7 @@ public final class Replica { // null if table without pk or haven't become leader private volatile @Nullable KvTablet kvTablet; private volatile @Nullable CloseableRegistry closeableRegistryForKv; - private @Nullable PeriodicSnapshotManager kvSnapshotManager; + private volatile @Nullable PeriodicSnapshotManager kvSnapshotManager; /** * Server-wide {@link ScannerManager}. Active sessions for this bucket are closed in {@link @@ -323,6 +325,15 @@ private void registerMetrics() { logicalStorageMetrics.gauge(MetricNames.LOCAL_STORAGE_KV_SIZE, this::logicalStorageKvSize); } + private void stopSnapshotManager() { + PeriodicSnapshotManager mgr = this.kvSnapshotManager; + if (mgr != null) { + this.kvSnapshotManager = null; + closeableRegistryForKv.unregisterCloseable(mgr); + mgr.close(); + } + } + public long logicalStorageLogSize() { if (isLeader()) { return logTablet.logicalStorageSize(); @@ -340,8 +351,12 @@ public long logicalStorageKvSize() { KvTablet currentKvTablet = kvTablet; return currentKvTablet == null ? 0L : currentKvTablet.liveSstFilesSize(); } - checkNotNull(kvSnapshotManager, "kvSnapshotManager is null"); - return kvSnapshotManager.getSnapshotSize(); + PeriodicSnapshotManager snapshots = kvSnapshotManager; + if (snapshots != null) { + return snapshots.getSnapshotSize(); + } + KvTablet kv = kvTablet; + return kv != null && kv.isLazyMode() ? getLatestKvSnapshotSize() : 0L; } else { // follower doesn't need to report the logical storage size. return 0L; @@ -771,12 +786,21 @@ private void createKv() { e); } + checkNotNull(kvManager); + if (kvManager.isLazyOpenEnabled()) { + createKvLazy(); + } else { + createKvEager(); + } + } + + private void createKvEager() { // init kv tablet and get the snapshot it uses to init if have any Optional snapshotUsed = Optional.empty(); Exception lastError = null; for (int i = 1; i <= INIT_KV_TABLET_MAX_RETRY_TIMES; i++) { try { - snapshotUsed = initKvTablet(); + snapshotUsed = initKvTabletEager(); lastError = null; break; } catch (Exception e) { @@ -802,24 +826,23 @@ private void createKv() { lastError); } if (!isHistoricalPartition()) { - startPeriodicKvSnapshot(snapshotUsed.orElse(null)); + startPeriodicKvSnapshot(kvTablet, snapshotUsed.orElse(null)); } } private void dropKv() { - // Release scanner leases first; otherwise resourceGuard.close() inside kvTablet.close() - // blocks waiting for them. Runs under leaderIsrUpdateLock(W), so no concurrent register. scannerManager.closeScannersForBucket(tableBucket); - + // close any closeable registry for kv if (closeableRegistry.unregisterCloseable(closeableRegistryForKv)) { IOUtils.closeQuietly(closeableRegistryForKv); } - if (kvTablet != null) { - bucketMetricGroup.unregisterRocksDBStatistics(); - - checkNotNull(kvManager); - kvManager.dropKv(tableBucket); - kvTablet = null; + KvTablet kv = this.kvTablet; + if (kv != null) { + if (!kv.isLazyMode()) { + bucketMetricGroup.unregisterRocksDBStatistics(); + } + checkNotNull(kvManager).dropKv(tableBucket); + this.kvTablet = null; } } @@ -855,36 +878,105 @@ private void onKvFlushComplete() { * * @return the snapshot used to init kv tablet, empty if no any snapshot. */ - private Optional initKvTablet() { + private Optional initKvTabletEager() { + InitKvResult result = initKvTablet(null); + kvManager.registerKv(tableBucket, result.tablet); + kvTablet = result.tablet; + if (kvTablet.getRocksDBStatistics() != null) { + bucketMetricGroup.registerRocksDBStatistics(kvTablet.getRocksDBStatistics()); + } + return result.snapshotUsed; + } + + private void createKvLazy() { + TableConfig tableConfig = getTableConfig(); + KvTablet sentinel = kvManager.createLazyTablet(physicalPath, tableBucket, logTablet); + KvTabletLazyLifecycle lifecycle = sentinel.getLifecycle(); + AtomicReference snapshotRef = new AtomicReference<>(); + lifecycle.configure( + () -> leaderEpoch, + () -> bucketEpoch, + hasLocal -> { + InitKvResult result = initKvTablet(hasLocal ? sentinel : null); + snapshotRef.set(result.snapshotUsed.orElse(null)); + return result.tablet; + }, + kv -> { + if (!isHistoricalPartition()) { + startPeriodicKvSnapshot(kv, snapshotRef.getAndSet(null)); + } + if (kv.getRocksDBStatistics() != null) { + bucketMetricGroup.registerRocksDBStatistics(kv.getRocksDBStatistics()); + } + }, + kv -> { + stopSnapshotManager(); + bucketMetricGroup.unregisterRocksDBStatistics(); + }); + + // Initialize cached row count from snapshot metadata if available (no RocksDB needed). + // This is best-effort — if ZK is unavailable, default to 0 rather than failing + // the entire lazy sentinel creation. + long initialRowCount; + if (isHistoricalPartition() || !supportsExactRowCount(tableConfig)) { + initialRowCount = -1L; + } else { + try { + Long snapshotRowCount = + getLatestSnapshot(tableBucket) + .map(CompletedSnapshot::getRowCount) + .orElse(null); + initialRowCount = snapshotRowCount != null ? snapshotRowCount : 0L; + } catch (Exception e) { + LOG.warn( + "Failed to read snapshot row count for {}, defaulting to 0", + tableBucket, + e); + initialRowCount = 0L; + } + } + lifecycle.initCachedRowCount(initialRowCount); + + // Register sentinel in KvManager registry + kvManager.registerKv(tableBucket, sentinel); + + this.kvTablet = sentinel; + } + + private InitKvResult initKvTablet(@Nullable KvTablet released) { checkNotNull(kvManager); TableConfig tableConfig = getTableConfig(); long startTime = clock.milliseconds(); LOG.info("Start to init kv tablet for {} of table {}.", tableBucket, physicalPath); - // todo: we may need to handle the following cases: - // case1: no kv files in local, restore from remote snapshot; and apply - // the log; - // case2: kv files in local - // - if no remote snapshot, restore from local and apply the log known to the local - // files. - // - have snapshot, if the known offset to the local files is much less than(maybe - // some value configured) - // the remote snapshot; restore from remote snapshot; - - // currently for simplicity, we'll always download the snapshot files and restore from - // the snapshots as kv files won't exist in our current implementation for - // when replica become follower, we'll always delete the kv files. - + // Released tablets retain local SSTs and the corresponding WAL replay position. + // Initial opens and failed local recovery use the latest durable snapshot instead. // get the offset from which, we should restore from. default is 0 - long restoreStartOffset = isHistoricalPartition() ? historicalRecoveryStartOffset() : 0; + long restoreStartOffset = + released == null && isHistoricalPartition() ? historicalRecoveryStartOffset() : 0; // Lake is the durable base for local historical KV state. Historical replicas therefore // never restore a normal KV snapshot, even if one exists from older code. Optional optCompletedSnapshot = - isHistoricalPartition() ? Optional.empty() : getLatestSnapshot(tableBucket); + released != null || isHistoricalPartition() + ? Optional.empty() + : getLatestSnapshot(tableBucket); + KvTablet kvTablet = null; try { Long rowCount; AutoIncIDRange autoIncIDRange; - if (optCompletedSnapshot.isPresent()) { + if (released != null) { + File directory = released.getKvTabletDir(); + if (!new File(directory, RocksDBKvBuilder.DB_INSTANCE_DIR_STRING).isDirectory()) { + throw new IOException("Missing local RocksDB directory for " + tableBucket); + } + kvTablet = + kvManager.loadKvUnregistered( + directory, schemaGetter, this::onKvFlushComplete); + restoreStartOffset = released.getFlushedLogOffset(); + long cachedRows = released.getLifecycle().getCachedRowCount(); + rowCount = cachedRows == KvTablet.ROW_COUNT_DISABLED ? null : cachedRows; + autoIncIDRange = null; + } else if (optCompletedSnapshot.isPresent()) { LOG.info( "Use snapshot {} to restore kv tablet for {} of table {}.", optCompletedSnapshot.get(), @@ -899,7 +991,9 @@ private Optional initKvTablet() { downloadKvSnapshots(completedSnapshot, tabletDir.toPath()); // as we have downloaded kv files into the tablet dir, now, we can load it - kvTablet = kvManager.loadKv(tabletDir, schemaGetter, this::onKvFlushComplete); + kvTablet = + kvManager.loadKvUnregistered( + tabletDir, schemaGetter, this::onKvFlushComplete); checkNotNull(kvTablet, "kv tablet should not be null."); restoreStartOffset = completedSnapshot.getLogOffset(); @@ -913,13 +1007,9 @@ private Optional initKvTablet() { tableBucket, physicalPath); - // Rebuild state in a fresh directory, as with snapshot recovery. - // Recovery retries can reuse the tablet opened by the first attempt. - if (kvTablet == null) { - kvManager.createTabletDir(logTablet.getDataDir(), physicalPath, tableBucket); - } + kvManager.createTabletDir(logTablet.getDataDir(), physicalPath, tableBucket); kvTablet = - kvManager.getOrCreateKv( + kvManager.createKvTabletUnregistered( physicalPath, tableBucket, logTablet, @@ -946,10 +1036,28 @@ private Optional initKvTablet() { logTablet.updateMinRetainOffset(restoreStartOffset); if (isHistoricalPartition()) { checkNotNull(kvTablet, "kv tablet should not be null.") - .advanceHistoricalCleanupOffset(restoreStartOffset); + .advanceHistoricalCleanupOffset( + released == null + ? restoreStartOffset + : Math.max(logTablet.getLakeLogEndOffset(), 0)); } - recoverKvTablet(restoreStartOffset, rowCount, autoIncIDRange); + recoverKvTablet(kvTablet, restoreStartOffset, rowCount, autoIncIDRange); } catch (Exception e) { + try { + if (kvTablet != null) { + kvTablet.close(); + } + if (released != null) { + FileUtils.deleteDirectory(released.getKvTabletDir()); + } + } catch (Exception cleanupError) { + e.addSuppressed(cleanupError); + throw new KvStorageException("KV recovery cleanup failed for " + tableBucket, e); + } + if (released != null) { + LOG.warn("Local KV recovery failed for {}; restoring snapshot", tableBucket, e); + return initKvTablet(null); + } throw new KvStorageException( String.format( "Fail to init kv tablet for %s of table %s.", @@ -963,14 +1071,17 @@ private Optional initKvTablet() { tableBucket, endTime - startTime); - if (kvTablet != null) { - // Register RocksDB statistics now that the kv tablet is fully initialized. - if (kvTablet.getRocksDBStatistics() != null) { - bucketMetricGroup.registerRocksDBStatistics(kvTablet.getRocksDBStatistics()); - } - } + return new InitKvResult(kvTablet, optCompletedSnapshot); + } + + private static class InitKvResult { + final KvTablet tablet; + final Optional snapshotUsed; - return optCompletedSnapshot; + InitKvResult(KvTablet tablet, Optional snapshotUsed) { + this.tablet = tablet; + this.snapshotUsed = snapshotUsed; + } } private void downloadKvSnapshots(CompletedSnapshot completedSnapshot, Path kvTabletDir) @@ -1037,6 +1148,7 @@ private Optional getLatestSnapshot(TableBucket tableBucket) { } private void recoverKvTablet( + KvTablet kvTablet, long startRecoverLogOffset, @Nullable Long rowCount, @Nullable AutoIncIDRange autoIncIDRange) { @@ -1113,7 +1225,8 @@ private long historicalRecoveryStartOffset() { return recoveryStartOffset; } - private void startPeriodicKvSnapshot(@Nullable CompletedSnapshot completedSnapshot) { + private void startPeriodicKvSnapshot( + KvTablet kvTablet, @Nullable CompletedSnapshot completedSnapshot) { checkNotNull(kvTablet); KvTabletSnapshotTarget kvTabletSnapshotTarget; try { @@ -1198,11 +1311,12 @@ private void startPeriodicKvSnapshot(@Nullable CompletedSnapshot completedSnapsh } public long getLatestKvSnapshotSize() { - if (kvSnapshotManager == null) { - return 0L; - } else { - return kvSnapshotManager.getSnapshotSize(); + PeriodicSnapshotManager snapshots = kvSnapshotManager; + if (snapshots != null) { + return snapshots.getSnapshotSize(); } + CompletedSnapshot latest = getLatestSnapshot(tableBucket).orElse(null); + return latest == null ? 0L : latest.getSnapshotSize(); } public long getLeaderEndOffsetSnapshot() { @@ -1277,9 +1391,9 @@ public LogAppendInfo putRecordsToLeader( MergeMode mergeMode, int requiredAcks) throws Exception { - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + requiredAcks, + kv -> { if (!isLeader()) { throw new NotLeaderOrFollowerException( String.format( @@ -1292,7 +1406,6 @@ public LogAppendInfo putRecordsToLeader( } validateInSyncReplicaSize(requiredAcks); - KvTablet kv = this.kvTablet; checkNotNull( kv, "KvTablet for the replica to put kv records shouldn't be null."); LogAppendInfo logAppendInfo; @@ -1324,11 +1437,10 @@ public LogAppendInfo putHistoricalRecordsToLeader( int requiredAcks) throws Exception { LocalValueLookupResult localLookupResult = - inReadLock( - leaderIsrUpdateLock, - () -> { + withGuardedLeaderKv( + requiredAcks, + kv -> { validateHistoricalWrite(expectedLeaderEpoch, requiredAcks); - KvTablet kv = this.kvTablet; checkNotNull( kv, "KvTablet for the historical replica shouldn't be null."); return kv.probeLocalPreviousValues( @@ -1339,11 +1451,10 @@ public LogAppendInfo putHistoricalRecordsToLeader( HistoricalValueLookup historicalValueLookup = localLookupResult.createValueLookup(lakeLookup); - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + requiredAcks, + kv -> { validateHistoricalWrite(expectedLeaderEpoch, requiredAcks); - KvTablet kv = this.kvTablet; checkNotNull(kv, "KvTablet for the historical replica shouldn't be null."); LogAppendInfo appendInfo = kv.putHistoricalAsLeader( @@ -1418,9 +1529,9 @@ private void validateHistoricalWrite(int expectedLeaderEpoch, int requiredAcks) /** Looks up keys from the local historical KV state of the leader replica. */ public List lookupHistoricalLocal( String originalPartitionName, List keys) throws Exception { - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { if (!isLeader()) { throw new NotLeaderOrFollowerException( String.format( @@ -1432,7 +1543,6 @@ public List lookupHistoricalLocal( "Historical lookup request must target a historical partition."); } - KvTablet kv = this.kvTablet; checkNotNull(kv, "KvTablet for the historical replica shouldn't be null."); List results = new ArrayList<>(keys.size()); for (byte[] key : keys) { @@ -1659,9 +1769,9 @@ public List lookups(List keys) { throw new NonPrimaryKeyTableException( "the primary key table not exists for " + tableBucket); } - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { try { if (!isLeader()) { throw new NotLeaderOrFollowerException( @@ -1671,7 +1781,7 @@ public List lookups(List keys) { } checkNotNull( kvTablet, "KvTablet for the replica to get key shouldn't be null."); - return kvTablet.multiGet(keys); + return kv.multiGet(keys); } catch (IOException e) { String errorMsg = String.format( @@ -1694,9 +1804,9 @@ public List lookupsFromBufferOrKv(List keys) { throw new NonPrimaryKeyTableException( "the primary key table not exists for " + tableBucket); } - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { try { if (!isLeader()) { throw new NotLeaderOrFollowerException( @@ -1706,7 +1816,7 @@ public List lookupsFromBufferOrKv(List keys) { } checkNotNull( kvTablet, "KvTablet for the replica to get key shouldn't be null."); - return kvTablet.multiGetFromBufferOrKv(keys); + return kv.multiGetFromBufferOrKv(keys); } catch (IOException e) { String errorMsg = String.format( @@ -1724,9 +1834,9 @@ public List prefixLookup(byte[] prefixKey) { "Try to do prefix lookup on a non primary key table: " + getTablePath()); } - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { try { if (!isLeader()) { throw new NotLeaderOrFollowerException( @@ -1736,7 +1846,7 @@ public List prefixLookup(byte[] prefixKey) { } checkNotNull( kvTablet, "KvTablet for the replica to get key shouldn't be null."); - return kvTablet.prefixLookup(prefixKey); + return kv.prefixLookup(prefixKey); } catch (IOException e) { String errorMsg = String.format( @@ -1754,9 +1864,9 @@ public DefaultValueRecordBatch limitKvScan(int limit) { "the primary key table not exists for " + tableBucket); } - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { try { if (!isLeader()) { throw new NotLeaderOrFollowerException( @@ -1767,7 +1877,7 @@ public DefaultValueRecordBatch limitKvScan(int limit) { checkNotNull( kvTablet, "KvTablet for the replica to limit scan shouldn't be null."); - List values = kvTablet.limitScan(limit); + List values = kv.limitScan(limit); DefaultValueRecordBatch.Builder builder = DefaultValueRecordBatch.builder(); for (ByteArraySlice value : values) { builder.append(value.array(), value.offset(), value.length()); @@ -1803,9 +1913,9 @@ public OpenScanResult openScan( "the primary key table not exists for " + tableBucket); } - return inReadLock( - leaderIsrUpdateLock, - () -> { + return withGuardedLeaderKv( + 0, + kv -> { if (!isLeader()) { throw new NotLeaderOrFollowerException( String.format( @@ -1814,8 +1924,7 @@ public OpenScanResult openScan( } checkNotNull( kvTablet, "KvTablet for the replica to open scan shouldn't be null."); - OpenScanResult result = - kvTablet.openScan(scannerId, limit, initialAccessTimeMs); + OpenScanResult result = kv.openScan(scannerId, limit, initialAccessTimeMs); ScannerContext context = result.getContext(); if (context == null) { return result; @@ -1830,6 +1939,59 @@ public OpenScanResult openScan( }); } + private T withGuardedLeaderKv( + int requiredAcks, GuardedKvOperation action) throws E { + while (true) { + KvTablet kv = + inReadLock( + leaderIsrUpdateLock, + () -> { + ensureLeaderForKvAccess(); + validateInSyncReplicaSize(requiredAcks); + return checkNotNull(kvTablet, "Leader KV tablet is missing."); + }); + boolean openingKv = kv.isLazyMode() && !kv.isLazyOpen(); + try (KvTablet.Guard guard = kv.acquireGuard()) { + // A leadership transition drains pins while holding the write lock. Never wait + // for that lock with a pin: release it and retry after the transition instead. + if (leaderIsrUpdateLock.readLock().tryLock()) { + try { + ensureLeaderForKvAccess(); + if (this.kvTablet != kv) { + throw new NotLeaderOrFollowerException( + "KV tablet changed while opening " + tableBucket); + } + validateInSyncReplicaSize(requiredAcks); + return action.apply(guard.getTablet()); + } finally { + leaderIsrUpdateLock.readLock().unlock(); + } + } + } finally { + // Recovery can leave WAL entries above the old high watermark in the buffer. + // Publish progress only after releasing the pin and replica read lock. + if (openingKv && kv.isLazyOpen()) { + onKvFlushComplete(); + } + } + inReadLock(leaderIsrUpdateLock, this::ensureLeaderForKvAccess); + } + } + + private void ensureLeaderForKvAccess() { + if (!isLeader()) { + throw new NotLeaderOrFollowerException( + String.format( + "Leader not local for bucket %s on tabletServer %d", + tableBucket, localTabletServerId)); + } + } + + @FunctionalInterface + private interface GuardedKvOperation { + T apply(KvTablet tablet) throws E; + } + public LogRecords limitLogScan(int limit) { return inReadLock( leaderIsrUpdateLock, @@ -1908,6 +2070,8 @@ public long getRowCount() { if (kv != null) { // return materialized row count for primary key table return kv.getRowCount(); + } else if (isKvTable()) { + return 0L; } else { // return log row count for non-primary key table return logTablet.getRowCount(); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java index 41ce69fce7..c11d2fe418 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/ReplicaManager.java @@ -125,9 +125,11 @@ import org.apache.fluss.server.zk.data.LeaderAndIsr; import org.apache.fluss.server.zk.data.lake.LakeTableSnapshot; import org.apache.fluss.utils.ByteArraySlice; +import org.apache.fluss.utils.ExecutorUtils; import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.FlussPaths; import org.apache.fluss.utils.clock.Clock; +import org.apache.fluss.utils.concurrent.ExecutorThreadFactory; import org.apache.fluss.utils.concurrent.FutureUtils; import org.apache.fluss.utils.concurrent.Scheduler; @@ -155,6 +157,9 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.Lock; @@ -183,6 +188,8 @@ public class ReplicaManager implements ServerReconfigurable { * safe. */ private static final short PUT_KV_VERSION_WITH_STORAGE_BACKPRESSURE = 2; + // ---- Idle release internal constants ---- + private static final long IDLE_RELEASE_CHECK_INTERVAL_MS = 60_000; private final Configuration conf; private final Scheduler scheduler; @@ -244,6 +251,8 @@ public class ReplicaManager implements ServerReconfigurable { private final ScannerManager scannerManager; private final HistoricalPartitionManager historicalPartitionManager; + private final long kvIdleTimeoutMs; + private @Nullable ScheduledExecutorService kvIdleReleaseScheduler; public ReplicaManager( Configuration conf, @@ -378,6 +387,13 @@ public ReplicaManager( dataDirVolumeBytes, scheduler); + this.kvIdleTimeoutMs = conf.get(ConfigOptions.KV_LAZY_OPEN_IDLE_TIMEOUT).toMillis(); + if (kvManager != null && kvManager.isLazyOpenEnabled()) { + checkArgument(kvIdleTimeoutMs > 0, "kv.lazy-open.idle-timeout must be positive"); + kvIdleReleaseScheduler = + Executors.newSingleThreadScheduledExecutor( + new ExecutorThreadFactory("kv-idle-release")); + } registerMetrics(); } @@ -395,6 +411,13 @@ public void startup() { // Start periodic disk usage monitoring (initial + periodic sampling) localDiskManager.startDiskUsageMonitor(scheduler); + if (kvIdleReleaseScheduler != null) { + kvIdleReleaseScheduler.scheduleWithFixedDelay( + () -> kvManager.releaseIdleTablets(kvIdleTimeoutMs, clock.milliseconds()), + IDLE_RELEASE_CHECK_INTERVAL_MS, + IDLE_RELEASE_CHECK_INTERVAL_MS, + TimeUnit.MILLISECONDS); + } } public RemoteLogManager getRemoteLogManager() { @@ -2677,6 +2700,10 @@ public Replica getReplica() { public static final class OfflineReplica implements HostedReplica {} public void shutdown() throws InterruptedException { + if (kvIdleReleaseScheduler != null) { + ExecutorUtils.gracefulShutdown(5, TimeUnit.SECONDS, kvIdleReleaseScheduler); + } + // Close the resources for snapshot kv kvSnapshotResource.close(); historicalPartitionManager.close(); 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 c153bfbfaf..6491739455 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 @@ -22,6 +22,7 @@ import org.apache.fluss.config.MemorySize; import org.apache.fluss.config.TableConfig; import org.apache.fluss.exception.FlussRuntimeException; +import org.apache.fluss.exception.KvStorageException; import org.apache.fluss.metadata.KvFormat; import org.apache.fluss.metadata.LogFormat; import org.apache.fluss.metadata.PhysicalTablePath; @@ -50,16 +51,20 @@ import org.apache.fluss.server.zk.ZooKeeperExtension; import org.apache.fluss.server.zk.data.TableRegistration; import org.apache.fluss.testutils.common.AllCallbackWrapper; +import org.apache.fluss.testutils.common.CommonTestUtils; import org.apache.fluss.types.RowType; import org.apache.fluss.utils.ByteArraySlice; import org.apache.fluss.utils.FlussPaths; +import org.apache.fluss.utils.clock.ManualClock; import org.apache.fluss.utils.clock.SystemClock; import org.apache.fluss.utils.concurrent.FlussScheduler; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInfo; import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.api.io.TempDir; import org.junit.jupiter.params.ParameterizedTest; @@ -80,17 +85,25 @@ import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.concurrent.Callable; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executor; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; +import java.util.concurrent.FutureTask; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; +import java.util.function.IntSupplier; import static org.apache.fluss.compression.ArrowCompressionInfo.DEFAULT_COMPRESSION; import static org.apache.fluss.record.TestData.DATA1_SCHEMA_PK; import static org.apache.fluss.record.TestData.DATA2_SCHEMA; import static org.apache.fluss.server.kv.KvTabletTestUtils.flushAndWait; +import static org.apache.fluss.testutils.DataTestUtils.compactedRow; import static org.apache.fluss.testutils.common.CommonTestUtils.waitUntil; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -133,10 +146,13 @@ static void baseBeforeAll() { } @BeforeEach - void setup() throws Exception { + void setup(TestInfo testInfo) throws Exception { conf = new Configuration(); conf.setString(ConfigOptions.DATA_DIR, tempDir.getAbsolutePath()); conf.set(ConfigOptions.TABLET_SERVER_ID, 1); + conf.set( + ConfigOptions.KV_LAZY_OPEN_ENABLED, + testInfo.getTestClass().orElse(KvManagerTest.class) == LazyLifecycleTest.class); String dbName = "db1"; tablePath1 = TablePath.of(dbName, "t1"); @@ -1128,4 +1144,543 @@ private void unblock() { unblock.countDown(); } } + + private static FutureTask start(Callable action) { + FutureTask task = new FutureTask<>(action); + new Thread(task, "lazy-lifecycle-test").start(); + return task; + } + + private static class BlockingCallback implements AutoCloseable { + private final CountDownLatch entered = new CountDownLatch(1); + private final CountDownLatch finish = new CountDownLatch(1); + + void run() { + entered.countDown(); + try { + finish.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + } + + void awaitEntered() throws Exception { + assertThat(entered.await(5, TimeUnit.SECONDS)).isTrue(); + } + + @Override + public void close() { + finish.countDown(); + } + } + + @Nested + class LazyLifecycleTest { + private PhysicalTablePath physicalPath; + private TableBucket tableBucket; + private LogTablet logTablet; + private ManualClock clock; + private KvTablet sentinel; + private KvTabletLazyLifecycle lifecycle; + private int baseLazy; + private int baseOpen; + private int baseFailed; + private KvTabletLazyLifecycle.OpenCallback openCallback; + private Consumer commitCallback = kv -> {}; + private Consumer releaseCallback = kv -> {}; + private IntSupplier leaderEpochSupplier = () -> 1; + private IntSupplier bucketEpochSupplier = () -> 1; + + @BeforeEach + void setupLifecycle() throws Exception { + assertThat(kvManager.isLazyOpenEnabled()).isTrue(); + TablePath tablePath = TablePath.of("db1", "t_lazy"); + physicalPath = + PhysicalTablePath.of( + tablePath.getDatabaseName(), tablePath.getTableName(), null); + tableBucket = new TableBucket(20001L, 0); + logTablet = + logManager.getOrCreateLog( + tempDir, physicalPath, tableBucket, LogFormat.ARROW, 1, true); + TabletServerMetricGroup metrics = TestingMetricGroups.TABLET_SERVER_METRICS; + baseLazy = metrics.kvTabletLazyCount().get(); + baseOpen = metrics.kvTabletOpenCount().get(); + baseFailed = metrics.kvTabletFailedCount().get(); + clock = new ManualClock(System.currentTimeMillis()); + openCallback = local -> openTablet(sentinel); + sentinel = createSentinel(); + lifecycle = configureLifecycle(sentinel); + } + + private KvTablet createSentinel() { + KvTablet tablet = kvManager.createLazyTablet(physicalPath, tableBucket, logTablet); + kvManager.registerKv(tableBucket, tablet); + return tablet; + } + + private KvTabletLazyLifecycle configureLifecycle(KvTablet sentinel) { + KvTabletLazyLifecycle lifecycle = sentinel.getLifecycle(); + lifecycle.configureTiming(clock, new java.util.concurrent.Semaphore(10), 100); + lifecycle.configure( + () -> leaderEpochSupplier.getAsInt(), + () -> bucketEpochSupplier.getAsInt(), + local -> openCallback.doOpen(local), + kv -> commitCallback.accept(kv), + kv -> releaseCallback.accept(kv)); + return lifecycle; + } + + private KvTablet openTablet(KvTablet owner) throws Exception { + return kvManager.createKvTabletUnregistered( + physicalPath, + owner.getTableBucket(), + owner.logTablet, + KvFormat.COMPACTED, + new TestingSchemaGetter(new SchemaInfo(DATA1_SCHEMA_PK, schemaId)), + new TableConfig(new Configuration()), + DEFAULT_COMPRESSION, + null); + } + + private KvTablet secondTablet() throws Exception { + TableBucket bucket = new TableBucket(tableBucket.getTableId(), 1); + LogTablet log = + logManager.getOrCreateLog( + tempDir, physicalPath, bucket, LogFormat.ARROW, 1, true); + KvTablet second = kvManager.createLazyTablet(physicalPath, bucket, log); + second.getLifecycle() + .configureTiming(clock, new java.util.concurrent.Semaphore(10), 100); + second.getLifecycle() + .configure(() -> 1, () -> 1, local -> openTablet(second), kv -> {}, kv -> {}); + kvManager.registerKv(bucket, second); + return second; + } + + @ParameterizedTest + @ValueSource(strings = {"active", "idle", "recent", "pinned", "failed"}) + void testIdleRelease(String condition) throws Exception { + KvTablet second = secondTablet(); + assertGauges(2, 0, 0); + try (KvTablet.Guard ignored = sentinel.acquireGuard(); + KvTablet.Guard other = second.acquireGuard()) { + writeRecord(); + flushAndAwait(sentinel); + } + assertGauges(0, 2, 0); + KvTablet.Guard pin = condition.equals("pinned") ? sentinel.acquireGuard() : null; + if (!condition.equals("active")) { + clock.advanceTime(61, TimeUnit.SECONDS); + } + if (condition.equals("recent")) { + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + assertThat(lifecycle.getLastAccessTimestamp()).isEqualTo(clock.milliseconds()); + } + } + if (condition.equals("failed")) { + releaseCallback = + kv -> { + throw new IllegalStateException("release failed"); + }; + } + try { + kvManager.releaseIdleTablets(60_000, clock.milliseconds()); + assertThat(sentinel.isLazyOpen()).isEqualTo(!condition.equals("idle")); + assertThat(second.isLazyOpen()).isEqualTo(condition.equals("active")); + int open = (sentinel.isLazyOpen() ? 1 : 0) + (second.isLazyOpen() ? 1 : 0); + assertGauges(2 - open, open, 0); + if (condition.equals("failed")) { + assertValue(); + releaseCallback = kv -> {}; + assertThat(sentinel.releaseKv()).isTrue(); + assertGauges(2, 0, 0); + assertValue(); + assertGauges(1, 1, 0); + } + } finally { + if (pin != null) { + pin.close(); + } + releaseCallback = kv -> {}; + } + } + + private void open() { + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + assertThat(sentinel.getRocksDBKv()).isNotNull(); + } + } + + private void awaitState(KvTabletLazyLifecycle.LazyState state) throws Exception { + CommonTestUtils.waitUntil( + () -> lifecycle.getLazyState() == state, + Duration.ofSeconds(5), + "Waiting for " + state); + } + + private void assertGauges(int lazy, int open, int failed) { + TabletServerMetricGroup metrics = TestingMetricGroups.TABLET_SERVER_METRICS; + assertThat(metrics.kvTabletLazyCount().get() - baseLazy).isEqualTo(lazy); + assertThat(metrics.kvTabletOpenCount().get() - baseOpen).isEqualTo(open); + assertThat(metrics.kvTabletFailedCount().get() - baseFailed).isEqualTo(failed); + } + + private void assertValue() throws Exception { + try (KvTablet.Guard guard = sentinel.acquireGuard()) { + assertThat(guard.getTablet().multiGet(Collections.singletonList("k1".getBytes()))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly( + ValueEncoder.encodeValue( + schemaId, + compactedRow(baseRowType, new Object[] {1, "a"}))); + } + } + + private void writeRecord() throws Exception { + sentinel.requireOpenedTablet() + .putAsLeader( + kvRecordBatchFactory.ofRecords( + Collections.singletonList( + kvRecordFactory.ofRecord( + "k1".getBytes(), new Object[] {1, "a"}))), + null); + } + + private void flushAndAwait(KvTablet tablet) throws Exception { + long offset = logTablet.localLogEndOffset(); + tablet.requestFlush(offset, NOPErrorHandler.INSTANCE); + CommonTestUtils.waitUntil( + () -> tablet.getFlushedLogOffset() >= offset, + Duration.ofSeconds(10), + "KV operation did not complete"); + } + + @ParameterizedTest + @CsvSource({"release,false", "release,true", "failure,true", "commit,true"}) + void testCloseWaitsForCallbacks(String phase, boolean delete) throws Exception { + boolean commit = phase.equals("commit"); + try (BlockingCallback blocked = new BlockingCallback()) { + if (commit) { + commitCallback = kv -> blocked.run(); + } else { + open(); + AtomicInteger cleanupCalls = new AtomicInteger(); + releaseCallback = + kv -> { + blocked.run(); + if (cleanupCalls.incrementAndGet() == 1 + && phase.equals("failure")) { + throw new IllegalStateException("release failed"); + } + }; + } + FutureTask work = + start( + () -> { + if (commit) { + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(KvStorageException.class); + return null; + } + return sentinel.releaseKv(); + }); + blocked.awaitEntered(); + FutureTask close = + start( + () -> { + if (delete) { + sentinel.drop(); + } else { + sentinel.close(); + } + return null; + }); + awaitState(KvTabletLazyLifecycle.LazyState.CLOSING); + clock.advanceTime(1, TimeUnit.SECONDS); + assertThatThrownBy(() -> close.get(200, TimeUnit.MILLISECONDS)) + .isInstanceOf(TimeoutException.class); + blocked.close(); + work.get(5, TimeUnit.SECONDS); + close.get(5, TimeUnit.SECONDS); + } + releaseCallback = kv -> {}; + assertThat(lifecycle.getLazyState()).isEqualTo(KvTabletLazyLifecycle.LazyState.CLOSED); + assertThat(sentinel.getRocksDBKv()).isNull(); + assertThat(sentinel.getKvTabletDir().exists()).isEqualTo(!delete); + assertThatThrownBy(sentinel::acquireGuard).isInstanceOf(KvStorageException.class); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testCloseWaitsForFencedOpenCleanup(boolean shutdown) throws Exception { + try (BlockingCallback blocked = new BlockingCallback()) { + AtomicReference candidate = new AtomicReference<>(); + AtomicReference lease = new AtomicReference<>(); + openCallback = + hasLocal -> { + KvTablet opened = openTablet(sentinel); + candidate.set(opened); + lease.set(opened.getRocksDBKv().getResourceGuard().acquireResource()); + blocked.run(); + return opened; + }; + FutureTask open = + start( + () -> { + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(RuntimeException.class); + return null; + }); + blocked.awaitEntered(); + FutureTask close = + start( + () -> { + if (shutdown) { + kvManager.shutdown(); + } else { + kvManager.dropKv(tableBucket); + } + return null; + }); + try { + awaitState(KvTabletLazyLifecycle.LazyState.CLOSING); + blocked.close(); + CommonTestUtils.waitUntil( + () -> candidate.get().getRocksDBKv().getResourceGuard().isClosed(), + Duration.ofSeconds(5), + "Candidate cleanup did not start"); + assertThatThrownBy(() -> close.get(200, TimeUnit.MILLISECONDS)) + .isInstanceOf(TimeoutException.class); + assertThat(kvManager.getKv(tableBucket)).contains(sentinel); + assertThatThrownBy(() -> kvManager.registerKv(tableBucket, sentinel)) + .isInstanceOf(IllegalStateException.class); + } finally { + blocked.close(); + if (lease.get() != null) { + lease.get().close(); + } + open.get(5, TimeUnit.SECONDS); + close.get(5, TimeUnit.SECONDS); + } + assertThat(lifecycle.getLazyState()) + .isEqualTo(KvTabletLazyLifecycle.LazyState.CLOSED); + assertThat(sentinel.getRocksDBKv()).isNull(); + assertThat(sentinel.getKvTabletDir()).doesNotExist(); + if (shutdown) { + kvManager = null; + } else { + assertThat(kvManager.getKv(tableBucket)).isEmpty(); + KvTablet replacement = createSentinel(); + openCallback = local -> openTablet(replacement); + configureLifecycle(replacement); + try (KvTablet.Guard ignored = replacement.acquireGuard()) { + assertThat(replacement.getKvTabletDir()).isDirectory(); + } + } + } + } + + @Test + void testDropAfterCloseDeletesLocalDataOnlyOnce() throws Exception { + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + assertThat(sentinel.getKvTabletDir()).isDirectory(); + } + sentinel.close(); + assertThat(sentinel.getKvTabletDir()).isDirectory(); + sentinel.drop(); + assertThat(sentinel.getKvTabletDir()).doesNotExist(); + kvManager.dropKv(tableBucket); + KvTablet replacement = createSentinel(); + configureLifecycle(replacement); + try (KvTablet.Guard ignored = replacement.acquireGuard()) { + sentinel.drop(); + assertThat(replacement.getKvTabletDir()).isDirectory(); + } + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testReleaseChecksWritesCompletedDuringDrain(boolean flush) throws Exception { + FutureTask release; + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + release = start(sentinel::releaseKv); + awaitState(KvTabletLazyLifecycle.LazyState.RELEASING); + writeRecord(); + if (flush) { + flushAndAwait(sentinel); + assertThat(sentinel.getFlushedLogOffset()) + .isEqualTo(logTablet.localLogEndOffset()); + } + } + assertThat(release.get(5, TimeUnit.SECONDS)).isEqualTo(flush); + assertThat(sentinel.isLazyOpen()).isEqualTo(!flush); + } + + @ParameterizedTest + @ValueSource(strings = {"LAZY", "OPEN", "FAILED"}) + void testDropIsTerminal(String state) throws Exception { + if (state.equals("OPEN")) { + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + assertThat(sentinel.getKvTabletDir()).isDirectory(); + } + } else if (state.equals("FAILED")) { + openCallback = + local -> { + throw new IllegalStateException("open failed"); + }; + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(IllegalStateException.class); + } + sentinel.drop(); + assertGauges(0, 0, 0); + assertThat(lifecycle.getLazyState()).isEqualTo(KvTabletLazyLifecycle.LazyState.CLOSED); + assertThat(sentinel.getRocksDBKv()).isNull(); + assertThat(sentinel.getKvTabletDir()).doesNotExist(); + assertThatThrownBy(sentinel::acquireGuard).isInstanceOf(KvStorageException.class); + } + + @Test + void testSnapshotExecutorDoesNotRunAgainstReopenedTablet() throws Exception { + Executor executor; + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + executor = sentinel.getGuardedExecutor(); + } + assertThat(sentinel.releaseKv()).isTrue(); + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + AtomicInteger executions = new AtomicInteger(); + executor.execute(executions::incrementAndGet); + assertThat(executions).hasValue(0); + sentinel.getGuardedExecutor().execute(executions::incrementAndGet); + assertThat(executions).hasValue(1); + } + } + + @Test + void testMaintenancePinsTabletUntilCompletion() throws Exception { + open(); + try (BlockingCallback blocked = new BlockingCallback()) { + FutureTask task = + start( + () -> { + sentinel.getGuardedExecutor().execute(blocked::run); + return null; + }); + blocked.awaitEntered(); + FutureTask release = start(sentinel::releaseKv); + awaitState(KvTabletLazyLifecycle.LazyState.RELEASING); + assertThat(release.isDone()).isFalse(); + blocked.close(); + task.get(5, TimeUnit.SECONDS); + assertThat(release.get(5, TimeUnit.SECONDS)).isTrue(); + assertThat(lifecycle.getActivePins()).isZero(); + } + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testEpochFencingRejectsStaleOpen(boolean leader) throws Exception { + AtomicInteger epoch = new AtomicInteger(1); + if (leader) { + leaderEpochSupplier = epoch::get; + } else { + bucketEpochSupplier = epoch::get; + } + try (BlockingCallback blocked = new BlockingCallback()) { + openCallback = + local -> { + blocked.run(); + return openTablet(sentinel); + }; + FutureTask task = + start( + () -> { + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(RuntimeException.class); + return null; + }); + blocked.awaitEntered(); + epoch.incrementAndGet(); + blocked.close(); + task.get(5, TimeUnit.SECONDS); + assertThat(lifecycle.getLazyState()) + .isEqualTo(KvTabletLazyLifecycle.LazyState.FAILED); + assertThat(sentinel.getKvTabletDir()).doesNotExist(); + } + } + + @Test + void testConcurrentRequestsShareOpenDespiteCommitCallbackFailure() throws Exception { + commitCallback = + kv -> { + throw new IllegalStateException("commit failed"); + }; + AtomicInteger opens = new AtomicInteger(); + try (BlockingCallback blocked = new BlockingCallback()) { + openCallback = + local -> { + opens.incrementAndGet(); + blocked.run(); + return openTablet(sentinel); + }; + List> requests = new java.util.ArrayList<>(); + for (int i = 0; i < 5; i++) { + requests.add( + start( + () -> { + open(); + return null; + })); + } + blocked.awaitEntered(); + blocked.close(); + for (FutureTask request : requests) { + request.get(5, TimeUnit.SECONDS); + } + assertThat(opens).hasValue(1); + } + try (KvTablet.Guard guard = sentinel.acquireGuard()) { + assertThat( + guard.getTablet() + .multiGet(Collections.singletonList("missing".getBytes()))) + .containsExactly((ByteArraySlice) null); + guard.close(); + guard.close(); + assertThat(lifecycle.getActivePins()).isZero(); + } + } + + @Test + void testReleaseDrainTimeoutRollsBackToOpen() throws Exception { + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + FutureTask release = start(sentinel::releaseKv); + awaitState(KvTabletLazyLifecycle.LazyState.RELEASING); + clock.advanceTime(1, TimeUnit.SECONDS); + assertThat(release.get(5, TimeUnit.SECONDS)).isFalse(); + assertThat(sentinel.isLazyOpen()).isTrue(); + } + assertThat(lifecycle.getActivePins()).isZero(); + } + + @Test + void testExponentialBackoffAndRecovery() throws Exception { + openCallback = + local -> { + throw new IllegalStateException("open failed"); + }; + for (int cooldown : new int[] {10, 20, 40}) { + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(IllegalStateException.class); + assertGauges(0, 0, 1); + clock.advanceTime(cooldown - 1, TimeUnit.SECONDS); + assertThatThrownBy(sentinel::acquireGuard) + .isInstanceOf(KvStorageException.class) + .hasMessageContaining("cooldown"); + clock.advanceTime(2, TimeUnit.SECONDS); + } + openCallback = local -> openTablet(sentinel); + try (KvTablet.Guard ignored = sentinel.acquireGuard()) { + assertThat(sentinel.getRocksDBKv()).isNotNull(); + } + } + } } diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaLazyOpenTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaLazyOpenTest.java new file mode 100644 index 0000000000..c7f7368700 --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaLazyOpenTest.java @@ -0,0 +1,303 @@ +/* + * 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.replica; + +import org.apache.fluss.config.ConfigOptions; +import org.apache.fluss.config.Configuration; +import org.apache.fluss.metadata.LogFormat; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.record.ChangeType; +import org.apache.fluss.record.DefaultValueRecordBatch; +import org.apache.fluss.record.KvRecordBatch; +import org.apache.fluss.row.encode.CompactedKeyEncoder; +import org.apache.fluss.row.encode.ValueEncoder; +import org.apache.fluss.rpc.protocol.MergeMode; +import org.apache.fluss.server.entity.NotifyLeaderAndIsrData; +import org.apache.fluss.server.kv.KvTablet; +import org.apache.fluss.server.kv.scan.ScannerContext; +import org.apache.fluss.server.kv.snapshot.CompletedSnapshot; +import org.apache.fluss.server.log.LogAppendInfo; +import org.apache.fluss.server.zk.NOPErrorHandler; +import org.apache.fluss.server.zk.data.LeaderAndIsr; +import org.apache.fluss.testutils.common.CommonTestUtils; +import org.apache.fluss.utils.ByteArraySlice; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.io.File; +import java.time.Duration; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.apache.fluss.compression.ArrowCompressionInfo.DEFAULT_COMPRESSION; +import static org.apache.fluss.record.LogRecordBatch.CURRENT_LOG_MAGIC_VALUE; +import static org.apache.fluss.record.LogRecordBatchFormat.NO_BATCH_SEQUENCE; +import static org.apache.fluss.record.LogRecordBatchFormat.NO_WRITER_ID; +import static org.apache.fluss.record.TestData.DATA1_PHYSICAL_TABLE_PATH_PK; +import static org.apache.fluss.record.TestData.DATA1_ROW_TYPE; +import static org.apache.fluss.record.TestData.DATA1_TABLE_ID_PK; +import static org.apache.fluss.record.TestData.DEFAULT_SCHEMA_ID; +import static org.apache.fluss.testutils.DataTestUtils.compactedRow; +import static org.apache.fluss.testutils.DataTestUtils.createBasicMemoryLogRecords; +import static org.apache.fluss.testutils.DataTestUtils.genKvRecordBatch; +import static org.apache.fluss.utils.FileUtils.deleteDirectory; +import static org.apache.fluss.utils.Preconditions.checkNotNull; +import static org.assertj.core.api.Assertions.assertThat; + +/** Integration tests for lazy-open behavior in {@link Replica}. */ +final class ReplicaLazyOpenTest extends ReplicaTestBase { + + @Override + protected Configuration getServerConf() { + Configuration conf = super.getServerConf(); + conf.set(ConfigOptions.KV_LAZY_OPEN_ENABLED, true); + conf.set(ConfigOptions.KV_LAZY_OPEN_IDLE_TIMEOUT, Duration.ofHours(1)); + return conf; + } + + @Test + void testPeriodicSnapshotDoesNotKeepIdleTabletOpen() throws Exception { + Replica replica = makeLazyLeaderReplica(); + replica.lookups(Collections.singletonList(key(1))); + KvTablet tablet = checkNotNull(replica.getKvTablet()); + manualClock.advanceTime(Duration.ofMinutes(10)); + checkNotNull(replica.getKvSnapshotManager()).triggerSnapshot(); + assertThat(tablet.canRelease(Duration.ofMinutes(10).toMillis(), manualClock.milliseconds())) + .isTrue(); + replica.lookups(Collections.singletonList(key(1))); + assertThat(tablet.canRelease(Duration.ofMinutes(10).toMillis(), manualClock.milliseconds())) + .isFalse(); + } + + @Test + void testFirstLookupFlushesRecoveredRecordsBeyondHighWatermark() throws Exception { + Replica replica = + makeKvReplica(DATA1_PHYSICAL_TABLE_PATH_PK, new TableBucket(DATA1_TABLE_ID_PK, 1)); + replica.appendRecordsToFollower( + createBasicMemoryLogRecords( + DATA1_ROW_TYPE, + DEFAULT_SCHEMA_ID, + 0L, + -1L, + CURRENT_LOG_MAGIC_VALUE, + NO_WRITER_ID, + NO_BATCH_SEQUENCE, + Collections.singletonList(ChangeType.INSERT), + Collections.singletonList(new Object[] {1, "a"}), + LogFormat.ARROW, + DEFAULT_COMPRESSION)); + replica.makeLeader(notifyLeaderAndIsr(0)); + assertThat(replica.getLogHighWatermark()).isZero(); + assertThat(replica.getLocalLogEndOffset()).isEqualTo(1L); + + replica.lookups(Collections.singletonList(key(1))); + + CommonTestUtils.waitUntil( + () -> replica.getLogHighWatermark() == 1L, + Duration.ofSeconds(10), + "Recovered WAL did not become visible after lazy open"); + assertThat(replica.lookups(Collections.singletonList(key(1)))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly(valueBytes(1, "a")); + } + + @Test + void testPutRecordsToLeaderTriggersLazyOpen() throws Exception { + Replica replica = makeLazyLeaderReplica(); + KvTablet kvTablet = checkNotNull(replica.getKvTablet()); + assertThat(kvTablet.isLazyOpen()).isFalse(); + assertThat(kvTablet.getRocksDBKv()).isNull(); + + LogAppendInfo appendInfo = + putRecordsToLeaderAndFlush(replica, genKvRecordBatch(new Object[] {1, "a"})); + + assertThat(appendInfo.lastOffset()).isEqualTo(0L); + CommonTestUtils.waitUntil( + () -> replica.getLogTablet().getHighWatermark() == appendInfo.lastOffset() + 1, + Duration.ofSeconds(10), + "KV operation did not complete"); + assertThat(kvTablet.isLazyOpen()).isTrue(); + assertThat(kvTablet.getRocksDBKv()).isNotNull(); + assertThat(replica.lookups(Collections.singletonList(key(1)))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly(valueBytes(1, "a")); + } + + @ParameterizedTest + @ValueSource(strings = {"lookup", "prefix", "scan"}) + void testReadReopensReleasedTablet(String operation) throws Exception { + Replica replica = makeLazyLeaderReplica(); + putRecordsToLeaderAndFlush( + replica, genKvRecordBatch(new Object[] {1, "a"}, new Object[] {2, "b"})); + long flushed = replica.getKvTablet().getFlushedLogOffset(); + releaseToLazy(replica); + assertThat(replica.getKvTablet().getFlushedLogOffset()).isEqualTo(flushed); + assertThat(replica.getRowCount()).isEqualTo(2L); + assertThat(replica.getKvTablet().isLazyOpen()).isFalse(); + if (operation.equals("scan")) { + DefaultValueRecordBatch.Builder builder = DefaultValueRecordBatch.builder(); + builder.append(DEFAULT_SCHEMA_ID, compactedRow(DATA1_ROW_TYPE, new Object[] {1, "a"})); + builder.append(DEFAULT_SCHEMA_ID, compactedRow(DATA1_ROW_TYPE, new Object[] {2, "b"})); + assertThat(replica.limitKvScan(2)).isEqualTo(builder.build()); + } else { + assertThat( + operation.equals("prefix") + ? replica.prefixLookup(new byte[0]) + : replica.lookups(Arrays.asList(key(1), key(2)))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly(valueBytes(1, "a"), valueBytes(2, "b")); + } + assertThat(replica.getKvTablet().isLazyOpen()).isTrue(); + } + + @Test + void testScannerPreventsIdleReleaseUntilClosed() throws Exception { + Replica replica = makeLazyLeaderReplica(); + putRecordsToLeaderAndFlush(replica, genKvRecordBatch(new Object[] {1, "a"})); + releaseToLazy(replica); + + ScannerContext scanner = + checkNotNull(scannerManager.createScanner(replica, null).getContext()); + KvTablet tablet = checkNotNull(replica.getKvTablet()); + try { + assertThat(tablet.isLazyOpen()).isTrue(); + assertThat(tablet.releaseKv()).isFalse(); + assertThat(scanner.isValid()).isTrue(); + assertThat(scanner.currentValue()).isEqualTo(valueBytes(1, "a")); + } finally { + scannerManager.closeScannersForBucket(replica.getTableBucket()); + } + releaseToLazy(replica); + assertThat(replica.lookups(Collections.singletonList(key(1)))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly(valueBytes(1, "a")); + } + + @Test + void testReleasedTabletUsesCommittedSnapshotSize() throws Exception { + Replica replica = makeLazyLeaderReplica(); + putRecordsToLeaderAndFlush( + replica, genKvRecordBatch(new Object[] {1, "a"}, new Object[] {2, "b"})); + checkNotNull(replica.getKvSnapshotManager()).triggerSnapshot(); + CompletedSnapshot snapshot = + snapshotReporter.waitUntilSnapshotComplete(replica.getTableBucket(), 0); + assertThat(snapshot.getSnapshotSize()).isPositive(); + releaseToLazy(replica); + assertThat(replica.getKvSnapshotManager()).isNull(); + assertThat(replica.getLatestKvSnapshotSize()).isEqualTo(snapshot.getSnapshotSize()); + assertThat(replica.logicalStorageKvSize()).isEqualTo(snapshot.getSnapshotSize()); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void testBecomingFollowerDropsTablet(boolean opened) throws Exception { + Replica replica = makeLazyLeaderReplica(); + if (opened) { + putRecordsToLeaderAndFlush(replica, genKvRecordBatch(new Object[] {1, "a"})); + } + replica.makeFollower(notifyLeaderAndIsr(1)); + assertThat(replica.getKvTablet()).isNull(); + assertThat(kvManager.getKv(replica.getTableBucket())).isEmpty(); + } + + @Test + void testReopenFallsBackToFullInitWhenLocalDataCorrupted() throws Exception { + Replica replica = makeLazyLeaderReplica(); + putRecordsToLeaderAndFlush( + replica, genKvRecordBatch(new Object[] {1, "a"}, new Object[] {2, "b"})); + + File tabletDir = replica.getKvTablet().getKvTabletDir(); + + releaseToLazy(replica); + + assertThat(tabletDir).exists(); + deleteDirectory(tabletDir); + assertThat(tabletDir).doesNotExist(); + + List values = replica.lookups(Collections.singletonList(key(1))); + + assertThat(checkNotNull(replica.getKvTablet()).isLazyOpen()).isTrue(); + assertThat(values.get(0).toByteArray()).isEqualTo(valueBytes(1, "a")); + + putRecordsToLeaderAndFlush(replica, genKvRecordBatch(new Object[] {3, "c"})); + assertThat(replica.lookups(Collections.singletonList(key(3)))) + .extracting(ByteArraySlice::toByteArray) + .containsExactly(valueBytes(3, "c")); + } + + private Replica makeLazyLeaderReplica() throws Exception { + Replica replica = + makeKvReplica( + DATA1_PHYSICAL_TABLE_PATH_PK, + new TableBucket(DATA1_TABLE_ID_PK, 1), + new TestSnapshotContext( + conf.getString(ConfigOptions.REMOTE_DATA_DIR), snapshotReporter)); + replica.makeLeader(notifyLeaderAndIsr(0)); + return replica; + } + + private NotifyLeaderAndIsrData notifyLeaderAndIsr(int epoch) { + return new NotifyLeaderAndIsrData( + DATA1_PHYSICAL_TABLE_PATH_PK, + new TableBucket(DATA1_TABLE_ID_PK, 1), + Collections.singletonList(TABLET_SERVER_ID), + new LeaderAndIsr( + TABLET_SERVER_ID, + epoch, + Collections.singletonList(TABLET_SERVER_ID), + Collections.emptyList(), + 0, + epoch)); + } + + private LogAppendInfo putRecordsToLeaderAndFlush(Replica replica, KvRecordBatch kvRecords) + throws Exception { + LogAppendInfo appendInfo = + replica.putRecordsToLeader(kvRecords, null, MergeMode.DEFAULT, 0); + KvTablet tablet = checkNotNull(replica.getKvTablet()); + tablet.requestFlush(replica.getLocalLogEndOffset(), NOPErrorHandler.INSTANCE); + CommonTestUtils.waitUntil( + () -> + tablet.getFlushedLogOffset() >= replica.getLocalLogEndOffset() + && replica.getLogHighWatermark() >= replica.getLocalLogEndOffset(), + Duration.ofSeconds(10), + "KV operation did not complete"); + return appendInfo; + } + + private void releaseToLazy(Replica replica) throws Exception { + KvTablet kvTablet = checkNotNull(replica.getKvTablet()); + CommonTestUtils.waitUntil( + kvTablet::releaseKv, Duration.ofSeconds(10), "KV operation did not complete"); + assertThat(kvTablet.isLazyOpen()).isFalse(); + assertThat(kvTablet.getRocksDBKv()).isNull(); + } + + private byte[] key(int id) { + return new CompactedKeyEncoder(DATA1_ROW_TYPE, new int[] {0}) + .encodeKey(compactedRow(DATA1_ROW_TYPE, new Object[] {id, ""})); + } + + private byte[] valueBytes(int id, String value) { + return ValueEncoder.encodeValue( + DEFAULT_SCHEMA_ID, compactedRow(DATA1_ROW_TYPE, new Object[] {id, value})); + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTest.java index 50a9dc79ec..fb25d48f48 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/ReplicaTest.java @@ -47,11 +47,13 @@ import org.apache.fluss.server.kv.snapshot.CompletedSnapshot; import org.apache.fluss.server.kv.snapshot.KvSnapshotDataDownloader; import org.apache.fluss.server.kv.snapshot.KvSnapshotDownloadSpec; +import org.apache.fluss.server.kv.snapshot.KvSnapshotHandle; import org.apache.fluss.server.kv.snapshot.TestingCompletedKvSnapshotCommitter; import org.apache.fluss.server.log.FetchParams; import org.apache.fluss.server.log.ListOffsetsParam; import org.apache.fluss.server.log.LogAppendInfo; import org.apache.fluss.server.log.LogReadInfo; +import org.apache.fluss.server.log.LogTablet; import org.apache.fluss.server.testutils.KvTestUtils; import org.apache.fluss.server.zk.data.LeaderAndIsr; import org.apache.fluss.testutils.DataTestUtils; @@ -1276,6 +1278,66 @@ private void verifyGetKeyValues( .containsExactlyElementsOf(expectValues); } + @Test + void testGetRowCountPkTableWithNullKvTablet() throws Exception { + + TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID_PK, 1); + Replica kvReplica = makeKvReplica(DATA1_PHYSICAL_TABLE_PATH_PK, tableBucket); + + assertThat(kvReplica.isKvTable()).isTrue(); + assertThat(kvReplica.getKvTablet()).isNull(); + + LogTablet logTablet = kvReplica.getLogTablet(); + MemoryLogRecords records = genMemoryLogRecordsByObject(DATA1); + logTablet.appendAsLeader(records); + logTablet.updateHighWatermark(logTablet.localLogEndOffset()); + assertThat(logTablet.getRowCount()).isGreaterThan(0); + + assertThat(kvReplica.getRowCount()).isEqualTo(0L); + } + + @Test + void testCreateKvRollbackOnAllRetriesFailed() throws Exception { + + TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID_PK, 1); + TestSnapshotContext failingSnapshotContext = + new TestSnapshotContext(conf.getString(ConfigOptions.REMOTE_DATA_DIR)) { + @Override + public FunctionWithException + getLatestCompletedSnapshotProvider() { + + return tb -> + new CompletedSnapshot( + tb, + 1L, + new FsPath("file:///non-existent-path/snapshot-1"), + KvSnapshotHandle.create( + Collections.emptyList(), + Collections.emptyList(), + 0), + 0L, + null, + null); + } + + @Override + public KvSnapshotDataDownloader getSnapshotDataDownloader() { + + throw new IllegalStateException("Snapshot download unavailable"); + } + }; + + Replica kvReplica = + makeKvReplica(DATA1_PHYSICAL_TABLE_PATH_PK, tableBucket, failingSnapshotContext); + + assertThatThrownBy(() -> makeKvReplicaAsLeader(kvReplica)) + .isInstanceOf(org.apache.fluss.exception.KvStorageException.class); + + assertThat(kvReplica.getKvTablet()).isNull(); + + assertThat(kvManager.getKv(tableBucket)).isEmpty(); + } + /** A scheduledExecutorService that will execute the scheduled task immediately. */ private static class ImmediateTriggeredScheduledExecutorService extends ManuallyTriggeredScheduledExecutorService { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java index 73d41737d8..07e87adec6 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java @@ -62,6 +62,7 @@ import org.apache.fluss.server.entity.NotifyLeaderAndIsrData; import org.apache.fluss.server.entity.NotifyLeaderAndIsrResultForBucket; import org.apache.fluss.server.entity.PutKvDataForBucket; +import org.apache.fluss.server.kv.KvManager; import org.apache.fluss.server.kv.KvStateLookupResult; import org.apache.fluss.server.kv.KvTablet; import org.apache.fluss.server.kv.historical.HistoricalKvKeyEncoder; @@ -75,12 +76,14 @@ import org.apache.fluss.server.metadata.PartitionMetadata; import org.apache.fluss.server.metadata.ServerInfo; import org.apache.fluss.server.metadata.TableMetadata; +import org.apache.fluss.server.metrics.group.TestingMetricGroups; import org.apache.fluss.server.replica.Replica; import org.apache.fluss.server.replica.ReplicaTestBase; import org.apache.fluss.server.zk.data.LeaderAndIsr; import org.apache.fluss.server.zk.data.TableRegistration; import org.apache.fluss.server.zk.data.lake.LakeTableHelper; import org.apache.fluss.server.zk.data.lake.LakeTableSnapshot; +import org.apache.fluss.testutils.common.CommonTestUtils; import org.apache.fluss.testutils.common.ManuallyTriggeredScheduledExecutorService; import org.apache.fluss.types.DataField; import org.apache.fluss.types.DataTypes; @@ -142,6 +145,73 @@ class HistoricalPartitionManagerTest extends ReplicaTestBase { new DataField("id", DataTypes.INT()), new DataField("region", DataTypes.STRING())); + @Test + void testHistoricalLazyOpenReleaseAndRecover() throws Exception { + replicaManager.shutdown(); + kvManager.shutdown(); + conf.set(ConfigOptions.KV_LAZY_OPEN_ENABLED, true); + kvManager = + KvManager.create( + conf, + zkClient, + logManager, + TestingMetricGroups.TABLET_SERVER_METRICS, + localDiskManager, + manualClock); + kvManager.startup(); + replicaManager = buildReplicaManager(testCoordinatorGateway); + replicaManager.startup(); + + TableInfo tableInfo = registerHistoricalTableAndBecomeLeader(); + Replica replica = replicaManager.getReplicaOrException(TABLE_BUCKET); + KvTablet tablet = replica.getKvTablet(); + assertThat(tablet).isNotNull(); + assertThat(tablet.isLazyOpen()).isFalse(); + TestingHistoricalLakeLookupManager lakeLookupManager = + new TestingHistoricalLakeLookupManager(lookupConfiguration()); + try (HistoricalPartitionManager manager = + createHistoricalPartitionManager(lakeLookupManager)) { + writeBatch( + manager, + replica, + ORIGINAL_PARTITION, + batch(tableInfo.getRowType(), upsert(1, "us", ORIGINAL_PARTITION, "v1"))); + flushAndWait(tablet, replica.getLocalLogEndOffset()); + assertThat(tablet.isLazyOpen()).isTrue(); + assertThat(replica.getKvSnapshotManager()).isNull(); + CommonTestUtils.waitUntil( + tablet::releaseKv, Duration.ofSeconds(10), "KV release did not finish"); + assertThat(tablet.getRocksDBKv()).isNull(); + replica.tryUpdateHistoricalCleanupOffset(replica.getLocalLogEndOffset(), () -> {}); + + assertThat( + replica.lookupHistoricalLocal( + ORIGINAL_PARTITION, Collections.singletonList(key(1, "us")))) + .hasSize(1); + assertHistoricalValue( + tablet, + ORIGINAL_PARTITION, + key(1, "us"), + tableInfo, + row(1, "us", ORIGINAL_PARTITION, "v1")); + assertThat(replica.getKvSnapshotManager()).isNull(); + assertThat(tablet.getHistoricalCleanupOffset()) + .isEqualTo(replica.getLocalLogEndOffset()); + writeBatch( + manager, + replica, + ORIGINAL_PARTITION, + batch(tableInfo.getRowType(), upsert(1, "us", ORIGINAL_PARTITION, "v2"))); + flushAndWait(tablet, replica.getLocalLogEndOffset()); + assertHistoricalValue( + tablet, + ORIGINAL_PARTITION, + key(1, "us"), + tableInfo, + row(1, "us", ORIGINAL_PARTITION, "v2")); + } + } + @Test void testResolvesMultipleLakeMissesWithoutPrewriteRollback() throws Exception { TableInfo tableInfo = registerHistoricalTableAndBecomeLeader(); @@ -1161,7 +1231,7 @@ TABLE_BUCKET, new FetchReqInfo(TABLE_ID, fetchOffset, Integer.MAX_VALUE)), null, future::complete); FetchLogResultForBucket result = future.get(10, TimeUnit.SECONDS).get(TABLE_BUCKET); - assertThat(result.failed()).isFalse(); + assertThat(result.failed()).as("Historical lookup failed: %s", result.getError()).isFalse(); return result.records(); } @@ -1341,7 +1411,7 @@ private static void assertHistoricalLookup( originalPartition), (lookupTimeNanos, lookupFileDownloaded) -> {}) .get(10, TimeUnit.SECONDS); - assertThat(result.failed()).isFalse(); + assertThat(result.failed()).as("Historical lookup failed: %s", result.getError()).isFalse(); assertThat(result.lookupValues()).hasSize(1); if (expectedRow == null) { assertThat(result.lookupValues().get(0)).isNull(); @@ -1363,7 +1433,10 @@ private static void assertHistoricalValue( TableInfo tableInfo, InternalRow expectedRow) throws Exception { - KvStateLookupResult result = kvTablet.lookupHistoricalLocal(originalPartition, primaryKey); + KvStateLookupResult result; + try (KvTablet.Guard guard = kvTablet.acquireGuard()) { + result = guard.getTablet().lookupHistoricalLocal(originalPartition, primaryKey); + } assertThat(result.isPresent()).isTrue(); BinaryValue value = new ValueDecoder( diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index fb88ab0a0a..8aa9c64073 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -620,6 +620,22 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM The cumulative number of lookup files evicted to enforce the shared TabletServer disk budget. Expiration, replacement, and explicit invalidation are excluded. Gauge + + kvLazyOpen + kvTabletOpenCount + The number of KvTablets currently in OPEN state (RocksDB loaded). + Gauge + + + kvTabletLazyCount + The number of KvTablets currently in LAZY state (RocksDB not loaded). + Gauge + + + kvTabletFailedCount + The number of KvTablets currently in FAILED state (open attempt failed, in backoff cooldown). + Gauge + logicalStorage logSize