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
6 changes: 6 additions & 0 deletions docs/docs/maintenance/filesystems.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -333,6 +333,12 @@ Download [paimon-jindo-@@VERSION@@.jar](https://repository.apache.org/snapshots/

</Unstable>

### Credentials Provider

`fs.oss.credentials.provider` can replace `fs.oss.accessKeyId` and `fs.oss.accessKeySecret` with a `com.aliyun.oss.common.auth.CredentialsProvider` class, such as `com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider`.

With a REST catalog, the data token is reloaded before it expires, also for streams that are already open.

### Server-Side Encryption

Paimon can stamp OSS server-side-encryption headers on the writes it performs (`PutObject`, server-side `CopyObject` used by rename/commit, and multipart-upload initiation) via three keys that map 1:1 to the OSS request headers:
Expand Down
115 changes: 102 additions & 13 deletions paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

package org.apache.paimon.rest;

import org.apache.paimon.annotation.VisibleForTesting;
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.data.BlobDescriptor;
Expand Down Expand Up @@ -47,6 +48,7 @@
import java.io.IOException;
import java.time.Duration;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

Expand All @@ -67,7 +69,7 @@ public class RESTTokenFileIO implements FileIO {
.defaultValue(false)
.withDescription("Whether to support data token provided by the REST server.");

private static final Cache<RESTToken, FileIO> FILE_IO_CACHE =
private static final Cache<DelegateKey, FileIO> FILE_IO_CACHE =
Caffeine.newBuilder()
.maximumSize(1000)
.expireAfterAccess(10, TimeUnit.HOURS)
Expand All @@ -92,6 +94,12 @@ public static void setFileIOCacheMaximumSize(long maximumSize) {
.setMaximum(maximumSize);
}

@VisibleForTesting
static void invalidateFileIOCache() {
FILE_IO_CACHE.invalidateAll();
FILE_IO_CACHE.cleanUp();
}

static long fileIOCacheMaximumSize() {
return FILE_IO_CACHE
.policy()
Expand All @@ -112,6 +120,9 @@ static long fileIOCacheMaximumSize() {
// Server again after serialization
private volatile RESTToken token;

// The table and catalog options a delegate refreshes its token with, computed once.
private transient volatile DelegateContext delegateContext;

public RESTTokenFileIO(
CatalogContext catalogContext, RESTApi apiInstance, Identifier identifier, Path path) {
this.catalogContext = catalogContext;
Expand Down Expand Up @@ -245,28 +256,33 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc
+ "REST credential lifetime after refresh.");
}

FileIO fileIO = FILE_IO_CACHE.getIfPresent(currentToken);
// A delegate reloads its token as a table and a catalog user, so it is shared only
// with callers that have the same ones.
DelegateKey key = new DelegateKey(currentToken, delegateContext());
FileIO fileIO = FILE_IO_CACHE.getIfPresent(key);
if (fileIO != null) {
return new FileIOWithToken(fileIO, currentToken);
}

synchronized (FILE_IO_CACHE) {
fileIO = FILE_IO_CACHE.getIfPresent(currentToken);
fileIO = FILE_IO_CACHE.getIfPresent(key);
if (fileIO != null) {
return new FileIOWithToken(fileIO, currentToken);
}

Options options = catalogContext.options();
options = new Options(RESTUtil.merge(options.toMap(), currentToken.token()));
options.set(FILE_IO_ALLOW_CACHE, false);
// Record whose data token the delegate holds, so it can reload it from the catalog.
RESTTokenRefresher.configure(options, tokenIdentifier(), currentToken.expireAtMillis());
CatalogContext context =
CatalogContext.create(
options,
catalogContext.hadoopConf(),
catalogContext.preferIO(),
catalogContext.fallbackIO());
fileIO = FileIO.get(path, context);
FILE_IO_CACHE.put(currentToken, fileIO);
FILE_IO_CACHE.put(key, fileIO);
return new FileIOWithToken(fileIO, currentToken);
}
}
Expand Down Expand Up @@ -300,6 +316,68 @@ private boolean shouldRefresh(long minimumValidityMillis) {
< Math.max(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, minimumValidityMillis);
}

/** The table and catalog options that a delegate reloads its token with. */
private static final class DelegateContext {

private final Identifier table;
private final Map<String, String> catalogOptions;
private final int hash;

private DelegateContext(Identifier table, Map<String, String> catalogOptions) {
this.table = table;
this.catalogOptions = catalogOptions;
this.hash = Objects.hash(table, catalogOptions);
}

@Override
public boolean equals(Object o) {
if (this == o) {
return true;
}
if (!(o instanceof DelegateContext)) {
return false;
}
DelegateContext that = (DelegateContext) o;
return hash == that.hash
&& table.equals(that.table)
&& catalogOptions.equals(that.catalogOptions);
}

@Override
public int hashCode() {
return hash;
}
}

/** Cache key of a delegate: its token and the context it reloads the token with. */
private static final class DelegateKey {

private final RESTToken token;
private final DelegateContext context;

private DelegateKey(RESTToken token, DelegateContext context) {
this.token = token;
this.context = context;
}

@Override
public boolean equals(Object o) {
if (this == o) {
return true;
}
if (!(o instanceof DelegateKey)) {
return false;
}
DelegateKey that = (DelegateKey) o;
return token.equals(that.token) && context.equals(that.context);
}

@Override
public int hashCode() {
return 31 * token.hashCode() + context.hashCode();
}
}

private static class FileIOWithToken {

private final FileIO fileIO;
Expand All @@ -316,15 +394,7 @@ private void refreshToken() {
if (apiInstance == null) {
apiInstance = new RESTApi(catalogContext.options(), false);
}
Identifier tableIdentifier = identifier;
if (identifier.isSystemTable()) {
tableIdentifier =
new Identifier(
identifier.getDatabaseName(),
identifier.getTableName(),
identifier.getBranchName());
}
GetTableTokenResponse response = apiInstance.loadTableToken(tableIdentifier);
GetTableTokenResponse response = apiInstance.loadTableToken(tokenIdentifier());
LOG.info(
"end refresh data token for identifier [{}] expiresAtMillis [{}]",
identifier,
Expand All @@ -336,6 +406,25 @@ private void refreshToken() {
response.getExpiresAtMillis());
}

private DelegateContext delegateContext() {
DelegateContext context = delegateContext;
if (context == null) {
context = new DelegateContext(tokenIdentifier(), catalogContext.options().toMap());
delegateContext = context;
}
return context;
}

private Identifier tokenIdentifier() {
if (identifier.isSystemTable()) {
return new Identifier(
identifier.getDatabaseName(),
identifier.getTableName(),
identifier.getBranchName());
}
return identifier;
}

private Map<String, String> mergeTokenWithCatalogOptions(Map<String, String> token) {
Map<String, String> newToken = Maps.newLinkedHashMap(token);
Options catalogOptions = catalogContext.options();
Expand Down
Loading
Loading