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

yuqi1129 pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/branch-1.3 by this push:
     new 518e2faed8 [Cherry-pick to branch-1.3] [MINOR] refactor(flink): Allow 
catalogs to customize Flink type conversion (#12961) (#12965)
518e2faed8 is described below

commit 518e2faed83da817fee880fb66b7019a9ba16c09
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Sep 8 09:36:44 2026 +0800

    [Cherry-pick to branch-1.3] [MINOR] refactor(flink): Allow catalogs to 
customize Flink type conversion (#12961) (#12965)
    
    **Cherry-pick Information:**
    - Original commit: 14f4995f9ccb1fd8c513db4a27eed694823d2558
    - Target branch: `branch-1.3`
    - Status: ⚠️ **Has conflicts - manual resolution required**
    
    **Do not merge** until conflict markers are resolved and the
    `cherry-pick-conflict` label is removed.
    
    Please review and resolve the conflicts before merging.
    
    Co-authored-by: Qi Yu <[email protected]>
---
 .../flink/connector/catalog/BaseCatalog.java       | 16 +++++++--
 .../flink/connector/catalog/TestBaseCatalog.java   | 42 ++++++++++++++++++++++
 2 files changed, 55 insertions(+), 3 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 ae806dcdbd..dbc7087649 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
@@ -98,6 +98,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 +128,16 @@ public abstract class BaseCatalog extends AbstractCatalog {
 
   protected abstract AbstractCatalog realCatalog();
 
+  /**
+   * Converts a Gravitino type to a Flink type, allowing catalog-specific 
native type mappings.
+   *
+   * @param type the Gravitino type
+   * @return the corresponding Flink data type
+   */
+  protected DataType toFlinkType(Type type) {
+    return TypeUtils.toFlinkType(type);
+  }
+
   @Override
   public void open() throws CatalogException {
     realCatalog().open();
@@ -1142,12 +1153,11 @@ public abstract class BaseCatalog extends 
AbstractCatalog {
    * @param columns the Gravitino column definitions
    * @return a Flink schema builder populated with the given columns
    */
-  protected static org.apache.flink.table.api.Schema.Builder 
buildSchemaFromColumns(
-      Column[] columns) {
+  protected org.apache.flink.table.api.Schema.Builder 
buildSchemaFromColumns(Column[] columns) {
     org.apache.flink.table.api.Schema.Builder builder =
         org.apache.flink.table.api.Schema.newBuilder();
     for (Column column : columns) {
-      DataType flinkType = TypeUtils.toFlinkType(column.dataType());
+      DataType flinkType = toFlinkType(column.dataType());
       builder
           .column(column.name(), column.nullable() ? flinkType.nullable() : 
flinkType.notNull())
           .withComment(column.comment());
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 3ba0edff0a..9db49f8933 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
@@ -38,6 +38,7 @@ import org.apache.flink.table.catalog.ResolvedCatalogView;
 import org.apache.flink.table.catalog.ResolvedSchema;
 import org.apache.flink.table.catalog.TableChange;
 import org.apache.flink.table.catalog.exceptions.TableNotExistException;
+import org.apache.flink.table.types.DataType;
 import org.apache.gravitino.Catalog;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
@@ -53,6 +54,7 @@ import org.apache.gravitino.rel.TableCatalog;
 import org.apache.gravitino.rel.ViewCatalog;
 import org.apache.gravitino.rel.ViewChange;
 import org.apache.gravitino.rel.expressions.distributions.Distributions;
+import org.apache.gravitino.rel.types.Type;
 import org.apache.gravitino.rel.types.Types;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -60,6 +62,46 @@ import org.mockito.Mockito;
 
 public class TestBaseCatalog {
 
+  @Test
+  void testDefaultFlinkTypeConversion() {
+    BaseCatalog catalog = new TestableBaseCatalog(null, null);
+    Assertions.assertEquals(DataTypes.INT(), 
catalog.toFlinkType(Types.IntegerType.get()));
+    Assertions.assertEquals(DataTypes.STRING(), 
catalog.toFlinkType(Types.StringType.get()));
+    Assertions.assertEquals(
+        DataTypes.DECIMAL(10, 2), catalog.toFlinkType(Types.DecimalType.of(10, 
2)));
+  }
+
+  @Test
+  void testSchemaUsesCatalogTypeConversion() {
+    Type nativeType = Types.ExternalType.of("native_text");
+    BaseCatalog catalog =
+        new TestableBaseCatalog(null, null) {
+          /** {@inheritDoc} */
+          @Override
+          protected DataType toFlinkType(Type type) {
+            return type.equals(nativeType) ? DataTypes.STRING() : 
super.toFlinkType(type);
+          }
+        };
+    org.apache.gravitino.rel.Column[] columns = {
+      org.apache.gravitino.rel.Column.of("text", nativeType, "source comment"),
+      org.apache.gravitino.rel.Column.of("required_text", nativeType, null, 
false, false, null),
+      org.apache.gravitino.rel.Column.of("id", Types.IntegerType.get(), null)
+    };
+    Schema schema = catalog.buildSchemaFromColumns(columns).build();
+    Schema.UnresolvedPhysicalColumn text =
+        (Schema.UnresolvedPhysicalColumn) schema.getColumns().get(0);
+    Schema.UnresolvedPhysicalColumn requiredText =
+        (Schema.UnresolvedPhysicalColumn) schema.getColumns().get(1);
+    Schema.UnresolvedPhysicalColumn id =
+        (Schema.UnresolvedPhysicalColumn) schema.getColumns().get(2);
+    Assertions.assertEquals("text", text.getName());
+    Assertions.assertEquals(DataTypes.STRING(), text.getDataType());
+    Assertions.assertEquals("source comment", text.getComment().orElseThrow());
+    Assertions.assertEquals(DataTypes.STRING().notNull(), 
requiredText.getDataType());
+    Assertions.assertEquals(DataTypes.INT(), id.getDataType());
+    Assertions.assertEquals(nativeType, columns[0].dataType());
+  }
+
   @Test
   public void testHiveSchemaChanges() {
     Map<String, String> currentProperties = ImmutableMap.of("key", "value", 
"key2", "value2");

Reply via email to