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 7dadbf76d6 [MINOR] refactor(flink-connector): Add overridable 
toGravitinoType hook to BaseCatalog (#12244)
7dadbf76d6 is described below

commit 7dadbf76d6a65e9dbc40b02a1980ce132dc34b09
Author: Yuhui <[email protected]>
AuthorDate: Wed Jul 29 14:20:59 2026 +0800

    [MINOR] refactor(flink-connector): Add overridable toGravitinoType hook to 
BaseCatalog (#12244)
    
    ### What changes were proposed in this pull request?
    Added an overridable `toGravitinoType(LogicalType)` hook in
    `BaseCatalog` (Flink connector), replacing the hardcoded static
    `TypeUtils.toGravitinoType` calls with calls to this instance method.
    
    ### Why are the changes needed?
    The Flink connector's type conversion (`TypeUtils.toGravitinoType`) had
    no per-catalog override point, unlike the Spark and Trino connectors,
    which already special-case certain JDBC catalogs for type-mapping quirks
    (e.g. `SparkTypeConverter`). This adds the same extension point to the
    Flink connector so catalog-specific subclasses can override type mapping
    where needed.
    
    ### Does this PR introduce any user-facing change?
    No public API change. `BaseCatalog` subclasses can now override
    `toGravitinoType` to customize column type mapping.
    
    ### How was this patch tested?
    Existing `TestBaseCatalog` suite updated to construct catalogs via the
    `TestableBaseCatalog` test helper instead of calling static methods, and
    passes.
---
 .../flink/connector/catalog/BaseCatalog.java       | 42 ++++++++++++++--------
 .../flink/connector/catalog/TestBaseCatalog.java   | 12 ++++---
 2 files changed, 36 insertions(+), 18 deletions(-)

diff --git 
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
 
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
index 59874aff4a..b116e93c69 100644
--- 
a/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
+++ 
b/flink-connector/flink-common/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
@@ -63,6 +63,7 @@ import 
org.apache.flink.table.catalog.stats.CatalogColumnStatistics;
 import org.apache.flink.table.catalog.stats.CatalogTableStatistics;
 import org.apache.flink.table.expressions.Expression;
 import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.gravitino.Catalog;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
@@ -98,6 +99,7 @@ import org.apache.gravitino.rel.expressions.sorts.SortOrder;
 import org.apache.gravitino.rel.expressions.transforms.Transform;
 import org.apache.gravitino.rel.indexes.Index;
 import org.apache.gravitino.rel.indexes.Indexes;
+import org.apache.gravitino.rel.types.Type;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -127,6 +129,18 @@ public abstract class BaseCatalog extends AbstractCatalog {
 
   protected abstract AbstractCatalog realCatalog();
 
+  /**
+   * Converts a Flink logical type to a Gravitino type. Subclasses may 
override this to special-case
+   * types whose default mapping does not fit a particular catalog (e.g. 
Oracle has no pure date
+   * type, so its catalog maps Flink's {@code DATE} to a Gravitino timestamp 
instead).
+   *
+   * @param logicalType the Flink logical type
+   * @return the corresponding Gravitino type
+   */
+  protected Type toGravitinoType(LogicalType logicalType) {
+    return TypeUtils.toGravitinoType(logicalType);
+  }
+
   @Override
   public void open() throws CatalogException {
     realCatalog().open();
@@ -404,7 +418,7 @@ public abstract class BaseCatalog extends AbstractCatalog {
     ResolvedCatalogBaseTable<?> resolvedTable = (ResolvedCatalogBaseTable<?>) 
table;
     Column[] columns =
         resolvedTable.getResolvedSchema().getColumns().stream()
-            .map(BaseCatalog::toGravitinoColumn)
+            .map(this::toGravitinoColumn)
             .toArray(Column[]::new);
     String comment = table.getComment();
     Map<String, String> flinkOptions = table.getOptions();
@@ -852,10 +866,10 @@ public abstract class BaseCatalog extends AbstractCatalog 
{
     return Optional.of(primaryKeyFieldList);
   }
 
-  private static Column 
toGravitinoColumn(org.apache.flink.table.catalog.Column column) {
+  private Column toGravitinoColumn(org.apache.flink.table.catalog.Column 
column) {
     return Column.of(
         column.getName(),
-        TypeUtils.toGravitinoType(column.getDataType().getLogicalType()),
+        toGravitinoType(column.getDataType().getLogicalType()),
         column.getComment().orElse(null),
         column.getDataType().getLogicalType().isNullable(),
         false,
@@ -877,17 +891,17 @@ public abstract class BaseCatalog extends AbstractCatalog 
{
     changes.add(TableChange.deleteColumn(new String[] 
{change.getColumnName()}, true));
   }
 
-  private static void addColumn(
+  private void addColumn(
       org.apache.flink.table.catalog.TableChange.AddColumn change, 
List<TableChange> changes) {
     changes.add(
         TableChange.addColumn(
             new String[] {change.getColumn().getName()},
-            
TypeUtils.toGravitinoType(change.getColumn().getDataType().getLogicalType()),
+            toGravitinoType(change.getColumn().getDataType().getLogicalType()),
             change.getColumn().getComment().orElse(null),
             TableUtils.toGravitinoColumnPosition(change.getPosition())));
   }
 
-  private static void modifyColumn(
+  private void modifyColumn(
       org.apache.flink.table.catalog.TableChange change, List<TableChange> 
changes) {
     if (change instanceof 
org.apache.flink.table.catalog.TableChange.ModifyColumnName) {
       org.apache.flink.table.catalog.TableChange.ModifyColumnName 
modifyColumnName =
@@ -903,7 +917,7 @@ public abstract class BaseCatalog extends AbstractCatalog {
       changes.add(
           TableChange.updateColumnType(
               new String[] {modifyColumnType.getOldColumn().getName()},
-              
TypeUtils.toGravitinoType(modifyColumnType.getNewType().getLogicalType())));
+              
toGravitinoType(modifyColumnType.getNewType().getLogicalType())));
     } else if (change instanceof 
org.apache.flink.table.catalog.TableChange.ModifyColumnPosition) {
       org.apache.flink.table.catalog.TableChange.ModifyColumnPosition 
modifyColumnPosition =
           (org.apache.flink.table.catalog.TableChange.ModifyColumnPosition) 
change;
@@ -937,7 +951,7 @@ public abstract class BaseCatalog extends AbstractCatalog {
   }
 
   @VisibleForTesting
-  static TableChange[] getGravitinoTableChanges(
+  TableChange[] getGravitinoTableChanges(
       List<org.apache.flink.table.catalog.TableChange> tableChanges) {
     List<TableChange> changes = Lists.newArrayList();
     for (org.apache.flink.table.catalog.TableChange change : tableChanges) {
@@ -1046,14 +1060,14 @@ public abstract class BaseCatalog extends 
AbstractCatalog {
   }
 
   @VisibleForTesting
-  static ViewChange[] toReplaceViewChange(
+  ViewChange[] toReplaceViewChange(
       CatalogBaseTable existingView, ResolvedCatalogView newView, String 
dialect) {
     return toReplaceViewChange(
         existingView, newView, buildSqlRepresentation(dialect, 
newView.getExpandedQuery()));
   }
 
   @VisibleForTesting
-  static ViewChange[] toReplaceViewChange(
+  ViewChange[] toReplaceViewChange(
       CatalogBaseTable existingView,
       ResolvedCatalogView newView,
       Representation[] representations) {
@@ -1064,7 +1078,7 @@ public abstract class BaseCatalog extends AbstractCatalog 
{
   }
 
   @VisibleForTesting
-  static ViewChange[] toReplaceViewChange(
+  ViewChange[] toReplaceViewChange(
       List<org.apache.flink.table.catalog.TableChange> tableChanges,
       ResolvedCatalogView newView,
       String dialect) {
@@ -1073,7 +1087,7 @@ public abstract class BaseCatalog extends AbstractCatalog 
{
   }
 
   @VisibleForTesting
-  static ViewChange[] toReplaceViewChange(
+  ViewChange[] toReplaceViewChange(
       List<org.apache.flink.table.catalog.TableChange> tableChanges,
       ResolvedCatalogView newView,
       Representation[] representations) {
@@ -1106,9 +1120,9 @@ public abstract class BaseCatalog extends AbstractCatalog 
{
     return changes.toArray(new ViewChange[0]);
   }
 
-  private static Column[] toGravitinoColumns(ResolvedCatalogView view) {
+  private Column[] toGravitinoColumns(ResolvedCatalogView view) {
     return view.getResolvedSchema().getColumns().stream()
-        .map(BaseCatalog::toGravitinoColumn)
+        .map(this::toGravitinoColumn)
         .toArray(Column[]::new);
   }
 
diff --git 
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
index e24dc13c6c..5b3e44d90b 100644
--- 
a/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
+++ 
b/flink-connector/flink-common/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
@@ -116,8 +116,9 @@ public class TestBaseCatalog {
             org.apache.gravitino.rel.TableChange.setProperty("key", "value"),
             org.apache.gravitino.rel.TableChange.removeProperty("key"));
 
+    BaseCatalog catalog = new TestableBaseCatalog(null, null);
     org.apache.gravitino.rel.TableChange[] gravitinoTableChanges =
-        BaseCatalog.getGravitinoTableChanges(tableChanges);
+        catalog.getGravitinoTableChanges(tableChanges);
     Assertions.assertArrayEquals(expected.toArray(), gravitinoTableChanges);
   }
 
@@ -185,8 +186,9 @@ public class TestBaseCatalog {
 
     Schema schema = Schema.newBuilder().column("id", DataTypes.INT()).build();
 
+    BaseCatalog catalog = new TestableBaseCatalog(null, null);
     ViewChange[] changes =
-        BaseCatalog.toReplaceViewChange(
+        catalog.toReplaceViewChange(
             tableChanges, resolveView(schema, "SELECT 1", "comment"), 
Dialects.FLINK);
 
     Assertions.assertEquals(2, changes.length);
@@ -206,8 +208,9 @@ public class TestBaseCatalog {
 
     Schema schema = Schema.newBuilder().column("id", DataTypes.INT()).build();
 
+    BaseCatalog catalog = new TestableBaseCatalog(null, null);
     ViewChange[] changes =
-        BaseCatalog.toReplaceViewChange(
+        catalog.toReplaceViewChange(
             tableChanges, resolveView(schema, "SELECT id FROM t", "new 
comment"), Dialects.FLINK);
 
     // Should have exactly one SetProperty and one ReplaceView (order may vary)
@@ -241,8 +244,9 @@ public class TestBaseCatalog {
     CatalogView existing =
         CatalogView.of(schema, "old comment", "SELECT 1", "SELECT 1", 
Collections.emptyMap());
 
+    BaseCatalog catalog = new TestableBaseCatalog(null, null);
     ViewChange[] changes =
-        BaseCatalog.toReplaceViewChange(
+        catalog.toReplaceViewChange(
             existing, resolveView(schema, "SELECT 2", "new comment"), 
Dialects.FLINK);
 
     Assertions.assertEquals(1, changes.length);

Reply via email to