diff --git a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/HadoopCompliantFileIO.java b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/HadoopCompliantFileIO.java index 64fc3e07c708..d365e691efcd 100644 --- a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/HadoopCompliantFileIO.java +++ b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/HadoopCompliantFileIO.java @@ -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; @@ -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; @@ -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"; @@ -68,14 +72,21 @@ public abstract class HadoopCompliantFileIO implements FileIO { protected boolean metaCacheEnabled = false; protected boolean readCacheEnabled = false; protected boolean writeCacheEnabled = false; + protected boolean existsCacheEnabled = false; protected transient volatile Map> fsMap; protected transient volatile Map> 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 cacheWhitelistPaths = new ArrayList<>(); boolean shouldCache(Path path) { + if (cacheRouting != null) { + return cacheRouting.targetOf(path) != null; + } if (cacheWhitelistPaths.isEmpty()) { return true; } @@ -90,28 +101,44 @@ 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(); + existsCacheEnabled = cacheRouting.existsCacheEnabled(); LOG.info( - "FileIO cache is enabled but JindoCache RPC address is not set, fallback to no-cache"); + "Cache endpoints enabled: meta {}, read {}, write {}, exists {}, {}", + metaCacheEnabled, + readCacheEnabled, + writeCacheEnabled, + existsCacheEnabled, + 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); + existsCacheEnabled = metaCacheEnabled; + String whitelist = options.get(IO_CACHE_WHITELIST_PATH); if (!whitelist.equals("*")) { cacheWhitelistPaths = Lists.newArrayList(whitelist.split(",")); } @@ -194,7 +221,7 @@ public FileStatus next() throws IOException { @Override public boolean exists(Path path) throws IOException { org.apache.hadoop.fs.Path hadoopPath = path(path); - boolean shouldCache = metaCacheEnabled && shouldCache(path); + boolean shouldCache = existsCacheEnabled && shouldCache(path); LOG.debug("Exists should cache {} for path {}", shouldCache, path); return getFileSystem(hadoopPath, shouldCache).exists(hadoopPath); } diff --git a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/IoCacheRouting.java b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/IoCacheRouting.java new file mode 100644 index 000000000000..6ba00f224a05 --- /dev/null +++ b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/IoCacheRouting.java @@ -0,0 +1,360 @@ +/* + * 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.paimon.jindo; + +import org.apache.paimon.fs.Path; +import org.apache.paimon.options.Options; +import org.apache.paimon.utils.FileType; +import org.apache.paimon.utils.StringUtils; + +import javax.annotation.Nullable; + +import java.io.Serializable; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.EnumMap; +import java.util.EnumSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.regex.Pattern; +import java.util.stream.Collectors; + +import static org.apache.paimon.CoreOptions.CHANGELOG_FILE_PREFIX; +import static org.apache.paimon.CoreOptions.DATA_FILE_PREFIX; +import static org.apache.paimon.rest.RESTCatalogOptions.DLF_OSS_ENDPOINT; +import static org.apache.paimon.rest.RESTCatalogOptions.IO_CACHE_ENABLED; +import static org.apache.paimon.rest.RESTCatalogOptions.IO_CACHE_POLICY; + +/** + * Selects cache targets for immutable Paimon files using {@code io-cache.targets}, {@code + * io-cache.routes} and {@code io-cache.whitelist}. The endpoint mode also carries the normalized + * {@code io-cache.policy} read, metadata, existence and write switches. + */ +final class IoCacheRouting implements Serializable { + + private static final long serialVersionUID = 1L; + + private static final String META_CACHE_ENABLED_TAG = "meta"; + private static final String READ_CACHE_ENABLED_TAG = "read"; + private static final String WRITE_CACHE_ENABLED_TAG = "write"; + // exists asks whether a file is still there, which a cache of deleted files can get wrong + private static final String EXISTS_CACHE_ENABLED_TAG = "exists"; + private static final String DISABLE_CACHE_TAG = "none"; + + private static final Pattern TARGET_NAME = Pattern.compile("[a-z][a-z0-9-]*"); + + // Data files are named {prefix}{uuid}-{count}.{extension}; this matches what follows the + // prefix. + private static final Pattern DATA_FILE_SUFFIX = + Pattern.compile( + "[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}-[0-9]+\\..+"); + + private static final String MANIFEST_SIDECAR_SUFFIX = ".avro.sidecar"; + + @Nullable private final String ossEndpoint; + private final List targets; + private final Map targetByType; + private final List dataFilePrefixes; + private final boolean metaCacheEnabled; + private final boolean readCacheEnabled; + private final boolean writeCacheEnabled; + private final boolean existsCacheEnabled; + + private IoCacheRouting( + @Nullable String ossEndpoint, + List targets, + Map targetByType, + List dataFilePrefixes, + boolean metaCacheEnabled, + boolean readCacheEnabled, + boolean writeCacheEnabled, + boolean existsCacheEnabled) { + this.ossEndpoint = ossEndpoint; + this.targets = targets; + this.targetByType = targetByType; + this.dataFilePrefixes = dataFilePrefixes; + this.metaCacheEnabled = metaCacheEnabled; + this.readCacheEnabled = readCacheEnabled; + this.writeCacheEnabled = writeCacheEnabled; + this.existsCacheEnabled = existsCacheEnabled; + } + + /** Creates endpoint routing from the options, or null when routing is not enabled. */ + @Nullable + static IoCacheRouting create(Options options) { + if (!options.get(IO_CACHE_ENABLED)) { + return null; + } + Set policy = parsePolicy(options.get(IO_CACHE_POLICY)); + boolean metaCache = policy.contains(META_CACHE_ENABLED_TAG); + boolean readCache = policy.contains(READ_CACHE_ENABLED_TAG); + boolean writeCache = policy.contains(WRITE_CACHE_ENABLED_TAG); + boolean existsCache = policy.contains(EXISTS_CACHE_ENABLED_TAG); + if (policy.contains(DISABLE_CACHE_TAG) + || !(metaCache || readCache || writeCache || existsCache)) { + return null; + } + // a client-side dlf.oss-endpoint replaces every endpoint the token vends + if (!StringUtils.isNullOrWhitespaceOnly(options.get(DLF_OSS_ENDPOINT.key()))) { + return null; + } + Map urls = declaredTargets(options); + if (urls.isEmpty()) { + return null; + } + Map byName = new LinkedHashMap<>(); + urls.forEach( + (name, url) -> { + if (!StringUtils.isNullOrWhitespaceOnly(url)) { + byName.put(name, target(options, name, url.trim())); + } + }); + + String whitelist = options.get("io-cache.whitelist"); + Set types = + whitelist == null ? EnumSet.allOf(FileType.class) : fileTypes(whitelist); + String routes = options.get("io-cache.routes"); + Map targetByType = new EnumMap<>(FileType.class); + for (FileType type : types) { + CacheTarget target = + routes == null + ? byName.values().stream().findFirst().orElse(null) + : byName.get(routedName(routes, type)); + if (target != null) { + targetByType.put(type, target); + } + } + + String origin = options.get("io-cache.origin.endpoint"); + return new IoCacheRouting( + StringUtils.isNullOrWhitespaceOnly(origin) + ? options.get("fs.oss.endpoint") + : origin, + new ArrayList<>(byName.values()), + targetByType, + dataFilePrefixes(options), + metaCache, + readCache, + writeCache, + existsCache); + } + + private static Set parsePolicy(@Nullable String value) { + return value == null + ? Collections.emptySet() + : Arrays.stream(value.split(",")) + .map(token -> token.trim().toLowerCase(Locale.ROOT)) + .collect(Collectors.toSet()); + } + + boolean metaCacheEnabled() { + return metaCacheEnabled; + } + + boolean readCacheEnabled() { + return readCacheEnabled; + } + + boolean writeCacheEnabled() { + return writeCacheEnabled; + } + + boolean existsCacheEnabled() { + return existsCacheEnabled; + } + + /** The OSS endpoint used by requests that do not use a cache target. */ + @Nullable + String ossEndpoint() { + return ossEndpoint; + } + + Collection targets() { + return targets; + } + + /** The cache target selected for a file, or null when it does not use a cache target. */ + @Nullable + CacheTarget targetOf(Path path) { + // Only OSS paths can use cache targets. + if (!"oss".equals(path.toUri().getScheme())) { + return null; + } + FileType type = cacheableType(path, dataFilePrefixes); + return type == null ? null : targetByType.get(type); + } + + /** The type of a file a cache may serve, or null when it must be read from OSS. */ + @Nullable + static FileType cacheableType(Path path, List dataFilePrefixes) { + String name = path.getName(); + if (isTemporary(name) || FileType.isMutable(path)) { + return null; + } + if (isDataFileName(name, dataFilePrefixes)) { + return name.endsWith(".index") ? FileType.FILE_INDEX : FileType.DATA; + } + if (isSequential(path)) { + return null; + } + FileType type = FileType.classify(path); + return type == FileType.DATA ? null : type; + } + + // Manifests, indexes and statistics share the uuid-count shape but have no extension. + private static boolean isDataFileName(String name, List dataFilePrefixes) { + if (name.endsWith(MANIFEST_SIDECAR_SUFFIX)) { + return false; + } + for (String prefix : dataFilePrefixes) { + if (name.startsWith(prefix) + && DATA_FILE_SUFFIX.matcher(name.substring(prefix.length())).matches()) { + return true; + } + } + return false; + } + + private static boolean isTemporary(String name) { + return name.endsWith(".tmp") || name.contains(".tmp-") || name.contains(".tmp."); + } + + // Sequential metadata can be read before it exists; a cached NotFound would hide new commits. + private static boolean isSequential(Path path) { + String name = path.getName(); + Path parent = path.getParent(); + return name.startsWith("snapshot-") + || name.startsWith("schema-") + || (name.startsWith("changelog-") + && parent != null + && "changelog".equals(parent.getName())); + } + + static List dataFilePrefixes(Options options) { + List prefixes = new ArrayList<>(); + prefixes.add(DATA_FILE_PREFIX.defaultValue()); + prefixes.add(CHANGELOG_FILE_PREFIX.defaultValue()); + for (String key : new String[] {DATA_FILE_PREFIX.key(), CHANGELOG_FILE_PREFIX.key()}) { + String prefix = options.get(key); + if (!StringUtils.isNullOrWhitespaceOnly(prefix) && !prefixes.contains(prefix)) { + prefixes.add(prefix); + } + } + return prefixes; + } + + // io-cache.targets with io-cache.target..endpoint, else io-cache.endpoint as "default" + private static Map declaredTargets(Options options) { + Map urls = new LinkedHashMap<>(); + String names = options.get("io-cache.targets"); + if (names == null) { + String url = options.get("io-cache.endpoint"); + if (url != null) { + urls.put("default", url); + } + return urls; + } + for (String name : names.toLowerCase(Locale.ROOT).split(",")) { + name = name.trim(); + if (TARGET_NAME.matcher(name).matches() && !urls.containsKey(name)) { + urls.put(name, options.get("io-cache.target." + name + ".endpoint")); + } + } + return urls; + } + + private static CacheTarget target(Options options, String name, String url) { + String prefix = "io-cache.target." + name + "."; + String region = options.get(prefix + "region"); + String pathStyle = options.get(prefix + "path-style-access"); + return new CacheTarget( + name, + url.contains("://") ? url : "https://" + url, + pathStyle != null && "true".equalsIgnoreCase(pathStyle.trim()), + StringUtils.isNullOrWhitespaceOnly(region) ? null : region); + } + + // io-cache.routes is "types=name;...": the first matching rule selects the target. + @Nullable + static String routedName(String routes, FileType type) { + for (String rule : routes.toLowerCase(Locale.ROOT).split(";")) { + int eq = rule.indexOf('='); + String name = eq > 0 ? rule.substring(eq + 1).trim() : ""; + if (!name.isEmpty() && fileTypes(rule.substring(0, eq)).contains(type)) { + return name; + } + } + return null; + } + + private static Set fileTypes(String value) { + return FileType.parseWhitelist(value.toLowerCase(Locale.ROOT)); + } + + @Override + public String toString() { + return "IoCacheRouting{" + targetByType + ", oss=" + ossEndpoint + "}"; + } + + /** A named cache target with its endpoint and addressing settings. */ + static final class CacheTarget implements Serializable { + + private static final long serialVersionUID = 1L; + + final String name; + final String url; + final boolean pathStyleAccess; + @Nullable final String region; + + CacheTarget(String name, String url, boolean pathStyleAccess, @Nullable String region) { + this.name = name; + this.url = url; + this.pathStyleAccess = pathStyleAccess; + this.region = region; + } + + @Override + public boolean equals(Object o) { + if (!(o instanceof CacheTarget)) { + return false; + } + CacheTarget that = (CacheTarget) o; + return name.equals(that.name) + && url.equals(that.url) + && pathStyleAccess == that.pathStyleAccess + && Objects.equals(region, that.region); + } + + @Override + public int hashCode() { + return Objects.hash(name, url, pathStyleAccess, region); + } + + @Override + public String toString() { + return name + "=" + url; + } + } +} diff --git a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java index 71468672af9f..0847a864ec28 100644 --- a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java +++ b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java @@ -24,6 +24,7 @@ import org.apache.paimon.fs.HadoopOptionsProvider; import org.apache.paimon.fs.Path; import org.apache.paimon.fs.TwoPhaseOutputStream; +import org.apache.paimon.jindo.IoCacheRouting.CacheTarget; import org.apache.paimon.options.Options; import org.apache.paimon.plugin.PluginLoader; import org.apache.paimon.utils.IOUtils; @@ -40,6 +41,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; @@ -72,6 +75,13 @@ public class JindoFileIO extends HadoopCompliantFileIO implements HadoopOptionsP private static final String OSS_ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; private static final String OSS_SECURITY_TOKEN = "fs.oss.securityToken"; private static final String OSS_SHOW_DIR_TIMESTAMP = "fs.oss.show-dir-timestamp"; + private static final String OSS_ENDPOINT = "fs.oss.endpoint"; + private static final String OSS_HTTPS_ENABLE = "fs.oss.https.enable"; + private static final String OSS_SECOND_LEVEL_DOMAIN_ENABLE = + "fs.oss.second.level.domain.enable"; + private static final String OSS_REGION = "fs.oss.region"; + // JindoSDK options of a DLF cache cluster, which the OSS endpoint must not use + private static final String OSS_DLF_CACHE_PREFIX = "fs.oss.dlf-cache."; private static final Map CASE_SENSITIVE_KEYS = new HashMap() { @@ -91,6 +101,8 @@ public class JindoFileIO extends HadoopCompliantFileIO implements HadoopOptionsP private Options hadoopOptions; private Options hadoopOptionsWithCache; + private Map cacheTargetOptions; + transient volatile Map> cacheTargetFsMap; private boolean allowCache = true; private transient BlobPresigner blobPresigner; @@ -163,6 +175,38 @@ public void configure(CatalogContext context) { hadoopOptionsWithCache.set("fs.oss.read.profile.columnar.use-pread", "false"); hadoopOptionsWithCache.set( "fs.jindocache.read.profile.columnar.readahead.pread.enable", "false"); + + if (cacheRouting != null) { + cacheTargetOptions = new HashMap<>(); + for (CacheTarget target : cacheRouting.targets()) { + Options options = withEndpoint(hadoopOptions, target.url); + options.set(OSS_SECOND_LEVEL_DOMAIN_ENABLE, String.valueOf(target.pathStyleAccess)); + if (target.region != null) { + options.set(OSS_REGION, target.region); + } + cacheTargetOptions.put(target.name, options); + } + // fs.oss.endpoint may name a cache for older clients, so set the OSS endpoint again + hadoopOptions = withEndpoint(hadoopOptions, cacheRouting.ossEndpoint()); + hadoopOptions.keySet().removeIf(key -> key.startsWith(OSS_DLF_CACHE_PREFIX)); + } + } + + /** Copies the options with an endpoint: its host[:port], and https from its scheme if any. */ + static Options withEndpoint(Options base, @Nullable String endpoint) { + Options options = new Options(base.toMap()); + if (endpoint == null) { + return options; + } + int schemeEnd = endpoint.indexOf("://"); + if (schemeEnd >= 0) { + String scheme = endpoint.substring(0, schemeEnd); + options.set(OSS_HTTPS_ENABLE, String.valueOf("https".equalsIgnoreCase(scheme))); + endpoint = endpoint.substring(schemeEnd + 3); + } + int pathStart = endpoint.indexOf('/'); + options.set(OSS_ENDPOINT, pathStart >= 0 ? endpoint.substring(0, pathStart) : endpoint); + return options; } /** @@ -181,29 +225,40 @@ public Options hadoopOptions(Path path, String opType) { shouldCache = writeCacheEnabled && shouldCache(path); } else if (opType.equalsIgnoreCase("meta")) { shouldCache = metaCacheEnabled && shouldCache(path); + } else if (opType.equalsIgnoreCase("exists")) { + shouldCache = existsCacheEnabled && shouldCache(path); } - if (shouldCache) { - return hadoopOptionsWithCache; - } else { + if (!shouldCache) { return hadoopOptions; + } else if (cacheRouting != null) { + return cacheTargetOptions.get(cacheRouting.targetOf(path).name); + } else { + return hadoopOptionsWithCache; } } @Override public TwoPhaseOutputStream newTwoPhaseOutputStream(Path path, boolean overwrite) throws IOException { - if (!overwrite && this.exists(path)) { + org.apache.hadoop.fs.Path hadoopPath = path(path); + // The check is part of the write, so a cache endpoint never sees the file before it exists. + boolean exists = + cacheRouting == null + ? this.exists(path) + : getFileSystem(hadoopPath, false).exists(hadoopPath); + if (!overwrite && exists) { throw new IOException("File " + path + " already exists."); } - org.apache.hadoop.fs.Path hadoopPath = path(path); - Pair pair = getFileSystemPair(hadoopPath, false); + boolean viaTarget = cacheRouting != null && writeCacheEnabled && shouldCache(path); + Pair pair = getFileSystemPair(hadoopPath, viaTarget); JindoHadoopSystem fs = pair.getKey(); JindoMpuStore mpuStore = fs.getMpuStore(hadoopPath); if (mpuStore == null) { LOG.debug( "Jindo multipart upload is unavailable for {}, falling back to rename commit.", path); - return super.newTwoPhaseOutputStream(path, overwrite); + // With cache routing the check above already ran on OSS; do not repeat it on a cache. + return super.newTwoPhaseOutputStream(path, overwrite || cacheRouting != null); } return new JindoTwoPhaseOutputStream( new JindoMultiPartUpload(mpuStore, fs.getWorkingDirectory()), hadoopPath, path); @@ -256,9 +311,44 @@ private static synchronized PluginLoader getLoader() { @Override protected Pair createFileSystem( org.apache.hadoop.fs.Path path, boolean enableCache) { + return createFileSystem(path, enableCache ? hadoopOptionsWithCache : hadoopOptions); + } + + Pair createFileSystem( + org.apache.hadoop.fs.Path path, CacheTarget target) { + return createFileSystem(path, cacheTargetOptions.get(target.name)); + } + + @Override + protected Pair getFileSystemPair( + org.apache.hadoop.fs.Path path, boolean enableCache) throws IOException { + if (!enableCache || cacheRouting == null) { + return super.getFileSystemPair(path, enableCache); + } + CacheTarget target = cacheRouting.targetOf(new Path(path.toUri())); + if (target == null) { + return super.getFileSystemPair(path, false); + } + if (cacheTargetFsMap == null) { + synchronized (this) { + if (cacheTargetFsMap == null) { + cacheTargetFsMap = new ConcurrentHashMap<>(); + } + } + } + String authority = path.toUri().getAuthority(); + String key = target.name + "/" + (authority == null ? "DEFAULT" : authority); + try { + return cacheTargetFsMap.computeIfAbsent(key, k -> createFileSystem(path, target)); + } catch (UncheckedIOException e) { + throw e.getCause(); + } + } + + private Pair createFileSystem( + org.apache.hadoop.fs.Path path, Options options) { final String scheme = path.toUri().getScheme(); final String authority = path.toUri().getAuthority(); - Options options = enableCache ? hadoopOptionsWithCache : hadoopOptions; Supplier> supplier = () -> { Configuration hadoopConf = new Configuration(false); @@ -314,8 +404,14 @@ public synchronized void close() { } } if (!allowCache) { - fsMap.values().stream().map(Pair::getKey).forEach(IOUtils::closeQuietly); - fsMap.clear(); + if (fsMap != null) { + fsMap.values().stream().map(Pair::getKey).forEach(IOUtils::closeQuietly); + fsMap.clear(); + } + if (cacheTargetFsMap != null) { + cacheTargetFsMap.values().stream().map(Pair::getKey).forEach(IOUtils::closeQuietly); + cacheTargetFsMap.clear(); + } } } diff --git a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/IoCacheRoutingTest.java b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/IoCacheRoutingTest.java new file mode 100644 index 000000000000..f53595f90ff2 --- /dev/null +++ b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/IoCacheRoutingTest.java @@ -0,0 +1,532 @@ +/* + * 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.paimon.jindo; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.Path; +import org.apache.paimon.jindo.IoCacheRouting.CacheTarget; +import org.apache.paimon.options.Options; +import org.apache.paimon.utils.FileType; +import org.apache.paimon.utils.InstantiationUtil; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +import java.util.Arrays; +import java.util.List; + +import static org.apache.paimon.utils.FileType.BUCKET_INDEX; +import static org.apache.paimon.utils.FileType.DATA; +import static org.apache.paimon.utils.FileType.FILE_INDEX; +import static org.apache.paimon.utils.FileType.GLOBAL_INDEX; +import static org.apache.paimon.utils.FileType.META; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link IoCacheRouting}. */ +public class IoCacheRoutingTest { + + private static final String TABLE_ROOT = "oss://bkt/db1.db/t1"; + private static final String UUID = "8b1f7c2e-3a4d-4e5f-9a0b-1c2d3e4f5a6b"; + private static final Path DATA_FILE = + new Path(TABLE_ROOT + "/dt=2026-09-30/bucket-0/data-" + UUID + "-0.orc"); + private static final Path MANIFEST = new Path(TABLE_ROOT + "/manifest/manifest-" + UUID + "-0"); + + // paths below are relative to TABLE_ROOT, and {uuid} stands for UUID + private static final String DATA_PATH = "dt=1/bucket-0/data-{uuid}-0.parquet"; + private static final String MANIFEST_PATH = "manifest/manifest-{uuid}-0"; + private static final String INDEX_PATH = "index/index-{uuid}-0"; + private static final String GLOBAL_INDEX_PATH = "index/btree-global-index-{uuid}.index"; + private static final String TEMP_PATH = "dt=1/bucket-0/.data-{uuid}-0.parquet.{uuid}.tmp"; + private static final String SNAPSHOT_PATH = "snapshot/snapshot-12"; + private static final String LATEST_PATH = "snapshot/LATEST"; + + private static final String OSS = "https://oss-cn-hangzhou-internal.aliyuncs.com"; + private static final String CACHE = "http://cache.example.com"; + private static final String ACCEL = "https://accelerator.example.com"; + private static final String CLUSTER = "http://cluster.example.com"; + private static final String WRITE = "io-cache.policy=meta,read,write"; + private static final String EXISTS = "io-cache.policy=meta,read,exists"; + + @Test + public void testOneCacheTarget() { + Options options = single(); + assertEndpoint(options, "read", DATA_PATH, CACHE); + assertEndpoint(options, "meta", DATA_PATH, CACHE); + // exists asks whether a file is still there, so it needs its own policy token + assertEndpoint(options, "exists", DATA_PATH, OSS); + assertEndpoint(single(EXISTS), "exists", DATA_PATH, CACHE); + assertEndpoint(single(EXISTS), "exists", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "read", MANIFEST_PATH, CACHE); + assertEndpoint(options, "read", "manifest/manifest-list-{uuid}-1", CACHE); + assertEndpoint(options, "read", "oss://other-bkt/db1.db/t1/" + DATA_PATH, CACHE); + assertEndpoint(options, "read", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "meta", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "exists", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "read", LATEST_PATH, OSS); + assertEndpoint(options, "exists", LATEST_PATH, OSS); + assertEndpoint(options, "read", TEMP_PATH, OSS); + assertEndpoint(options, "read", "dt=1/bucket-0/000000_0", OSS); + assertEndpoint(options, "read", "dls://bkt/db1.db/t1/" + DATA_PATH, OSS); + // the whitelist has no index + assertEndpoint(options, "read", INDEX_PATH, OSS); + assertEndpoint(options, "write", DATA_PATH, OSS); + for (String op : Arrays.asList("list", "delete", "rename", "mkdirs")) { + assertEndpoint(options, op, DATA_PATH, OSS); + } + } + + @Test + public void testPolicyAndWhitelist() { + assertEndpoint(single("-io-cache.enabled"), "read", DATA_PATH, CACHE); + assertEndpoint(single("io-cache.enabled=false"), "exists", DATA_PATH, CACHE); + assertEndpoint(single("-io-cache.policy"), "exists", DATA_PATH, CACHE); + assertEndpoint(single("io-cache.policy=none"), "exists", DATA_PATH, CACHE); + + Options readOnly = single("io-cache.policy=read"); + assertEndpoint(readOnly, "read", DATA_PATH, CACHE); + assertEndpoint(readOnly, "meta", DATA_PATH, OSS); + assertEndpoint(readOnly, "exists", DATA_PATH, OSS); + Options writeOnly = single("io-cache.policy=write"); + assertEndpoint(writeOnly, "write", DATA_PATH, CACHE); + assertEndpoint(writeOnly, "exists", DATA_PATH, OSS); + Options existsOnly = single("io-cache.policy=exists"); + assertEndpoint(existsOnly, "exists", DATA_PATH, CACHE); + assertEndpoint(existsOnly, "meta", DATA_PATH, OSS); + + assertEndpoint(single("-io-cache.whitelist"), "read", INDEX_PATH, CACHE); + assertEndpoint(single("io-cache.whitelist=*"), "read", GLOBAL_INDEX_PATH, CACHE); + assertEndpoint(single("io-cache.whitelist=*"), "read", DATA_PATH + ".index", CACHE); + } + + @Test + public void testEndpoints() { + assertEndpoint(single("io-cache.endpoint= "), "read", DATA_PATH, OSS); + assertEndpoint( + single("io-cache.endpoint=", "io-cache.routes=*=default"), "read", DATA_PATH, OSS); + // the origin defaults to fs.oss.endpoint + String oss = "fs.oss.endpoint=" + OSS; + assertEndpoint(single(oss, "-io-cache.origin.endpoint"), "list", DATA_PATH, OSS); + assertEndpoint(single(oss, "io-cache.origin.endpoint= "), "list", DATA_PATH, OSS); + // a client-side endpoint turns routing off + String override = "https://oss-cn-hangzhou.aliyuncs.com"; + assertEndpoint(single("dlf.oss-endpoint=" + override), "read", DATA_PATH, override); + } + + @Test + public void testWritePolicy() { + Options options = single(WRITE); + assertEndpoint(options, "write", DATA_PATH, CACHE); + assertEndpoint(options, "write", MANIFEST_PATH, CACHE); + assertEndpoint(options, "write", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "write", LATEST_PATH, OSS); + assertEndpoint(options, "write", TEMP_PATH, OSS); + for (String op : Arrays.asList("list", "delete", "rename", "mkdirs")) { + assertEndpoint(options, op, DATA_PATH, OSS); + } + assertEndpoint(single(WRITE, "io-cache.whitelist=meta"), "write", DATA_PATH, OSS); + + assertEndpoint(multi(WRITE), "write", DATA_PATH, CLUSTER); + assertEndpoint(multi(WRITE), "write", MANIFEST_PATH, ACCEL); + } + + @Test + public void testTwoCacheTargets() { + Options options = multi(); + assertEndpoint(options, "read", MANIFEST_PATH, ACCEL); + assertEndpoint(options, "meta", MANIFEST_PATH, ACCEL); + assertEndpoint(options, "exists", MANIFEST_PATH, OSS); + assertEndpoint(multi(EXISTS), "exists", MANIFEST_PATH, ACCEL); + assertEndpoint(options, "read", DATA_PATH, CLUSTER); + assertEndpoint(multi(EXISTS), "exists", DATA_PATH, CLUSTER); + assertEndpoint(options, "read", INDEX_PATH, CLUSTER); + assertEndpoint(options, "read", GLOBAL_INDEX_PATH, CLUSTER); + assertEndpoint(options, "read", SNAPSHOT_PATH, OSS); + assertEndpoint(options, "write", MANIFEST_PATH, OSS); + for (String op : Arrays.asList("list", "delete", "rename", "mkdirs")) { + assertEndpoint(options, op, MANIFEST_PATH, OSS); + } + assertEndpoint(multi("io-cache.whitelist=meta"), "read", DATA_PATH, OSS); + } + + @Test + public void testTargetsAndRoutes() { + // without routes the first target with an endpoint takes every type + String noRoutes = "-io-cache.routes"; + assertEndpoint(multi(noRoutes), "read", DATA_PATH, ACCEL); + assertEndpoint(multi(noRoutes, "io-cache.endpoint=" + CACHE), "read", DATA_PATH, ACCEL); + assertEndpoint( + multi(noRoutes, "-io-cache.target.accel.endpoint"), "read", DATA_PATH, CLUSTER); + assertEndpoint( + multi(noRoutes, "io-cache.target.accel.endpoint= "), "read", DATA_PATH, CLUSTER); + assertEndpoint( + multi(noRoutes, "io-cache.targets=Bad_Name,cluster"), "read", DATA_PATH, CLUSTER); + assertEndpoint( + multi("io-cache.target.cluster.endpoint= " + CLUSTER + " "), + "read", + DATA_PATH, + CLUSTER); + + // a type without a usable route goes to the origin + assertEndpoint(multi("io-cache.routes=data=x"), "read", DATA_PATH, OSS); + assertEndpoint(multi("-io-cache.target.cluster.endpoint"), "read", DATA_PATH, OSS); + String metaAndData = "io-cache.routes=meta=accel;data=cluster"; + assertEndpoint(multi(metaAndData), "read", DATA_PATH + ".index", OSS); + + assertEndpoint(multi("io-cache.routes=*=cluster"), "read", MANIFEST_PATH, CLUSTER); + String firstWins = "io-cache.routes=data=accel;data,meta=cluster"; + assertEndpoint(multi(firstWins), "read", DATA_PATH, ACCEL); + String malformed = "io-cache.routes==cluster;meta=accel;junk"; + assertEndpoint(multi(malformed), "read", MANIFEST_PATH, ACCEL); + assertEndpoint(multi(malformed), "read", DATA_PATH, OSS); + } + + @Test + public void testDataFilesNamedLikeMetadata() { + Options manifest = multi("data-file.prefix=manifest-"); + assertEndpoint(manifest, "read", "dt=1/bucket-0/manifest-{uuid}-0.orc", CLUSTER); + assertEndpoint(manifest, "read", "dt=1/bucket-0/7f3a/manifest-{uuid}-0.orc", CLUSTER); + assertEndpoint(manifest, "read", MANIFEST_PATH, ACCEL); + assertEndpoint(manifest, "read", MANIFEST_PATH + ".avro.sidecar", ACCEL); + Options snapshot = multi("data-file.prefix=snapshot-"); + assertEndpoint(snapshot, "read", "dt=1/bucket-0/snapshot-{uuid}-0.orc", CLUSTER); + Options stat = multi("data-file.prefix=stat-"); + assertEndpoint(stat, "read", "dt=1/bucket-0/stat-{uuid}-0.orc", CLUSTER); + Options index = + multi("io-cache.routes=data=cluster;bucket-index=accel", "data-file.prefix=index-"); + assertEndpoint(index, "read", "dt=1/bucket-0/index-{uuid}-0.orc", CLUSTER); + assertEndpoint(index, "read", INDEX_PATH, ACCEL); + + String custom = "dt=1/bucket-0/custom-{uuid}-0.orc"; + assertEndpoint(single(), "read", custom, OSS); + assertEndpoint(single("data-file.prefix=custom-"), "read", custom, CACHE); + } + + @Test + public void testCacheableType() { + // named by sequence id or rewritten in place + for (String path : + Arrays.asList( + "snapshot/snapshot-12", + "snapshot/LATEST", + "snapshot/EARLIEST", + "branch/branch-dev/snapshot/LATEST", + "schema/schema-3", + "changelog/changelog-5", + "changelog/LATEST", + "tag/tag-2026-09-30", + "tag/tag-success-file/t1_SUCCESS", + "consumer/consumer-job1", + "service/service-primary-key-lookup", + "dt=1/_SUCCESS", + "oss://bkt/bucket-0/db/t/changelog/changelog-5")) { + assertType(path, null); + } + // temporary or not written by Paimon + for (String path : + Arrays.asList( + "snapshot/.snapshot-13.{uuid}.tmp", + "dt=1/bucket-0/.data-{uuid}-0.parquet.{uuid}.tmp", + "dt=1/bucket-0/data-{uuid}-0.parquet.tmp-{uuid}", + "dt=1/bucket-0/data-{uuid}-0.parquet.tmp.{uuid}", + "metadata/version-hint.text", + "metadata/v3.metadata.json", + "dt=1/bucket-0/000000_0", + "dt=1/part-00000-1a2b.snappy.parquet", + "dt=1/bucket-0/data-1.parquet", + "dt=1/bucket-0/custom-{uuid}-0.orc", + "README")) { + assertType(path, null); + } + + assertType("manifest/manifest-{uuid}-0", META); + assertType("manifest/manifest-list-{uuid}-1", META); + assertType("manifest/index-manifest-{uuid}-0", META); + assertType("manifest/manifest-{uuid}-0.avro.sidecar", META); + assertType("statistics/stat-{uuid}-0", META); + assertType("dt=1/bucket-0/data-{uuid}-0.parquet", DATA); + assertType("dt=1/bucket-0/changelog-{uuid}-0.parquet", DATA); + assertType("dt=1/bucket-0/data-{uuid}-0.blob", DATA); + assertType("dt=1/8b1f7c2e/data-{uuid}-0.parquet", DATA); + assertType("dt=1/bucket-0/data-{uuid}-0.parquet.index", FILE_INDEX); + assertType("index/index-{uuid}-0", BUCKET_INDEX); + assertType("index/btree-global-index-{uuid}.index", GLOBAL_INDEX); + assertType("index/lumina-global-index-{uuid}.index", GLOBAL_INDEX); + } + + @Test + public void testCacheableTypeWithFilePrefixes() { + assertType("dt=1/bucket-0/custom-{uuid}-0.orc", DATA, "data-file.prefix=custom-"); + assertType("dt=1/bucket-0/data-{uuid}-0.orc", DATA, "data-file.prefix=custom-"); + assertType("dt=1/bucket-0/cl-{uuid}-0.orc", DATA, "changelog-file.prefix=cl-"); + assertType("dt=1/bucket-0/ unknown.parquet", null, "data-file.prefix= "); + + // the name decides, so data files may be named like metadata + String manifest = "data-file.prefix=manifest-"; + assertType("dt=1/bucket-0/manifest-{uuid}-0.orc", DATA, manifest); + assertType("dt=1/bucket-0/7f3a/manifest-{uuid}-0.orc", DATA, manifest); + assertType("dt=1/bucket-0/manifest-{uuid}-0.orc.index", FILE_INDEX, manifest); + assertType("manifest/manifest-{uuid}-0", META, manifest); + assertType("manifest/manifest-{uuid}-0.avro.sidecar", META, manifest); + String index = "data-file.prefix=index-"; + assertType("dt=1/bucket-0/index-{uuid}-0.orc", DATA, index); + assertType("bucket-postpone/index-{uuid}-0.orc", DATA, index); + assertType("index/index-{uuid}-0", BUCKET_INDEX, index); + assertType("dt=1/bucket-0/index-{uuid}-0", BUCKET_INDEX, index); + String stat = "data-file.prefix=stat-"; + assertType("dt=1/bucket-0/stat-{uuid}-0.orc", DATA, stat); + assertType("postpone/stat-{uuid}-0.orc", DATA, stat); + assertType("statistics/stat-{uuid}-0", META, stat); + String snapshot = "data-file.prefix=snapshot-"; + assertType("dt=1/bucket-0/snapshot-{uuid}-0.orc", DATA, snapshot); + assertType("dt=1/bucket-0/snapshot-{uuid}-0.orc.index", FILE_INDEX, snapshot); + assertType("dt=1/bucket-0/schema-{uuid}-0.orc", DATA, "data-file.prefix=schema-"); + assertType( + "dt=1/bucket-0/global-index-{uuid}-0.orc.index", + FILE_INDEX, + "data-file.prefix=global-index-"); + + // names rewritten in place never use a cache, whatever the prefix + assertType("dt=1/bucket-0/tag-{uuid}-0.orc", null, "data-file.prefix=tag-"); + assertType("dt=1/bucket-0/consumer-{uuid}-0.orc", null, "data-file.prefix=consumer-"); + String warehouse = "oss://bkt/warehouse/bucket-0/db/t1/"; + assertType(warehouse + "tag/tag-file", null, "data-file.prefix=tag-"); + assertType(warehouse + "consumer/consumer-file", null, "data-file.prefix=consumer-"); + assertType(warehouse + "service/service-file", null, "data-file.prefix=service-"); + assertType(warehouse + "metadata/version-hint.text", null, "data-file.prefix=version-"); + assertType(warehouse + "branch/branch-feature", null, "data-file.prefix=branch-"); + } + + @Test + public void testTargetNames() { + Options options = singleTargetOptions(); + options.set("io-cache.targets", " Accel , cluster-1,accel,1x,a_b,,-x"); + options.set("io-cache.target.accel.endpoint", "http://a"); + options.set("io-cache.target.cluster-1.endpoint", "http://c"); + assertThat(IoCacheRouting.create(options).targets()) + .extracting(target -> target.name) + .containsExactly("accel", "cluster-1"); + + Options invalidNames = singleTargetOptions(); + invalidNames.set("io-cache.targets", "Bad_Name"); + assertThat(IoCacheRouting.create(invalidNames)).isNull(); + } + + @ParameterizedTest(name = "{0}") + @MethodSource("policyCases") + public void testCachePolicyUsesCompleteCaseInsensitiveTokens( + String policy, boolean readCache, boolean metaCache) { + Options options = singleTargetOptions(); + options.set("io-cache.policy", policy); + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(options)); + assertThat(host(fileIO.hadoopOptions(DATA_FILE, "read").get("fs.oss.endpoint"))) + .isEqualTo(readCache ? "cache:8080" : "oss-cn-hangzhou.aliyuncs.com"); + assertThat(host(fileIO.hadoopOptions(DATA_FILE, "meta").get("fs.oss.endpoint"))) + .isEqualTo( + metaCache + ? "cache:8080" + : readCache + ? "origin.example.com" + : "oss-cn-hangzhou.aliyuncs.com"); + } + + private static List policyCases() { + return Arrays.asList( + Arguments.of(" READ , Meta ", true, true), + Arguments.of("read,NONE", false, false), + Arguments.of("thread,metadata", false, false), + Arguments.of("read,nonetheless", true, false)); + } + + @Test + public void testRoutedName() { + String routes = " meta = Accel ; data,BUCKET-INDEX=cluster;=x;junk;file-index=;*=all;x=y"; + assertThat(IoCacheRouting.routedName(routes, FileType.META)).isEqualTo("accel"); + assertThat(IoCacheRouting.routedName(routes, FileType.DATA)).isEqualTo("cluster"); + assertThat(IoCacheRouting.routedName(routes, FileType.BUCKET_INDEX)).isEqualTo("cluster"); + assertThat(IoCacheRouting.routedName(routes, FileType.FILE_INDEX)).isEqualTo("all"); + assertThat(IoCacheRouting.routedName("meta=accel", FileType.DATA)).isNull(); + } + + @Test + public void testTargetSettings() { + Options options = twoTargetOptions(); + options.set("io-cache.target.accel.region", "cn-hangzhou"); + options.set("io-cache.target.cluster.path-style-access", " TRUE "); + options.set("io-cache.target.off.endpoint", ""); + options.set("io-cache.targets", "accel,off,cluster,missing"); + IoCacheRouting routing = IoCacheRouting.create(options); + + assertThat(routing.targets()) + .containsExactly( + new CacheTarget("accel", "https://accel.example.com", false, "cn-hangzhou"), + new CacheTarget("cluster", "http://10.0.0.1:8080", true, null)); + assertThat(routing.targetOf(DATA_FILE).name).isEqualTo("cluster"); + assertThat(routing.targetOf(MANIFEST).name).isEqualTo("accel"); + assertThat(routing.targetOf(new Path(TABLE_ROOT + "/snapshot/snapshot-1"))).isNull(); + + Options bare = singleTargetOptions(); + bare.set("io-cache.endpoint", "cache.example.com"); + assertThat(IoCacheRouting.create(bare).targetOf(DATA_FILE).url) + .isEqualTo("https://cache.example.com"); + } + + @Test + public void testNoIoCacheRouting() { + Options override = singleTargetOptions(); + override.set("dlf.oss-endpoint", "oss-cn-hangzhou.aliyuncs.com"); + assertThat(IoCacheRouting.create(override)).isNull(); + + Options disabled = singleTargetOptions(); + disabled.set("io-cache.enabled", "false"); + assertThat(IoCacheRouting.create(disabled)).isNull(); + + Options writeOnly = singleTargetOptions(); + writeOnly.set("io-cache.policy", "write"); + IoCacheRouting writes = IoCacheRouting.create(writeOnly); + assertThat(writes.writeCacheEnabled()).isTrue(); + assertThat(writes.readCacheEnabled() || writes.metaCacheEnabled()).isFalse(); + + Options none = singleTargetOptions(); + none.remove("io-cache.endpoint"); + assertThat(IoCacheRouting.create(none)).isNull(); + + // A declared target without an endpoint sends every request to origin. + Options empty = singleTargetOptions(); + empty.set("io-cache.endpoint", ""); + IoCacheRouting routing = IoCacheRouting.create(empty); + assertThat(routing.targets()).isEmpty(); + assertThat(routing.targetOf(DATA_FILE)).isNull(); + assertThat(routing.ossEndpoint()).isEqualTo("https://origin.example.com"); + } + + @Test + public void testOssEndpoint() { + Options options = singleTargetOptions(); + assertThat(IoCacheRouting.create(options).ossEndpoint()) + .isEqualTo("https://origin.example.com"); + options.remove("io-cache.origin.endpoint"); + assertThat(IoCacheRouting.create(options).ossEndpoint()) + .isEqualTo("oss-cn-hangzhou.aliyuncs.com"); + } + + @Test + public void testSerializable() throws Exception { + IoCacheRouting routing = + InstantiationUtil.clone( + IoCacheRouting.create(twoTargetOptions()), getClass().getClassLoader()); + assertThat(routing.targetOf(MANIFEST).name).isEqualTo("accel"); + assertThat(routing.targetOf(DATA_FILE).name).isEqualTo("cluster"); + } + + private static Options singleTargetOptions() { + Options options = new Options(); + options.set("fs.oss.endpoint", "oss-cn-hangzhou.aliyuncs.com"); + options.set("io-cache.enabled", "true"); + options.set("io-cache.endpoint", "http://cache:8080"); + options.set("io-cache.origin.endpoint", "https://origin.example.com"); + options.set("io-cache.policy", "meta,read"); + return options; + } + + private static Options twoTargetOptions() { + Options options = singleTargetOptions(); + options.set("io-cache.targets", "accel,cluster"); + options.set("io-cache.target.accel.endpoint", "https://accel.example.com"); + options.set("io-cache.target.cluster.endpoint", "http://10.0.0.1:8080"); + options.set("io-cache.routes", "meta=accel;data,bucket-index=cluster"); + return options; + } + + /** One cache target, which older clients also get as fs.oss.endpoint. */ + private static Options single(String... changes) { + Options options = new Options(); + options.set("fs.oss.endpoint", CACHE); + options.set("io-cache.enabled", "true"); + options.set("io-cache.endpoint", CACHE); + options.set("io-cache.origin.endpoint", OSS); + options.set("io-cache.policy", "meta,read"); + options.set("io-cache.whitelist", "meta,data"); + return with(options, changes); + } + + /** Metadata on an accelerator, data and indexes on a cache cluster. */ + private static Options multi(String... changes) { + Options options = new Options(); + options.set("fs.oss.endpoint", ACCEL); + options.set("io-cache.enabled", "true"); + options.set("io-cache.origin.endpoint", OSS); + options.set("io-cache.targets", "accel,cluster"); + options.set("io-cache.target.accel.endpoint", ACCEL); + options.set("io-cache.target.accel.region", "cn-hangzhou"); + options.set("io-cache.target.cluster.endpoint", CLUSTER); + options.set("io-cache.target.cluster.path-style-access", "true"); + options.set("io-cache.policy", "meta,read"); + options.set("io-cache.whitelist", "*"); + options.set( + "io-cache.routes", "meta=accel;data,bucket-index,global-index,file-index=cluster"); + return with(options, changes); + } + + // "key=value" sets an option and "-key" removes it + private static Options with(Options options, String... changes) { + for (String change : changes) { + if (change.startsWith("-")) { + options.remove(change.substring(1)); + } else { + int eq = change.indexOf('='); + options.set(change.substring(0, eq), change.substring(eq + 1)); + } + } + return options; + } + + private static void assertEndpoint(Options options, String op, String path, String expect) { + Options copy = new Options(options.toMap()); + // RESTTokenFileIO puts a client-side dlf.oss-endpoint into fs.oss.endpoint + String override = copy.get("dlf.oss-endpoint"); + if (override != null) { + copy.set("fs.oss.endpoint", override); + } + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(copy)); + assertThat(host(fileIO.hadoopOptions(path(path), op).get("fs.oss.endpoint"))) + .as("%s %s with %s", op, path, options.toMap()) + .isEqualTo(host(expect)); + } + + private static void assertType(String path, FileType type, String... changes) { + List prefixes = IoCacheRouting.dataFilePrefixes(with(new Options(), changes)); + assertThat(IoCacheRouting.cacheableType(path(path), prefixes)) + .as("%s with %s", path, Arrays.toString(changes)) + .isEqualTo(type); + } + + private static Path path(String path) { + String resolved = path.replace("{uuid}", UUID); + return new Path(resolved.contains("://") ? resolved : TABLE_ROOT + "/" + resolved); + } + + private static String host(String endpoint) { + String host = endpoint.contains("://") ? endpoint.split("://", 2)[1] : endpoint; + return host.endsWith("/") ? host.substring(0, host.length() - 1) : host; + } +} diff --git a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoIoCacheRoutingTest.java b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoIoCacheRoutingTest.java new file mode 100644 index 000000000000..f0ed08eb3746 --- /dev/null +++ b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoIoCacheRoutingTest.java @@ -0,0 +1,490 @@ +/* + * 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.paimon.jindo; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.Path; +import org.apache.paimon.jindo.IoCacheRouting.CacheTarget; +import org.apache.paimon.options.Options; +import org.apache.paimon.utils.Pair; + +import com.aliyun.jindodata.api.spec.JdoException; +import com.aliyun.jindodata.common.JindoHadoopSystem; +import org.apache.hadoop.fs.FSDataInputStream; +import org.apache.hadoop.fs.FSDataOutputStream; +import org.apache.hadoop.fs.FileStatus; +import org.apache.hadoop.fs.PositionedReadable; +import org.apache.hadoop.fs.Seekable; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.io.ByteArrayInputStream; +import java.io.FileNotFoundException; +import java.io.IOException; +import java.io.UncheckedIOException; +import java.net.SocketTimeoutException; +import java.net.URI; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ConcurrentHashMap; + +import static org.apache.paimon.options.CatalogOptions.FILE_IO_ALLOW_CACHE; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** Tests for the cache endpoints of {@link JindoFileIO}. */ +public class JindoIoCacheRoutingTest { + + private static final String TABLE = "oss://bkt/db1.db/t1"; + private static final String UUID = "8b1f7c2e-3a4d-4e5f-9a0b-1c2d3e4f5a6b"; + private static final Path DATA = new Path(TABLE + "/dt=1/bucket-0/data-" + UUID + "-1.orc"); + private static final Path MANIFEST = new Path(TABLE + "/manifest/manifest-1"); + private static final Path SNAPSHOT = new Path(TABLE + "/snapshot/snapshot-2"); + private static final Path LATEST = new Path(TABLE + "/snapshot/LATEST"); + private static final String ACCEL_HOST = "cn-hangzhou-j-internal.oss-data-acc.aliyuncs.com"; + private static final String CLUSTER_HOST = "10.0.0.1:8080"; + private static final String OSS_HOST = "oss-cn-hangzhou-internal.aliyuncs.com"; + private static final String DLF_CACHE_KEY = "fs.oss.dlf-cache.consistent-hash.enabled"; + + private JindoHadoopSystem accelFs; + private JindoHadoopSystem clusterFs; + private JindoHadoopSystem ossFs; + + @BeforeEach + public void before() throws IOException { + accelFs = mockFileSystem(); + clusterFs = mockFileSystem(); + ossFs = mockFileSystem(); + } + + @Test + public void testOptionsOfEachCacheTarget() { + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(cacheTargetOptions())); + + Options cluster = fileIO.hadoopOptions(DATA, "read"); + assertThat(cluster.get("fs.oss.endpoint")).isEqualTo(CLUSTER_HOST); + assertThat(cluster.get("fs.oss.https.enable")).isEqualTo("false"); + assertThat(cluster.get("fs.oss.second.level.domain.enable")).isEqualTo("true"); + assertThat(cluster.get("fs.oss.region")).isEqualTo("cn-hangzhou"); + assertThat(cluster.get(DLF_CACHE_KEY)).isEqualTo("true"); + assertThat(cluster.containsKey("fs.xengine")).isFalse(); + assertThat(fileIO.hadoopOptions(DATA, "meta")).isSameAs(cluster); + + Options accel = fileIO.hadoopOptions(MANIFEST, "meta"); + assertThat(accel.get("fs.oss.endpoint")).isEqualTo(ACCEL_HOST); + assertThat(accel.get("fs.oss.https.enable")).isEqualTo("true"); + assertThat(accel.get("fs.oss.second.level.domain.enable")).isEqualTo("false"); + assertThat(accel.get("fs.oss.region")).isEqualTo("cn-shanghai"); + + Options oss = fileIO.hadoopOptions(DATA, "write"); + assertThat(oss.get("fs.oss.endpoint")).isEqualTo(OSS_HOST); + assertThat(oss.get("fs.oss.https.enable")).isEqualTo("true"); + assertThat(oss.get("fs.oss.region")).isEqualTo("cn-hangzhou"); + assertThat(oss.containsKey(DLF_CACHE_KEY)).isFalse(); + for (Path path : + new Path[] {SNAPSHOT, LATEST, new Path(TABLE + "/dt=1/bucket-0/000000_0")}) { + assertThat(fileIO.hadoopOptions(path, "read")).as(path.toString()).isSameAs(oss); + } + assertThat(fileIO.hadoopOptions(DATA, "list")).isSameAs(oss); + } + + @Test + public void testOnlyReadsAndStatusUseIoCacheRouting() throws IOException { + JindoFileIO fileIO = configuredFileIO(cacheTargetOptions()); + Path copy = new Path(TABLE + "/dt=1/bucket-0/data-" + UUID + "-2.orc"); + + fileIO.newInputStream(DATA).close(); + fileIO.getFileSize(DATA); + fileIO.getFileStatus(MANIFEST); + fileIO.exists(DATA); + fileIO.exists(MANIFEST); + verify(clusterFs).open(any(org.apache.hadoop.fs.Path.class)); + verify(clusterFs).getFileStatus(any()); + verify(clusterFs).exists(any()); + verify(accelFs).getFileStatus(any()); + verify(accelFs).exists(any()); + + fileIO.exists(SNAPSHOT); + fileIO.newInputStream(SNAPSHOT).close(); + fileIO.getFileStatus(SNAPSHOT); + fileIO.newOutputStream(DATA, false); + fileIO.listStatus(DATA.getParent()); + fileIO.delete(DATA, false); + fileIO.mkdirs(DATA.getParent()); + fileIO.rename(DATA, copy); + verify(ossFs).exists(any()); + verify(ossFs).open(any(org.apache.hadoop.fs.Path.class)); + verify(ossFs).getFileStatus(any()); + verify(ossFs).create(any(), anyBoolean()); + verify(ossFs).listStatus(any(org.apache.hadoop.fs.Path.class)); + verify(ossFs).delete(any(), anyBoolean()); + verify(ossFs).mkdirs(any()); + verify(ossFs).rename(any(), any()); + + // the two-phase write checks a routable data file on OSS, not on its cache target + fileIO.tryToWriteAtomic(new Path(TABLE + "/snapshot/snapshot-3"), "{}"); + fileIO.newTwoPhaseOutputStream(copy, false); + verify(ossFs, times(3)).create(any(), anyBoolean()); + verify(ossFs, times(2)).exists(any()); + verify(ossFs, times(2)).rename(any(), any()); + verify(ossFs).getMpuStore(any()); + + verify(clusterFs).exists(any()); + for (JindoHadoopSystem fs : new JindoHadoopSystem[] {accelFs, clusterFs}) { + verify(fs, never()).create(any(), anyBoolean()); + verify(fs, never()).listStatus(any(org.apache.hadoop.fs.Path.class)); + verify(fs, never()).rename(any(), any()); + } + } + + @Test + public void testExistsUsesOssWithoutTheExistsToken() throws IOException { + Options options = cacheTargetOptions(); + options.set("io-cache.policy", "meta,read"); + JindoFileIO fileIO = configuredFileIO(options); + + fileIO.exists(DATA); + fileIO.exists(MANIFEST); + fileIO.getFileStatus(DATA); + verify(ossFs, times(2)).exists(any()); + verify(clusterFs).getFileStatus(any()); + verify(clusterFs, never()).exists(any()); + verify(accelFs, never()).exists(any()); + } + + @Test + public void testCacheTargetErrorsAreThrown() throws IOException { + JindoFileIO fileIO = configuredFileIO(cacheTargetOptions()); + IOException unavailable = jindoError(6503, "open failed: 503 Service Unavailable"); + FileNotFoundException notFound = new FileNotFoundException("404"); + when(clusterFs.open(any(org.apache.hadoop.fs.Path.class))).thenThrow(unavailable); + when(accelFs.getFileStatus(any())).thenThrow(notFound); + + assertThatThrownBy(() -> fileIO.newInputStream(DATA)).isSameAs(unavailable); + assertThatThrownBy(() -> fileIO.getFileStatus(MANIFEST)).isSameAs(notFound); + verify(ossFs, never()).open(any(org.apache.hadoop.fs.Path.class)); + verify(ossFs, never()).getFileStatus(any()); + } + + @Test + public void testCacheTargetInitializationErrorsStayChecked() throws IOException { + SocketTimeoutException timeout = new SocketTimeoutException("connect timed out"); + JindoFileIO fileIO = + new TestingJindoFileIO(accelFs, clusterFs, ossFs) { + @Override + Pair createFileSystem( + org.apache.hadoop.fs.Path path, CacheTarget target) { + // JindoFileIO wraps an initialization IOException the same way + throw new UncheckedIOException(timeout); + } + }; + fileIO.configure(CatalogContext.create(cacheTargetOptions())); + + assertThatThrownBy(() -> fileIO.newInputStream(DATA)).isSameAs(timeout); + assertThatThrownBy(() -> fileIO.getFileStatus(MANIFEST)).isSameAs(timeout); + assertThatThrownBy(() -> fileIO.exists(MANIFEST)).isSameAs(timeout); + verify(ossFs, never()).open(any(org.apache.hadoop.fs.Path.class)); + } + + @Test + public void testWritesFollowTheRouteWithWritePolicy() throws IOException { + Options options = cacheTargetOptions(); + options.set("io-cache.policy", "meta,read,write"); + JindoFileIO fileIO = configuredFileIO(options); + Path next = new Path(TABLE + "/dt=1/bucket-0/data-" + UUID + "-2.orc"); + + fileIO.newOutputStream(DATA, false); + fileIO.newOutputStream(MANIFEST, false); + fileIO.newOutputStream(SNAPSHOT, false); + fileIO.newTwoPhaseOutputStream(next, false); + assertThat(fileIO.hadoopOptions(DATA, "write").get("fs.oss.endpoint")) + .isEqualTo(CLUSTER_HOST); + verify(clusterFs).create(any(), anyBoolean()); + verify(accelFs).create(any(), anyBoolean()); + verify(clusterFs).getMpuStore(any()); + // the snapshot, the existence check before the two-phase write and the temp file of its + // rename fallback (no multipart upload here) stay on OSS + verify(ossFs, times(2)).create(any(), anyBoolean()); + verify(ossFs).exists(any()); + verify(clusterFs, never()).exists(any()); + + fileIO.delete(DATA, false); + fileIO.rename(DATA, next); + fileIO.listStatus(DATA.getParent()); + verify(ossFs).delete(any(), anyBoolean()); + verify(ossFs).rename(any(), any()); + verify(ossFs).listStatus(any(org.apache.hadoop.fs.Path.class)); + } + + @Test + public void testSingleCacheTarget() { + Options options = baseOptions(); + options.set("io-cache.endpoint", "http://" + CLUSTER_HOST); + options.set("io-cache.target.default.path-style-access", "true"); + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(options)); + + Options read = fileIO.hadoopOptions(DATA, "read"); + assertThat(read.get("fs.oss.endpoint")).isEqualTo(CLUSTER_HOST); + assertThat(read.get("fs.oss.second.level.domain.enable")).isEqualTo("true"); + assertThat(fileIO.hadoopOptions(MANIFEST, "meta")).isSameAs(read); + assertThat(fileIO.hadoopOptions(DATA, "write").get("fs.oss.endpoint")).isEqualTo(OSS_HOST); + } + + @Test + public void testJindoCacheTakesPrecedence() { + Options options = cacheTargetOptions(); + options.set("fs.jindocache.namespace.rpc.address", "rpc:8101"); + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(options)); + + assertThat(fileIO.cacheRouting).isNull(); + Options read = fileIO.hadoopOptions(DATA, "read"); + assertThat(read.get("fs.xengine")).isEqualTo("jindocache"); + assertThat(read.get("fs.oss.endpoint")).isEqualTo("http://" + CLUSTER_HOST); + Options write = fileIO.hadoopOptions(DATA, "write"); + assertThat(write.containsKey("fs.xengine")).isFalse(); + // legacy whitelist-path matching: "manifest" is in the default whitelist + assertThat(fileIO.hadoopOptions(MANIFEST, "meta")).isSameAs(read); + assertThat(fileIO.hadoopOptions(LATEST, "read")).isSameAs(write); + } + + @Test + public void testWithoutIoCacheRouting() { + Options override = cacheTargetOptions(); + override.set("dlf.oss-endpoint", "oss-cn-hangzhou.aliyuncs.com"); + override.set("fs.oss.endpoint", "oss-cn-hangzhou.aliyuncs.com"); + Options notEnabled = cacheTargetOptions(); + notEnabled.remove("io-cache.enabled"); + Options noPolicy = cacheTargetOptions(); + noPolicy.set("io-cache.policy", "prefetch"); + + for (Options options : new Options[] {override, notEnabled, noPolicy}) { + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(options)); + assertThat(fileIO.cacheRouting).isNull(); + for (String op : new String[] {"read", "meta", "write"}) { + Options hadoopOptions = fileIO.hadoopOptions(DATA, op); + assertThat(hadoopOptions.get("fs.oss.endpoint")) + .isEqualTo(options.get("fs.oss.endpoint")); + assertThat(hadoopOptions.containsKey("fs.oss.https.enable")).isFalse(); + } + } + } + + @Test + public void testWithEndpoint() { + Options base = new Options(); + base.set("fs.oss.https.enable", "false"); + + Options bare = JindoFileIO.withEndpoint(base, "oss.example.com"); + assertThat(bare.get("fs.oss.endpoint")).isEqualTo("oss.example.com"); + assertThat(bare.get("fs.oss.https.enable")).isEqualTo("false"); + + Options https = JindoFileIO.withEndpoint(base, "HTTPS://cache.example.com:443/"); + assertThat(https.get("fs.oss.endpoint")).isEqualTo("cache.example.com:443"); + assertThat(https.get("fs.oss.https.enable")).isEqualTo("true"); + + Options http = JindoFileIO.withEndpoint(base, "http://10.0.0.1:8080"); + assertThat(http.get("fs.oss.endpoint")).isEqualTo("10.0.0.1:8080"); + assertThat(http.get("fs.oss.https.enable")).isEqualTo("false"); + + assertThat(JindoFileIO.withEndpoint(base, null).toMap()).isEqualTo(base.toMap()); + } + + @Test + public void testFileSystemPerCacheTargetAndBucket() throws IOException { + List created = new ArrayList<>(); + JindoFileIO fileIO = + new JindoFileIO() { + @Override + Pair createFileSystem( + org.apache.hadoop.fs.Path path, CacheTarget target) { + created.add(target.name + "/" + path.toUri().getAuthority()); + return Pair.of(mock(JindoHadoopSystem.class), "oss"); + } + }; + fileIO.configure(CatalogContext.create(cacheTargetOptions())); + + Pair first = + fileIO.getFileSystemPair(hadoopPath(TABLE + "/manifest/manifest-a"), true); + assertThat(fileIO.getFileSystemPair(hadoopPath(TABLE + "/manifest/manifest-b"), true)) + .isSameAs(first); + fileIO.getFileSystemPair(hadoopPath("oss://other/t/manifest/manifest-c"), true); + fileIO.getFileSystemPair( + hadoopPath(TABLE + "/dt=1/bucket-0/data-" + UUID + "-3.orc"), true); + assertThat(created).containsExactly("accel/bkt", "accel/other", "cluster/bkt"); + } + + @Test + public void testCloseAfterOnlyCacheReads() throws IOException { + Options options = cacheTargetOptions(); + options.set(FILE_IO_ALLOW_CACHE, false); + JindoFileIO fileIO = new TestingJindoFileIO(accelFs, clusterFs, ossFs); + fileIO.configure(CatalogContext.create(options)); + fileIO.newInputStream(DATA).close(); + + fileIO.close(); + + verify(clusterFs).close(); + } + + @Test + public void testCloseClosesCacheTargetFileSystems() throws IOException { + Options options = cacheTargetOptions(); + options.set(FILE_IO_ALLOW_CACHE, false); + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(CatalogContext.create(options)); + fileIO.fsMap = new ConcurrentHashMap<>(); + fileIO.fsMap.put("bkt", Pair.of(ossFs, "oss")); + fileIO.cacheTargetFsMap = new ConcurrentHashMap<>(); + fileIO.cacheTargetFsMap.put("accel/bkt", Pair.of(accelFs, "oss")); + fileIO.cacheTargetFsMap.put("cluster/bkt", Pair.of(clusterFs, "oss")); + + fileIO.close(); + + for (JindoHadoopSystem fs : new JindoHadoopSystem[] {ossFs, accelFs, clusterFs}) { + verify(fs).close(); + } + assertThat(fileIO.cacheTargetFsMap).isEmpty(); + } + + private static Options baseOptions() { + Options options = new Options(); + options.set("fs.oss.endpoint", "http://" + CLUSTER_HOST); + options.set("fs.oss.accessKeyId", "ak"); + options.set("fs.oss.accessKeySecret", "sk"); + options.set("fs.oss.region", "cn-hangzhou"); + options.set("io-cache.enabled", "true"); + options.set("io-cache.origin.endpoint", "https://" + OSS_HOST); + options.set("io-cache.policy", "meta,read"); + options.set(DLF_CACHE_KEY, "true"); + return options; + } + + private static Options cacheTargetOptions() { + Options options = baseOptions(); + options.set("io-cache.targets", "accel,cluster"); + options.set("io-cache.target.accel.endpoint", "https://" + ACCEL_HOST); + options.set("io-cache.target.accel.region", "cn-shanghai"); + options.set("io-cache.target.cluster.endpoint", "http://" + CLUSTER_HOST); + options.set("io-cache.target.cluster.path-style-access", "true"); + options.set("io-cache.routes", "meta=accel;data,bucket-index=cluster"); + options.set("io-cache.policy", "meta,read,exists"); + return options; + } + + private JindoFileIO configuredFileIO(Options options) { + JindoFileIO fileIO = new TestingJindoFileIO(accelFs, clusterFs, ossFs); + fileIO.configure(CatalogContext.create(options)); + return fileIO; + } + + private static IOException jindoError(int code, String message) { + return new IOException(new JdoException(code, message)); + } + + // Hadoop's Path(String) needs commons-lang, which is not on this module's test classpath. + private static org.apache.hadoop.fs.Path hadoopPath(String path) { + return new org.apache.hadoop.fs.Path(URI.create(path)); + } + + private static JindoHadoopSystem mockFileSystem() throws IOException { + JindoHadoopSystem fs = mock(JindoHadoopSystem.class); + when(fs.open(any(org.apache.hadoop.fs.Path.class))) + .thenAnswer(invocation -> new FSDataInputStream(new BytesInput())); + when(fs.create(any(), anyBoolean())).thenReturn(mock(FSDataOutputStream.class)); + when(fs.getFileStatus(any())).thenReturn(mock(FileStatus.class)); + when(fs.rename(any(), any())).thenReturn(true); + return fs; + } + + private static class BytesInput extends ByteArrayInputStream + implements Seekable, PositionedReadable { + + private BytesInput() { + super(new byte[] {1, 2, 3}); + } + + @Override + public void seek(long position) { + pos = (int) position; + } + + @Override + public long getPos() { + return pos; + } + + @Override + public boolean seekToNewSource(long targetPos) { + return false; + } + + @Override + public int read(long position, byte[] buffer, int offset, int length) { + throw new UnsupportedOperationException(); + } + + @Override + public void readFully(long position, byte[] buffer, int offset, int length) { + throw new UnsupportedOperationException(); + } + + @Override + public void readFully(long position, byte[] buffer) { + throw new UnsupportedOperationException(); + } + } + + private static class TestingJindoFileIO extends JindoFileIO { + + private final JindoHadoopSystem accelFs; + private final JindoHadoopSystem clusterFs; + private final JindoHadoopSystem ossFs; + + private TestingJindoFileIO( + JindoHadoopSystem accelFs, JindoHadoopSystem clusterFs, JindoHadoopSystem ossFs) { + this.accelFs = accelFs; + this.clusterFs = clusterFs; + this.ossFs = ossFs; + } + + @Override + protected Pair createFileSystem( + org.apache.hadoop.fs.Path path, boolean enableCache) { + assertThat(enableCache).isFalse(); + return Pair.of(ossFs, "oss"); + } + + @Override + Pair createFileSystem( + org.apache.hadoop.fs.Path path, CacheTarget target) { + return Pair.of("accel".equals(target.name) ? accelFs : clusterFs, "oss"); + } + } +} diff --git a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/TestJindoCacheEnable.java b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/TestJindoCacheEnable.java index 19a3da7bd533..c09b9ce71153 100644 --- a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/TestJindoCacheEnable.java +++ b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/TestJindoCacheEnable.java @@ -147,6 +147,22 @@ public void testCacheDisabledWhenPolicyNull() { verifyCacheFlags(fileIO, false, false, false); } + @Test + public void testLegacyPolicyMatchingRemainsUnchanged() { + // JindoCache RPC policies intentionally keep their historical substring matching. + JindoFileIO fileIO = new JindoFileIO(); + fileIO.configure(createCatalogContext(true, "read,nonetheless", true, null)); + verifyCacheFlags(fileIO, false, false, false); + + fileIO = new JindoFileIO(); + fileIO.configure(createCatalogContext(true, "thread,metadata", true, null)); + verifyCacheFlags(fileIO, true, true, false); + + fileIO = new JindoFileIO(); + fileIO.configure(createCatalogContext(true, " READ , Meta ", true, null)); + verifyCacheFlags(fileIO, false, false, false); + } + @Test public void testCacheWhitelist() { // default config