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 @@ -28,7 +28,6 @@
import java.util.Optional;
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.apache.flink.configuration.ReadableConfig;
import org.apache.flink.table.catalog.AbstractCatalog;
import org.apache.flink.table.catalog.CatalogBaseTable;
import org.apache.flink.table.catalog.CatalogTable;
Expand All @@ -45,6 +44,7 @@
import org.apache.gravitino.flink.connector.PartitionConverter;
import org.apache.gravitino.flink.connector.SchemaAndTablePropertiesConverter;
import org.apache.gravitino.flink.connector.catalog.BaseCatalog;
import org.apache.gravitino.flink.connector.utils.PropertyUtils;
import org.apache.gravitino.rel.Dialects;
import org.apache.gravitino.rel.Representation;
import org.apache.gravitino.rel.SQLRepresentation;
Expand All @@ -53,9 +53,14 @@
import org.apache.gravitino.rel.expressions.distributions.Distribution;
import org.apache.gravitino.rel.expressions.distributions.Distributions;
import org.apache.gravitino.rel.expressions.distributions.Strategy;
import org.apache.hadoop.conf.Configuration;
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.flink.FlinkCatalog;
import org.apache.paimon.flink.FlinkCatalogFactory;
import org.apache.paimon.flink.FlinkFileIOLoader;
import org.apache.paimon.options.Options;
import org.apache.paimon.utils.HadoopUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -76,8 +81,8 @@ public class GravitinoPaimonCatalog extends BaseCatalog {
private static final Logger LOG = LoggerFactory.getLogger(GravitinoPaimonCatalog.class);

private final CatalogFactory.Context context;
// Mutable copy shared with BaseCatalog.catalogOptions so credential injection in open() is
// visible to the inner Paimon catalog context.
// Mutable copy shared with BaseCatalog.catalogOptions. The inner Paimon catalog is created from a
// sanitized copy so Hadoop-prefixed options can be moved to Hadoop configuration before logging.
private final Map<String, String> mutableOptions;
private AbstractCatalog paimonCatalog;

Expand Down Expand Up @@ -129,43 +134,77 @@ public void open() throws CatalogException {
super.open();
return;
}
Map<String, String> paimonOptions = new HashMap<>(mutableOptions);
Configuration hadoopConf =
HadoopUtils.getHadoopConfiguration(Options.fromMap(withoutHadoopOptions(paimonOptions)));
moveHadoopOptionsToConf(paimonOptions, hadoopConf);
try {
// Keep vended credentials as Paimon options because Paimon native S3/OSS FileIO loaders
// declare them as required options.
CredentialPropertyUtils.applyPaimonCredentials(
CredentialPropertyUtils.getCredentials(catalog()), mutableOptions);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CredentialPropertyUtils.applyPaimonCredentials may put credentials to mutableOptions.

We need to confirm one thing: should these properties be removed from paimonOptions?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CredentialPropertyUtils.applyPaimonCredentials may put credentials to mutableOptions.

We need to confirm one thing: should these properties be removed from paimonOptions?

Good point. These two sets of properties are intentionally handled differently.

The Hadoop filesystem options from catalog properties (hadoop.*, fs.*, and dfs.*) are removed from paimonOptions and moved into Hadoop Configuration.

The credentials injected by CredentialPropertyUtils.applyPaimonCredentials are vended credentials, and they need to stay in paimonOptions because Paimon native S3/OSS FileIO loaders declare these keys as required options. If we remove them from paimonOptions, native Paimon S3/OSS FileIO may fail to initialize.

I added a comment to make this distinction explicit.

CredentialPropertyUtils.getCredentials(catalog()), paimonOptions);
} catch (NoSuchCatalogException e) {
LOG.warn(
"Catalog '{}' not found in Gravitino during open(); credential injection skipped."
+ " This is expected during CREATE CATALOG.",
getName(),
e);
}
CatalogFactory.Context contextWithCredentials =
new CatalogFactory.Context() {
@Override
public String getName() {
return context.getName();
}

@Override
public Map<String, String> getOptions() {
return mutableOptions;
}

@Override
public ReadableConfig getConfiguration() {
return context.getConfiguration();
}

@Override
public ClassLoader getClassLoader() {
return context.getClassLoader();
}
};
this.paimonCatalog =
(AbstractCatalog) new FlinkCatalogFactory().createCatalog(contextWithCredentials);
this.paimonCatalog = createInnerCatalog(paimonOptions, hadoopConf);
super.open();
}

