From 32bcd9fe36f5b556a77ecbf10bc882565d067903 Mon Sep 17 00:00:00 2001 From: jarredhj0214 Date: Tue, 21 Jul 2026 16:43:30 +0800 Subject: [PATCH 1/3] [#12051] fix(flink): move Paimon Hadoop options to Hadoop conf --- .../paimon/GravitinoPaimonCatalog.java | 93 +++++++++----- .../paimon/TestGravitinoPaimonCatalog.java | 116 ++++++++++++++++++ 2 files changed, 181 insertions(+), 28 deletions(-) diff --git a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java index 5290490f046..77c7e9d46fc 100644 --- a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java +++ b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java @@ -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; @@ -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; @@ -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; @@ -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 mutableOptions; private AbstractCatalog paimonCatalog; @@ -129,9 +134,13 @@ public void open() throws CatalogException { super.open(); return; } + Map paimonOptions = new HashMap<>(mutableOptions); + Configuration hadoopConf = + HadoopUtils.getHadoopConfiguration(Options.fromMap(withoutHadoopOptions(paimonOptions))); + moveHadoopOptionsToConf(paimonOptions, hadoopConf); try { CredentialPropertyUtils.applyPaimonCredentials( - CredentialPropertyUtils.getCredentials(catalog()), mutableOptions); + CredentialPropertyUtils.getCredentials(catalog()), paimonOptions); } catch (NoSuchCatalogException e) { LOG.warn( "Catalog '{}' not found in Gravitino during open(); credential injection skipped." @@ -139,33 +148,61 @@ public void open() throws CatalogException { getName(), e); } - CatalogFactory.Context contextWithCredentials = - new CatalogFactory.Context() { - @Override - public String getName() { - return context.getName(); - } - - @Override - public Map 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 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 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 withoutHadoopOptions(Map options) { + Map 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 // --------------------------------------------------------------------------- diff --git a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java index 3e29e30e2e6..cbaec4f21ec 100644 --- a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java +++ b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java @@ -38,6 +38,11 @@ import org.apache.flink.table.catalog.exceptions.TableNotExistException; import org.apache.flink.table.factories.CatalogFactory; import org.apache.gravitino.Catalog; +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; @@ -166,6 +171,36 @@ private static Map defaultPaimonOptions() { } } + private static class CapturingPaimonCatalog extends GravitinoPaimonCatalog { + + private final AbstractCatalog innerCatalog = mock(AbstractCatalog.class); + private final Catalog injectedCatalog; + private Map capturedOptions; + private org.apache.hadoop.conf.Configuration capturedHadoopConf; + + CapturingPaimonCatalog(Map 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 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); @@ -409,6 +444,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 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 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 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 // --------------------------------------------------------------------------- From 65c0db5c85c4cf76c6a1a6d5e763b6db57e0cb27 Mon Sep 17 00:00:00 2001 From: jarredhj0214 Date: Tue, 21 Jul 2026 17:53:19 +0800 Subject: [PATCH 2/3] [#12051] test(flink): import Paimon constants --- .../flink/connector/paimon/TestGravitinoPaimonCatalog.java | 1 + 1 file changed, 1 insertion(+) diff --git a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java index cbaec4f21ec..e3a4bca29f1 100644 --- a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java +++ b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/paimon/TestGravitinoPaimonCatalog.java @@ -38,6 +38,7 @@ 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; From 9e36de3c99b686ab43bf593c0be3f1dc964f09e5 Mon Sep 17 00:00:00 2001 From: jarredhj0214 Date: Wed, 22 Jul 2026 17:26:48 +0800 Subject: [PATCH 3/3] [#12051] docs(flink): clarify Paimon credential options --- .../flink/connector/paimon/GravitinoPaimonCatalog.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java index 77c7e9d46fc..08fe23798f2 100644 --- a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java +++ b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java @@ -139,6 +139,8 @@ public void open() throws CatalogException { 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()), paimonOptions); } catch (NoSuchCatalogException e) {