This is an automated email from the ASF dual-hosted git repository.
yuqi1129 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 7b7ec06582 [MINOR] refactor(flink-connector): Align JDBC connector
common logic (#11849)
7b7ec06582 is described below
commit 7b7ec06582701ed1783b2d7aac1013e41b8f6ba7
Author: Yuhui <[email protected]>
AuthorDate: Thu Jul 2 09:47:03 2026 +0800
[MINOR] refactor(flink-connector): Align JDBC connector common logic
(#11849)
### What changes were proposed in this pull request?
Align common Flink JDBC connector logic and related tests.
This PR:
- Exposes shared JDBC catalog mutable options and credential injection
to subclasses.
- Adds the common `jdbc-database` property constant for Flink JDBC
catalog conversion.
- Splits Flink common IT schema lifecycle checks from get-schema
capability checks.
- Adds coverage for existing Gravitino-managed catalog type routing.
### Why are the changes needed?
These changes keep common Flink JDBC connector behavior reusable without
adding database-specific logic to the OSS implementation.
They also allow catalogs that do not support schema lifecycle operations
to reuse the shared Flink integration test base.
Fix: #N/A
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Targeted Flink common tests.
---
.../flink/connector/jdbc/GravitinoJdbcCatalog.java | 12 +++++++++++-
.../connector/jdbc/JdbcPropertiesConstants.java | 1 +
.../connector/integration/test/FlinkCommonIT.java | 21 ++++++++++++++++++++-
.../store/TestGravitinoSessionCatalogStore.java | 10 ++++++++++
4 files changed, 42 insertions(+), 2 deletions(-)
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/GravitinoJdbcCatalog.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/GravitinoJdbcCatalog.java
index a50af98483..53b93dfb98 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/GravitinoJdbcCatalog.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/GravitinoJdbcCatalog.java
@@ -150,6 +150,16 @@ public class GravitinoJdbcCatalog extends BaseCatalog {
return Optional.of(new JdbcDynamicTableFactory());
}
+ /**
+ * Returns the mutable options map shared with {@link BaseCatalog}, allowing
subclasses to inject
+ * credentials into catalog options after construction.
+ *
+ * @return the mutable catalog options map
+ */
+ protected Map<String, String> getMutableOptions() {
+ return mutableOptions;
+ }
+
/**
* Overwrites the Flink JDBC user and password in {@code options} with
credentials obtained from
* the server via credential vending, if available. Falls back to the
existing options if the
@@ -159,7 +169,7 @@ public class GravitinoJdbcCatalog extends BaseCatalog {
* @param options the mutable Flink catalog options map to update
*/
@VisibleForTesting
- static void applyJdbcCredential(Catalog catalog, Map<String, String>
options) {
+ protected static void applyJdbcCredential(Catalog catalog, Map<String,
String> options) {
Credential[] credentials;
try {
credentials = catalog.supportsCredentials().getCredentials();
diff --git
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
index 4bea348e20..3a6de0ebd9 100644
---
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
+++
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/jdbc/JdbcPropertiesConstants.java
@@ -30,6 +30,7 @@ public class JdbcPropertiesConstants {
public static final String GRAVITINO_JDBC_PASSWORD = "jdbc-password";
public static final String GRAVITINO_JDBC_URL = "jdbc-url";
public static final String GRAVITINO_JDBC_DRIVER = "jdbc-driver";
+ public static final String GRAVITINO_JDBC_DATABASE = "jdbc-database";
public static final String FLINK_JDBC_URL = "base-url";
public static final String FLINK_JDBC_USER = "username";
diff --git
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/FlinkCommonIT.java
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/FlinkCommonIT.java
index 159fe00775..8cf59c7e60 100644
---
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/FlinkCommonIT.java
+++
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/integration/test/FlinkCommonIT.java
@@ -87,6 +87,23 @@ public abstract class FlinkCommonIT extends FlinkEnvIT {
protected abstract boolean supportDropCascade();
+ /**
+ * Returns {@code true} if the catalog supports schema CREATE/DROP via the
Gravitino API. Controls
+ * only the pure schema-lifecycle tests ({@code testCreateSchema}).
+ */
+ protected boolean supportsSchemaLifecycle() {
+ return true;
+ }
+
+ /**
+ * Returns {@code true} if get-schema tests can run. The base-class test
creates the schema via
+ * the Gravitino API, so this still depends on {@link
#supportsSchemaLifecycle()}. Subclasses that
+ * override {@code testGetSchemaWithoutCommentAndOption} directly should
ignore this flag.
+ */
+ protected boolean supportsSchemaLifecycleAndGetSchema() {
+ return supportsSchemaLifecycle() &&
supportGetSchemaWithoutCommentAndOption();
+ }
+
protected boolean supportsPrimaryKey() {
return true;
}
@@ -117,6 +134,7 @@ public abstract class FlinkCommonIT extends FlinkEnvIT {
}
@Test
+ @EnabledIf("supportsSchemaLifecycle")
public void testCreateSchema() {
doWithCatalog(
currentCatalog(),
@@ -134,7 +152,7 @@ public abstract class FlinkCommonIT extends FlinkEnvIT {
}
@Test
- @EnabledIf("supportGetSchemaWithoutCommentAndOption")
+ @EnabledIf("supportsSchemaLifecycleAndGetSchema")
public void testGetSchemaWithoutCommentAndOption() {
doWithCatalog(
currentCatalog(),
@@ -190,6 +208,7 @@ public abstract class FlinkCommonIT extends FlinkEnvIT {
}
@Test
+ @EnabledIf("supportsSchemaLifecycle")
public void testListSchema() {
doWithCatalog(
currentCatalog(),
diff --git
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/store/TestGravitinoSessionCatalogStore.java
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/store/TestGravitinoSessionCatalogStore.java
index 1b6011757a..d564fc7415 100644
---
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/store/TestGravitinoSessionCatalogStore.java
+++
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/store/TestGravitinoSessionCatalogStore.java
@@ -69,6 +69,16 @@ public class TestGravitinoSessionCatalogStore {
new GravitinoSessionCatalogStore(gravitinoCatalogStore,
memoryCatalogStore);
}
+ @Test
+ void testKnownGravitinoCatalogTypesAreGravitinoManaged() {
+ Assertions.assertTrue(isGravitinoManagedCatalogType("gravitino-hive"));
+
Assertions.assertTrue(isGravitinoManagedCatalogType("gravitino-jdbc-mysql"));
+
Assertions.assertTrue(isGravitinoManagedCatalogType("gravitino-jdbc-postgresql"));
+
Assertions.assertFalse(isGravitinoManagedCatalogType("gravitino-jdbc-custom"));
+ Assertions.assertFalse(isGravitinoManagedCatalogType("jdbc-mysql"));
+ Assertions.assertFalse(isGravitinoManagedCatalogType(null));
+ }
+
// -------------------------------------------------------------------------
// storeCatalog
// -------------------------------------------------------------------------