/**
* Creates the inner Paimon Flink catalog from sanitized Paimon options and Hadoop configuration.
*
* @param paimonOptions Paimon catalog options without Hadoop-prefixed sensitive properties.
* @param hadoopConf Hadoop configuration carrying filesystem credentials and Hadoop options.
* @return the created inner Paimon catalog.
*/
@VisibleForTesting
protected AbstractCatalog createInnerCatalog(
Map<String, String> paimonOptions, Configuration hadoopConf) {
CatalogContext catalogContext =
CatalogContext.create(
Options.fromMap(paimonOptions), hadoopConf, new FlinkFileIOLoader(), null);
return (AbstractCatalog)
FlinkCatalogFactory.createCatalog(
context.getName(), catalogContext, context.getClassLoader());
}

private static void moveHadoopOptionsToConf(
Map<String, String> paimonOptions, Configuration hadoopConf) {
paimonOptions
.entrySet()
.removeIf(
entry -> {
String hadoopKey = toHadoopConfKey(entry.getKey());
if (hadoopKey == null) {
return false;
}

hadoopConf.set(hadoopKey, entry.getValue());
return true;
});
}

private static Map<String, String> withoutHadoopOptions(Map<String, String> options) {
Map<String, String> result = new HashMap<>(options);
result.keySet().removeIf(key -> toHadoopConfKey(key) != null);
return result;
}

private static String toHadoopConfKey(String key) {
if (key.startsWith(PropertyUtils.HADOOP_PREFIX)) {
return key.substring(PropertyUtils.HADOOP_PREFIX.length());
} else if (key.startsWith(PropertyUtils.FS_PREFIX)
|| key.startsWith(PropertyUtils.DFS_PREFIX)) {
return key;
}

return null;
}

// ---------------------------------------------------------------------------
// Lifecycle — keep paimonCatalog in sync with the outer catalog
// ---------------------------------------------------------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,12 @@
import org.apache.flink.table.catalog.exceptions.TableNotExistException;
import org.apache.flink.table.factories.CatalogFactory;
import org.apache.gravitino.Catalog;
import org.apache.gravitino.catalog.lakehouse.paimon.PaimonConstants;
import org.apache.gravitino.credential.Credential;
import org.apache.gravitino.credential.JdbcCredential;
import org.apache.gravitino.credential.OSSSecretKeyCredential;
import org.apache.gravitino.credential.S3SecretKeyCredential;
import org.apache.gravitino.credential.SupportsCredentials;
import org.apache.gravitino.flink.connector.DefaultPartitionConverter;
import org.apache.gravitino.flink.connector.catalog.BaseCatalog;
import org.apache.gravitino.rel.Table;
Expand Down Expand Up @@ -166,6 +172,36 @@ private static Map<String, String> defaultPaimonOptions() {
}
}

private static class CapturingPaimonCatalog extends GravitinoPaimonCatalog {

private final AbstractCatalog innerCatalog = mock(AbstractCatalog.class);
private final Catalog injectedCatalog;
private Map<String, String> capturedOptions;
private org.apache.hadoop.conf.Configuration capturedHadoopConf;

CapturingPaimonCatalog(Map<String, String> options, Catalog injectedCatalog) {
super(
new MockCatalogContext("test-paimon", options),
"default",
PaimonPropertiesConverter.INSTANCE,
DefaultPartitionConverter.INSTANCE);
this.injectedCatalog = injectedCatalog;
}

@Override
protected Catalog catalog() {
return injectedCatalog;
}

@Override
protected AbstractCatalog createInnerCatalog(
Map<String, String> paimonOptions, org.apache.hadoop.conf.Configuration hadoopConf) {
capturedOptions = new HashMap<>(paimonOptions);
capturedHadoopConf = new org.apache.hadoop.conf.Configuration(hadoopConf);
return innerCatalog;
}
}

@BeforeEach
void setUp() {
mockPaimonCatalog = mock(AbstractCatalog.class);
Expand Down Expand Up @@ -409,6 +445,87 @@ public void testDropTableInvalidatesNativeCacheAfterSuccessfulPurge() throws Exc
verify(mockInnerCatalog).invalidateTable(expected);
}

