Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions LICENSE
Original file line number Diff line number Diff line change
Expand Up @@ -368,6 +368,7 @@ Apache Kafka
./fluss-server/src/main/java/org/apache/fluss/server/utils/timer/TimingWheel.java

Apache Paimon
./fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonLocalTableQuery.java
./fluss-common/src/main/java/org/apache/fluss/predicate/And.java
./fluss-common/src/main/java/org/apache/fluss/predicate/CompareUtils.java
./fluss-common/src/main/java/org/apache/fluss/predicate/CompoundPredicate.java
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -429,7 +429,7 @@ public class ConfigOptions {
.defaultValue(0.10)
.withDescription(
"The maximum fraction of the total capacity of the volume containing the first available data directory allocated to historical partition lookup caches on a TabletServer. "
+ "Up to ten table lookupers are cached, and each receives one tenth of this capacity. Historical lookup cache files are stored under that data directory; additional data volumes are not used. "
+ "All table lookupers share this capacity, with eviction at file granularity. Historical lookup cache files are stored under that data directory; additional data volumes are not used. "
+ "The valid range is (0.0, 1.0].");

public static final ConfigOption<Duration>
Expand All @@ -438,7 +438,9 @@ public class ConfigOptions {
.durationType()
.defaultValue(Duration.ofHours(3))
.withDescription(
"The duration after which an idle historical partition table lookuper is removed from the cache.");
"The duration after which an idle historical partition table lookuper or an idle lookup file is removed from its cache. "
+ "Lookuper and file access times are tracked independently. This setting replaces the Paimon table-level lookup.cache-file-retention option for historical lookups. "
+ "Dynamic changes apply to both caches without replacing active lookupers.");

public static final ConfigOption<Double> SERVER_DATA_DISK_WRITE_LIMIT_RATIO =
key("server.data-disk.write-limit-ratio")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,11 @@
package org.apache.fluss.lake.lakestorage;

import org.apache.fluss.annotation.PublicEvolving;
import org.apache.fluss.config.TableConfig;
import org.apache.fluss.lake.lakestorage.LakeTableLookuperManager.LookupCacheOptions;
import org.apache.fluss.lake.source.LakeSource;
import org.apache.fluss.lake.writer.LakeTieringFactory;
import org.apache.fluss.metadata.LakeLookupMode;
import org.apache.fluss.metadata.TablePath;

import static org.apache.fluss.utils.Preconditions.checkArgument;
import static org.apache.fluss.utils.Preconditions.checkNotNull;

/**
* The LakeStorage interface defines how to implement lakehouse storage system such as Paimon and
* Iceberg. It provides a method to create a lake tiering factory.
Expand Down Expand Up @@ -57,71 +53,15 @@ public interface LakeStorage {
LakeSource<?> createLakeSource(TablePath tablePath);

/**
* Creates a table-level point lookuper for the specified lake table.
* Creates a TabletServer-scoped manager for lake table lookupers and their shared resources.
*
* @param tablePath the logical path identifying the table in the lakehouse storage
* @param context runtime context for creating the lookuper
* @return a table-level point lookuper
* @param ioTmpDir local directory shared by lookupers for temporary files
* @param options initial runtime resource settings
* @return the lookuper manager
*/
default LakeTableLookuper createLakeTableLookuper(
TablePath tablePath, LookuperContext context) {
default LakeTableLookuperManager createLakeTableLookuperManager(
String ioTmpDir, LookupCacheOptions options) {
throw new UnsupportedOperationException(
"Point lookup is not supported for this lake storage.");
}

/** Runtime context for creating a lake table lookuper. */
final class LookuperContext {
private final String ioTmpDir;
private final TableConfig tableConfig;
private final long lookupCacheMaxDiskBytes;
private final Runnable diskWriteGuard;
private final LakeLookupMode lookupMode;

/**
* Creates a lookuper context.
*
* @param ioTmpDir local directory for temporary files used by the lookuper
* @param tableConfig configuration of the Fluss table
* @param lookupCacheMaxDiskBytes maximum local lookup cache size in bytes
* @param diskWriteGuard guard invoked before creating a local lookup cache file
*/
public LookuperContext(
String ioTmpDir,
TableConfig tableConfig,
long lookupCacheMaxDiskBytes,
Runnable diskWriteGuard) {
this.ioTmpDir = checkNotNull(ioTmpDir, "ioTmpDir must not be null.");
this.tableConfig = checkNotNull(tableConfig, "tableConfig must not be null.");
checkArgument(
lookupCacheMaxDiskBytes > 0, "lookupCacheMaxDiskBytes must be greater than 0.");
this.lookupCacheMaxDiskBytes = lookupCacheMaxDiskBytes;
this.diskWriteGuard = checkNotNull(diskWriteGuard, "diskWriteGuard must not be null.");
this.lookupMode = tableConfig.getHistoricalLookupMode();
}

/** Returns the local directory for temporary files used by the lookuper. */
public String ioTmpDir() {
return ioTmpDir;
}

/** Returns the configuration of the Fluss table. */
public TableConfig tableConfig() {
return tableConfig;
}

/** Returns the maximum local lookup cache size in bytes. */
public long lookupCacheMaxDiskBytes() {
return lookupCacheMaxDiskBytes;
}

/** Returns the guard invoked before creating a local lookup cache file. */
public Runnable diskWriteGuard() {
return diskWriteGuard;
}

/** Returns the mode used to look up historical data. */
public LakeLookupMode lookupMode() {
return lookupMode;
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
/*
* 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.lake.lakestorage;

import org.apache.fluss.annotation.PublicEvolving;
import org.apache.fluss.config.Configuration;
import org.apache.fluss.config.TableConfig;
import org.apache.fluss.metadata.TablePath;

import java.time.Duration;

import static org.apache.fluss.utils.Preconditions.checkArgument;
import static org.apache.fluss.utils.Preconditions.checkNotNull;

/**
* Creates lake table lookupers and manages their shared resources within one TabletServer.
*
* <p>The caller must close all created lookupers after their requests have finished before closing
* this manager. Closing a lookuper must not release resources shared with other lookupers.
*
* @since 1.1
*/
@PublicEvolving
public interface LakeTableLookuperManager extends AutoCloseable {

/**
* Creates a table-level point lookuper for the specified lake table.
*
* @param tablePath the logical path identifying the table in the lakehouse storage
* @param context runtime context for creating the lookuper
* @return a table-level point lookuper
*/
LakeTableLookuper createLakeTableLookuper(TablePath tablePath, Context context);

/**
* Applies a new snapshot of the runtime resource settings to existing and future lookupers
* without replacing them.
*
* <p>This method may be called concurrently with lookuper creation and lookups. Implementations
* must apply the settings in a thread-safe manner. Lake-format and table-specific configuration
* is outside the scope of this method.
*
* @param options the new runtime resource settings
*/
void reconfigure(LookupCacheOptions options);

/** Returns the cumulative number of lookup files evicted by the shared disk-space budget. */
default long fileCacheCapacityEvictions() {
return 0L;
}

/**
* Immutable snapshot of the format-independent resource settings for a lookup runtime.
*
* <p>These settings apply to resources shared by all table lookupers in one runtime.
* Lake-format and table-specific configuration is supplied separately when creating a lookuper.
*/
final class LookupCacheOptions {

private final long localCacheMaxBytes;
private final Duration expireAfterAccess;

/**
* Creates runtime resource settings.
*
* @param localCacheMaxBytes positive disk-space budget in bytes for local caches shared by
* all lookupers in the runtime; implementations without local disk caches may ignore
* this budget
* @param expireAfterAccess positive idle expiration for individual cached lookup files
*/
public LookupCacheOptions(long localCacheMaxBytes, Duration expireAfterAccess) {
checkArgument(localCacheMaxBytes > 0, "localCacheMaxBytes must be greater than 0.");
this.localCacheMaxBytes = localCacheMaxBytes;
this.expireAfterAccess =
checkNotNull(expireAfterAccess, "expireAfterAccess must not be null.");
checkArgument(
!expireAfterAccess.isNegative() && !expireAfterAccess.isZero(),
"expireAfterAccess must be greater than 0.");
}

/** Returns the runtime-wide disk-space budget for local caches, in bytes. */
public long localCacheMaxBytes() {
return localCacheMaxBytes;
}

/** Returns the idle expiration applied independently to each cached lookup file. */
public Duration expireAfterAccess() {
return expireAfterAccess;
}
}

/** Runtime context for creating a lake table lookuper. */
final class Context {
private final Configuration lakeConfiguration;
private final String cacheNamespace;
private final TableConfig tableConfig;
private final Runnable diskWriteGuard;

/**
* Creates a lookuper context.
*
* @param lakeConfiguration configuration of the lake storage for this lookuper
* @param cacheNamespace namespace identifying cache entries owned by this lookuper
* @param tableConfig configuration of the Fluss table
* @param diskWriteGuard guard invoked before creating a local lookup cache file
*/
public Context(
Configuration lakeConfiguration,
String cacheNamespace,
TableConfig tableConfig,
Runnable diskWriteGuard) {
this.lakeConfiguration =
checkNotNull(lakeConfiguration, "lakeConfiguration must not be null.");
this.cacheNamespace = checkNotNull(cacheNamespace, "cacheNamespace must not be null.");
this.tableConfig = checkNotNull(tableConfig, "tableConfig must not be null.");
this.diskWriteGuard = checkNotNull(diskWriteGuard, "diskWriteGuard must not be null.");
}

/** Returns the lake storage configuration for this lookuper. */
public Configuration lakeConfiguration() {
return lakeConfiguration;
}

/** Returns the namespace identifying cache entries owned by this lookuper. */
public String cacheNamespace() {
return cacheNamespace;
}

/** Returns the configuration of the Fluss table. */
public TableConfig tableConfig() {
return tableConfig;
}

/** Returns the guard invoked before creating a local lookup cache file. */
public Runnable diskWriteGuard() {
return diskWriteGuard;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import org.apache.fluss.config.Configuration;
import org.apache.fluss.exception.TableAlreadyExistException;
import org.apache.fluss.exception.TableNotExistException;
import org.apache.fluss.lake.lakestorage.LakeTableLookuperManager.LookupCacheOptions;
import org.apache.fluss.lake.source.LakeSource;
import org.apache.fluss.lake.writer.LakeTieringFactory;
import org.apache.fluss.metadata.TableChange;
Expand Down Expand Up @@ -139,13 +140,60 @@ public LakeSource<?> createLakeSource(TablePath tablePath) {
}

@Override
public LakeTableLookuper createLakeTableLookuper(
TablePath tablePath, LookuperContext context) {
public LakeTableLookuperManager createLakeTableLookuperManager(
String ioTmpDir, LookupCacheOptions options) {
try (TemporaryClassLoaderContext ignored = TemporaryClassLoaderContext.of(loader)) {
return new ClassLoaderFixingLakeTableLookuperManager(
inner.createLakeTableLookuperManager(ioTmpDir, options), loader);
}
}
}

static class ClassLoaderFixingLakeTableLookuperManager
implements LakeTableLookuperManager, WrappingProxy<LakeTableLookuperManager> {

private final LakeTableLookuperManager inner;
private final ClassLoader loader;

private ClassLoaderFixingLakeTableLookuperManager(
LakeTableLookuperManager inner, ClassLoader loader) {
this.inner = inner;
this.loader = loader;
}

@Override
public LakeTableLookuper createLakeTableLookuper(TablePath tablePath, Context context) {
try (TemporaryClassLoaderContext ignored = TemporaryClassLoaderContext.of(loader)) {
return new ClassLoaderFixingLakeTableLookuper(
inner.createLakeTableLookuper(tablePath, context), loader);
}
}

@Override
public void reconfigure(LookupCacheOptions options) {
try (TemporaryClassLoaderContext ignored = TemporaryClassLoaderContext.of(loader)) {
inner.reconfigure(options);
}
}

@Override
public long fileCacheCapacityEvictions() {
try (TemporaryClassLoaderContext ignored = TemporaryClassLoaderContext.of(loader)) {
return inner.fileCacheCapacityEvictions();
}
}

@Override
public void close() throws Exception {
try (TemporaryClassLoaderContext ignored = TemporaryClassLoaderContext.of(loader)) {
inner.close();
}
}

@Override
public LakeTableLookuperManager getWrappedDelegate() {
return inner;
}
}

static class ClassLoaderFixingLakeTableLookuper
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,8 +112,8 @@ public class MetricNames {
// for historical lookup cache
public static final String HISTORICAL_LOOKUP_CACHE_DISK_SIZE = "lookupCacheDiskSize";
public static final String HISTORICAL_LOOKUP_CACHE_TABLE_COUNT = "lookupCacheTableCount";
public static final String HISTORICAL_LOOKUP_CACHE_CAPACITY_EVICTIONS =
"lookupCacheCapacityEvictions";
public static final String HISTORICAL_LOOKUP_CACHE_FILE_CAPACITY_EVICTIONS =
"lookupCacheFileCapacityEvictions";

// --------------------------------------------------------------------------------------------
// metrics for user
Expand Down
Loading
Loading