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");