/** Verifies that Hadoop-prefixed catalog options are moved out of Paimon options. */
@Test
public void testOpenMovesFileSystemOptionsToHadoopConf() {
Catalog mockCatalog = catalogWithCredentials();
Map<String, String> options = new HashMap<>();
options.put("warehouse", "oss://bucket/path");
options.put("hadoop." + PaimonConstants.OSS_ACCESS_KEY, "catalog-access-key");
options.put("hadoop." + PaimonConstants.OSS_SECRET_KEY, "catalog-secret-key");
options.put("hadoop.fs.oss.endpoint", "oss-endpoint");
options.put("fs.bos.access.key", "bos-access-key");
options.put("fs.bos.secret.access.key", "bos-secret-key");

CapturingPaimonCatalog catalog = new CapturingPaimonCatalog(options, mockCatalog);
catalog.open();

Assertions.assertFalse(
catalog.capturedOptions.containsKey("hadoop." + PaimonConstants.OSS_ACCESS_KEY));
Assertions.assertFalse(
catalog.capturedOptions.containsKey("hadoop." + PaimonConstants.OSS_SECRET_KEY));
Assertions.assertFalse(catalog.capturedOptions.containsKey("hadoop.fs.oss.endpoint"));
Assertions.assertFalse(catalog.capturedOptions.containsKey("fs.bos.access.key"));
Assertions.assertFalse(catalog.capturedOptions.containsKey("fs.bos.secret.access.key"));
Assertions.assertEquals(
"catalog-access-key", catalog.capturedHadoopConf.get(PaimonConstants.OSS_ACCESS_KEY));
Assertions.assertEquals(
"catalog-secret-key", catalog.capturedHadoopConf.get(PaimonConstants.OSS_SECRET_KEY));
Assertions.assertEquals("oss-endpoint", catalog.capturedHadoopConf.get("fs.oss.endpoint"));
Assertions.assertEquals("bos-access-key", catalog.capturedHadoopConf.get("fs.bos.access.key"));
Assertions.assertEquals(
"bos-secret-key", catalog.capturedHadoopConf.get("fs.bos.secret.access.key"));
}

/** Verifies that vended filesystem credentials remain available to Paimon native FileIOs. */
@Test
public void testOpenKeepsStorageCredentialsInPaimonOptions() {
Catalog mockCatalog =
catalogWithCredentials(
new S3SecretKeyCredential("s3-key", "s3-secret"),
new OSSSecretKeyCredential("oss-key", "oss-secret"));
Map<String, String> options = new HashMap<>();
options.put("warehouse", "oss://bucket/path");
options.put(PaimonConstants.S3_ACCESS_KEY, "stale-s3-key");
options.put(PaimonConstants.S3_SECRET_KEY, "stale-s3-secret");
options.put(PaimonConstants.OSS_ACCESS_KEY, "stale-oss-key");
options.put(PaimonConstants.OSS_SECRET_KEY, "stale-oss-secret");

CapturingPaimonCatalog catalog = new CapturingPaimonCatalog(options, mockCatalog);
catalog.open();

Assertions.assertEquals("s3-key", catalog.capturedOptions.get(PaimonConstants.S3_ACCESS_KEY));
Assertions.assertEquals(
"s3-secret", catalog.capturedOptions.get(PaimonConstants.S3_SECRET_KEY));
Assertions.assertEquals("oss-key", catalog.capturedOptions.get(PaimonConstants.OSS_ACCESS_KEY));
Assertions.assertEquals(
"oss-secret", catalog.capturedOptions.get(PaimonConstants.OSS_SECRET_KEY));
}

/** Verifies that JDBC backend credentials remain Paimon catalog options. */
@Test
public void testOpenKeepsJdbcCredentialsInPaimonOptions() {
Catalog mockCatalog = catalogWithCredentials(new JdbcCredential("jdbc-user", "jdbc-password"));
Map<String, String> options = new HashMap<>();
options.put("warehouse", "file:/tmp/test-paimon-warehouse");

CapturingPaimonCatalog catalog = new CapturingPaimonCatalog(options, mockCatalog);
catalog.open();

Assertions.assertEquals(
"jdbc-user", catalog.capturedOptions.get(PaimonConstants.PAIMON_JDBC_USER));
Assertions.assertEquals(
"jdbc-password", catalog.capturedOptions.get(PaimonConstants.PAIMON_JDBC_PASSWORD));
}

private static Catalog catalogWithCredentials(Credential... credentials) {
Catalog catalog = mock(Catalog.class);
SupportsCredentials supportsCredentials = mock(SupportsCredentials.class);
when(catalog.supportsCredentials()).thenReturn(supportsCredentials);
when(supportsCredentials.getCredentials()).thenReturn(credentials);
return catalog;
}

// ---------------------------------------------------------------------------
// Helper: minimal CatalogFactory.Context implementation for constructor
// ---------------------------------------------------------------------------
Expand Down
Loading