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
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.paimon.fs.RemoteIterator;
import org.apache.paimon.fs.SeekableInputStream;
import org.apache.paimon.fs.VectoredReadable;
import org.apache.paimon.options.Options;
import org.apache.paimon.utils.Pair;

import org.apache.paimon.shade.guava30.com.google.common.collect.Lists;
Expand All @@ -37,6 +38,8 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import javax.annotation.Nullable;

import java.io.IOException;
import java.io.UncheckedIOException;
import java.net.URI;
Expand All @@ -58,6 +61,7 @@ public abstract class HadoopCompliantFileIO implements FileIO {
private static final Logger LOG = LoggerFactory.getLogger(HadoopCompliantFileIO.class);

private static final long serialVersionUID = 1L;
private static final String JINDO_CACHE_RPC_ADDRESS = "fs.jindocache.namespace.rpc.address";

/// Detailed cache strategies are retrieved from REST server.
private static final String META_CACHE_ENABLED_TAG = "meta";
Expand All @@ -72,10 +76,16 @@ public abstract class HadoopCompliantFileIO implements FileIO {
protected transient volatile Map<String, Pair<JindoHadoopSystem, String>> fsMap;
protected transient volatile Map<String, Pair<JindoHadoopSystem, String>> jindoCacheFsMap;

// Non-null when endpoint routing is configured instead of JindoCache RPC.
@Nullable IoCacheRouting cacheRouting;

// Only enable cache for path which is generated with uuid
private List<String> cacheWhitelistPaths = new ArrayList<>();

boolean shouldCache(Path path) {
if (cacheRouting != null) {
return cacheRouting.targetOf(path) != null;
}
if (cacheWhitelistPaths.isEmpty()) {
return true;
}
Expand All @@ -90,28 +100,41 @@ boolean shouldCache(Path path) {

@Override
public void configure(CatalogContext context) {
// Process file io cache configuration
if (!context.options().get(IO_CACHE_ENABLED)
|| context.options().get(IO_CACHE_POLICY) == null
|| context.options().get(IO_CACHE_POLICY).contains(DISABLE_CACHE_TAG)) {
Options options = context.options();
// Configure cache for the selected backend.
if (!options.get(IO_CACHE_ENABLED) || options.get(IO_CACHE_POLICY) == null) {
LOG.debug(
"Cache is disabled with io-cache.enabled={}, io-cache.policy={}",
context.options().get(IO_CACHE_ENABLED),
context.options().get(IO_CACHE_POLICY));
options.get(IO_CACHE_ENABLED),
options.get(IO_CACHE_POLICY));
return;
}
// Enable file io cache
if (context.options().get("fs.jindocache.namespace.rpc.address") == null) {
if (options.get(JINDO_CACHE_RPC_ADDRESS) == null) {
cacheRouting = IoCacheRouting.create(options);
if (cacheRouting == null) {
LOG.debug(
"FileIO cache endpoint routing is not configured for io-cache.policy={}",
options.get(IO_CACHE_POLICY));
return;
}
metaCacheEnabled = cacheRouting.metaCacheEnabled();
readCacheEnabled = cacheRouting.readCacheEnabled();
writeCacheEnabled = cacheRouting.writeCacheEnabled();
LOG.info(
"FileIO cache is enabled but JindoCache RPC address is not set, fallback to no-cache");
"Cache endpoints enabled: meta {}, read {}, write {}, {}",
metaCacheEnabled,
readCacheEnabled,
writeCacheEnabled,
cacheRouting);
} else {
metaCacheEnabled =
context.options().get(IO_CACHE_POLICY).contains(META_CACHE_ENABLED_TAG);
readCacheEnabled =
context.options().get(IO_CACHE_POLICY).contains(READ_CACHE_ENABLED_TAG);
writeCacheEnabled =
context.options().get(IO_CACHE_POLICY).contains(WRITE_CACHE_ENABLED_TAG);
String whitelist = context.options().get(IO_CACHE_WHITELIST_PATH);
// Keep legacy JindoCache policy matching unchanged; endpoint routing uses exact tokens.
if (options.get(IO_CACHE_POLICY).contains(DISABLE_CACHE_TAG)) {
return;
}
metaCacheEnabled = options.get(IO_CACHE_POLICY).contains(META_CACHE_ENABLED_TAG);
readCacheEnabled = options.get(IO_CACHE_POLICY).contains(READ_CACHE_ENABLED_TAG);
writeCacheEnabled = options.get(IO_CACHE_POLICY).contains(WRITE_CACHE_ENABLED_TAG);
String whitelist = options.get(IO_CACHE_WHITELIST_PATH);
if (!whitelist.equals("*")) {
cacheWhitelistPaths = Lists.newArrayList(whitelist.split(","));
}
Expand Down
Loading
Loading