This is an automated email from the ASF dual-hosted git repository.

diqiu50 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new b368aaf2b0 [#12051] fix(flink-connector): Move Paimon Hadoop options 
to Hadoop conf (#12121)
b368aaf2b0 is described below

commit b368aaf2b037200b5e3a65560cbbd2c808d88a98
Author: jarred0214 <[email protected]>
AuthorDate: Tue Jul 28 09:51:55 2026 +0800

    [#12051] fix(flink-connector): Move Paimon Hadoop options to Hadoop conf 
(#12121)
    
    ### What changes were proposed in this pull request?
    
    Move Paimon Flink connector Hadoop filesystem options (`hadoop.*`,
    `fs.*`, and `dfs.*`) from Paimon catalog options into Hadoop
    `Configuration` before creating the inner Paimon catalog.
    
    This creates the inner Paimon catalog with a `CatalogContext` that
    carries the sanitized Paimon options plus Hadoop configuration, and
    keeps vended Paimon S3/OSS/JDBC credentials available as Paimon options.
    
    ### Why are the changes needed?
    
    Paimon can log dynamic/catalog options and expose sensitive filesystem
    credentials such as OSS/BOS access keys when they remain in Paimon
    options. Moving Hadoop filesystem credentials to Hadoop configuration
    follows the Flink Hive connector pattern and avoids passing these
    sensitive Hadoop properties through Paimon options.
    
    Fix: #12051
    
    ### Does this PR introduce _any_ user-facing change?
    
    No API changes. Existing Hadoop filesystem catalog properties continue
    to be accepted, but the Flink Paimon connector forwards them via Hadoop
    `Configuration` instead of Paimon catalog options.
    
    ### How was this patch tested?
    
    Added unit coverage for:
    - moving `hadoop.*`/`fs.*` filesystem options into Hadoop configuration
    - keeping vended S3/OSS credentials in Paimon options for native Paimon
    FileIOs
    - keeping JDBC credentials in Paimon options
    
    Local checks:
    - `git diff --check`
    
    Could not run Gradle tests locally because only JDK 8 is installed and
    the build requires JDK 17.
---
 .../connector/paimon/GravitinoPaimonCatalog.java   |  95 ++++++++++----
 .../flink/connector/utils/PropertyUtils.java       |  50 +++++++
 .../paimon/TestGravitinoPaimonCatalog.java         | 146 +++++++++++++++++++++
 .../flink/connector/utils/TestPropertyUtils.java   |  49 +++++++
 4 files changed, 312 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 5290490f04..b4f190744c 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.Map;
 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.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 +53,14 @@ import org.apache.gravitino.rel.expressions.NamedReference;
 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.fs.FileIOLoader;
+import org.apache.paimon.options.Options;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -74,10 +79,15 @@ import org.slf4j.LoggerFactory;
 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<String, String> mutableOptions;
   private AbstractCatalog paimonCatalog;
 
@@ -129,9 +139,11 @@ public class GravitinoPaimonCatalog extends BaseCatalog {
       super.open();
       return;
     }
+    Map<String, String> paimonOptions = new HashMap<>(mutableOptions);
+    Configuration hadoopConf = 
PropertyUtils.extractHadoopConfiguration(paimonOptions);
     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 +151,60 @@ public class GravitinoPaimonCatalog extends BaseCatalog {
           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);
+    movePaimonStorageOptionsToConf(paimonOptions, hadoopConf);
+    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) {
+    FileIOLoader preferIOLoader = null;
+    FileIOLoader fallbackIOLoader = new FlinkFileIOLoader();
+    CatalogContext catalogContext =
+        CatalogContext.create(
+            Options.fromMap(paimonOptions), hadoopConf, preferIOLoader, 
fallbackIOLoader);
+    return (AbstractCatalog)
+        FlinkCatalogFactory.createCatalog(
+            context.getName(), catalogContext, context.getClassLoader());
+  }
+
+  private static void movePaimonStorageOptionsToConf(
+      Map<String, String> 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);
+    moveOptionToHadoopConf(
+        paimonOptions, hadoopConf, PaimonConstants.OSS_ACCESS_KEY, 
PaimonConstants.OSS_ACCESS_KEY);
+    moveOptionToHadoopConf(
+        paimonOptions, hadoopConf, PaimonConstants.OSS_SECRET_KEY, 
PaimonConstants.OSS_SECRET_KEY);
+  }
+
+  private static void moveOptionToHadoopConf(
+      Map<String, String> paimonOptions,
+      Configuration hadoopConf,
+      String paimonKey,
+      String hadoopKey) {
+    String value = paimonOptions.remove(paimonKey);
+    if (value != null) {
+      hadoopConf.set(hadoopKey, value);
+    }
+  }
+
   // 
---------------------------------------------------------------------------
   // Lifecycle — keep paimonCatalog in sync with the outer catalog
   // 
---------------------------------------------------------------------------
diff --git 
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/utils/PropertyUtils.java
 
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/utils/PropertyUtils.java
index 9fc24a6bf5..f2dd246414 100644
--- 
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/utils/PropertyUtils.java
+++ 
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/utils/PropertyUtils.java
@@ -19,9 +19,14 @@
 package org.apache.gravitino.flink.connector.utils;
 
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.Map;
 import java.util.stream.Collectors;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.utils.HadoopUtils;
 
+/** Utility methods for Flink connector properties. */
 public class PropertyUtils {
 
   public static final String HIVE_PREFIX = "hive.";
@@ -29,6 +34,12 @@ public class PropertyUtils {
   public static final String FS_PREFIX = "fs.";
   public static final String DFS_PREFIX = "dfs.";
 
+  /**
+   * Gets Hadoop and Hive properties.
+   *
+   * @param properties the source properties
+   * @return Hadoop and Hive properties
+   */
   public static Map<String, String> getHadoopAndHiveProperties(Map<String, 
String> properties) {
     if (properties == null) {
       return Collections.emptyMap();
@@ -43,4 +54,43 @@ public class PropertyUtils {
                     || entry.getKey().startsWith(HIVE_PREFIX))
         .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
   }
+
+  /**
+   * Extracts Hadoop configuration from the given options and removes 
Hadoop-prefixed options.
+   *
+   * <p>Options with the {@code hadoop.} prefix are written to Hadoop 
configuration without the
+   * prefix. Options with the {@code fs.} or {@code dfs.} prefix are written 
as-is.
+   *
+   * @param options mutable options to extract Hadoop configuration from
+   * @return Hadoop configuration containing extracted Hadoop options
+   */
+  public static Configuration extractHadoopConfiguration(Map<String, String> 
options) {
+    Map<String, String> hadoopProps = new HashMap<>();
+    options
+        .entrySet()
+        .removeIf(
+            entry -> {
+              String hadoopKey = toHadoopConfKey(entry.getKey());
+              if (hadoopKey == null) {
+                return false;
+              }
+
+              hadoopProps.put(hadoopKey, entry.getValue());
+              return true;
+            });
+
+    Configuration conf = 
HadoopUtils.getHadoopConfiguration(Options.fromMap(options));
+    hadoopProps.forEach(conf::set);
+    return conf;
+  }
+
+  private static String toHadoopConfKey(String key) {
+    if (key.startsWith(HADOOP_PREFIX)) {
+      return key.substring(HADOOP_PREFIX.length());
+    } else if (key.startsWith(FS_PREFIX) || key.startsWith(DFS_PREFIX)) {
+      return key;
+    }
+
+    return null;
+  }
 }
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 3e29e30e2e..89b7bc3fec 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.CatalogException;
 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 @@ public class TestGravitinoPaimonCatalog {
     }
   }
 
+  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);
@@ -409,6 +445,116 @@ public class TestGravitinoPaimonCatalog {
     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 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<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.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<String, String> 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<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
   // 
---------------------------------------------------------------------------
diff --git 
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/utils/TestPropertyUtils.java
 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/utils/TestPropertyUtils.java
new file mode 100644
index 0000000000..8c7711950e
--- /dev/null
+++ 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/utils/TestPropertyUtils.java
@@ -0,0 +1,49 @@
+/*
+ * 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.gravitino.flink.connector.utils;
+
+import java.util.HashMap;
+import java.util.Map;
+import org.apache.hadoop.conf.Configuration;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+/** Unit tests for {@link PropertyUtils}. */
+public class TestPropertyUtils {
+
+  @Test
+  public void testExtractHadoopConfigurationRemovesHadoopOptions() {
+    Map<String, String> options = new HashMap<>();
+    options.put("warehouse", "file:/tmp/warehouse");
+    options.put("hadoop.fs.oss.endpoint", "oss-endpoint");
+    options.put("fs.s3a.access.key", "s3-key");
+    options.put("dfs.client.use.datanode.hostname", "true");
+
+    Configuration configuration = 
PropertyUtils.extractHadoopConfiguration(options);
+
+    Assertions.assertEquals("file:/tmp/warehouse", options.get("warehouse"));
+    Assertions.assertFalse(options.containsKey("hadoop.fs.oss.endpoint"));
+    Assertions.assertFalse(options.containsKey("fs.s3a.access.key"));
+    
Assertions.assertFalse(options.containsKey("dfs.client.use.datanode.hostname"));
+    Assertions.assertEquals("oss-endpoint", 
configuration.get("fs.oss.endpoint"));
+    Assertions.assertEquals("s3-key", configuration.get("fs.s3a.access.key"));
+    Assertions.assertEquals("true", 
configuration.get("dfs.client.use.datanode.hostname"));
+  }
+}

Reply via email to