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..659da30801b 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; @@ -40,11 +39,16 @@ import org.apache.flink.table.factories.Factory; import org.apache.gravitino.NameIdentifier; import org.apache.gravitino.catalog.lakehouse.paimon.PaimonConstants; +import org.apache.gravitino.credential.Credential; import org.apache.gravitino.credential.CredentialPropertyUtils; +import org.apache.gravitino.credential.JdbcCredential; +import org.apache.gravitino.credential.OSSSecretKeyCredential; +import org.apache.gravitino.credential.S3SecretKeyCredential; import org.apache.gravitino.exceptions.NoSuchCatalogException; 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 +57,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; @@ -74,10 +83,15 @@ public class GravitinoPaimonCatalog extends BaseCatalog { private static final Logger LOG = LoggerFactory.getLogger(GravitinoPaimonCatalog.class); + private static final String PAIMON_S3_ACCESS_KEY_ALIAS = "s3.access.key"; + private static final String PAIMON_S3_SECRET_KEY_ALIAS = "s3.secret.key"; + private static final String HADOOP_S3_ENDPOINT = "fs.s3a.endpoint"; + private static final String HADOOP_S3_ACCESS_KEY = "fs.s3a.access.key"; + private static final String HADOOP_S3_SECRET_KEY = "fs.s3a.secret.key"; 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 +143,14 @@ public void open() throws CatalogException { super.open(); return; } + Map paimonOptions = new HashMap<>(mutableOptions); + Configuration hadoopConf = + HadoopUtils.getHadoopConfiguration(Options.fromMap(withoutHadoopOptions(paimonOptions))); + moveHadoopOptionsToConf(paimonOptions, hadoopConf); + movePaimonStorageOptionsToConf(paimonOptions, hadoopConf); try { - CredentialPropertyUtils.applyPaimonCredentials( - CredentialPropertyUtils.getCredentials(catalog()), mutableOptions); + applyCredentialsToCatalog( + CredentialPropertyUtils.getCredentials(catalog()), paimonOptions, hadoopConf); } catch (NoSuchCatalogException e) { LOG.warn( "Catalog '{}' not found in Gravitino during open(); credential injection skipped." @@ -139,33 +158,109 @@ 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; + } + + private static void movePaimonStorageOptionsToConf( + Map paimonOptions, Configuration hadoopConf) { + moveOptionToHadoopConf( + paimonOptions, hadoopConf, PaimonConstants.S3_ENDPOINT, HADOOP_S3_ENDPOINT); + moveOptionToHadoopConf( + paimonOptions, hadoopConf, PAIMON_S3_ACCESS_KEY_ALIAS, HADOOP_S3_ACCESS_KEY); + moveOptionToHadoopConf( + paimonOptions, hadoopConf, PaimonConstants.S3_ACCESS_KEY, HADOOP_S3_ACCESS_KEY); + moveOptionToHadoopConf( + paimonOptions, hadoopConf, PAIMON_S3_SECRET_KEY_ALIAS, HADOOP_S3_SECRET_KEY); + moveOptionToHadoopConf( + paimonOptions, hadoopConf, PaimonConstants.S3_SECRET_KEY, HADOOP_S3_SECRET_KEY); + } + + private static void moveOptionToHadoopConf( + Map paimonOptions, + Configuration hadoopConf, + String paimonKey, + String hadoopKey) { + String value = paimonOptions.remove(paimonKey); + if (value != null) { + hadoopConf.set(hadoopKey, value); + } + } + + private static void applyCredentialsToCatalog( + Credential[] credentials, Map paimonOptions, Configuration hadoopConf) { + for (Credential credential : credentials) { + if (credential instanceof JdbcCredential) { + JdbcCredential jdbc = (JdbcCredential) credential; + paimonOptions.put(PaimonConstants.PAIMON_JDBC_USER, jdbc.jdbcUser()); + paimonOptions.put(PaimonConstants.PAIMON_JDBC_PASSWORD, jdbc.jdbcPassword()); + } else if (credential instanceof S3SecretKeyCredential) { + S3SecretKeyCredential s3 = (S3SecretKeyCredential) credential; + hadoopConf.set(HADOOP_S3_ACCESS_KEY, s3.accessKeyId()); + hadoopConf.set(HADOOP_S3_SECRET_KEY, s3.secretAccessKey()); + } else if (credential instanceof OSSSecretKeyCredential) { + OSSSecretKeyCredential oss = (OSSSecretKeyCredential) credential; + hadoopConf.set(PaimonConstants.OSS_ACCESS_KEY, oss.accessKeyId()); + hadoopConf.set(PaimonConstants.OSS_SECRET_KEY, oss.secretAccessKey()); + } else { + LOG.warn( + "Received unrecognized credential type '{}' for Paimon catalog, skipping", + credential.getClass().getName()); + } + } + } + // --------------------------------------------------------------------------- // 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..89b7bc3fec6 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,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; @@ -166,6 +172,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 +445,116 @@ 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 are moved to Hadoop configuration. */ + @Test + public void testOpenMovesStorageCredentialsToHadoopConf() { + 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.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.S3_ACCESS_KEY)); + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.S3_SECRET_KEY)); + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.OSS_ACCESS_KEY)); + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.OSS_SECRET_KEY)); + Assertions.assertEquals("s3-key", catalog.capturedHadoopConf.get("fs.s3a.access.key")); + Assertions.assertEquals("s3-secret", catalog.capturedHadoopConf.get("fs.s3a.secret.key")); + Assertions.assertEquals( + "oss-key", catalog.capturedHadoopConf.get(PaimonConstants.OSS_ACCESS_KEY)); + Assertions.assertEquals( + "oss-secret", catalog.capturedHadoopConf.get(PaimonConstants.OSS_SECRET_KEY)); + } + + /** Verifies that Paimon native storage options are moved to Hadoop configuration. */ + @Test + public void testOpenMovesPaimonStorageOptionsToHadoopConf() { + Catalog mockCatalog = catalogWithCredentials(); + Map options = new HashMap<>(); + options.put("warehouse", "s3://bucket/path"); + options.put(PaimonConstants.S3_ENDPOINT, "s3-endpoint"); + options.put(PaimonConstants.S3_ACCESS_KEY, "s3-key"); + options.put(PaimonConstants.S3_SECRET_KEY, "s3-secret"); + options.put("s3.access.key", "s3-key-alias"); + options.put("s3.secret.key", "s3-secret-alias"); + + CapturingPaimonCatalog catalog = new CapturingPaimonCatalog(options, mockCatalog); + catalog.open(); + + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.S3_ENDPOINT)); + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.S3_ACCESS_KEY)); + Assertions.assertFalse(catalog.capturedOptions.containsKey(PaimonConstants.S3_SECRET_KEY)); + Assertions.assertFalse(catalog.capturedOptions.containsKey("s3.access.key")); + Assertions.assertFalse(catalog.capturedOptions.containsKey("s3.secret.key")); + Assertions.assertEquals("s3-endpoint", catalog.capturedHadoopConf.get("fs.s3a.endpoint")); + Assertions.assertEquals("s3-key", catalog.capturedHadoopConf.get("fs.s3a.access.key")); + Assertions.assertEquals("s3-secret", catalog.capturedHadoopConf.get("fs.s3a.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 // ---------------------------------------------------------------------------