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

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new eb2eb1ddc8 [Fix][Connector-V2] Support PostgreSQL enum columns in the 
Postgres catalog and CDC (#12588)
eb2eb1ddc8 is described below

commit eb2eb1ddc8da4987f8fb0f06c28797b9f78656fe
Author: Goutam Adwant <[email protected]>
AuthorDate: Sun Oct 4 23:00:16 2026 +0000

    [Fix][Connector-V2] Support PostgreSQL enum columns in the Postgres catalog 
and CDC (#12588)
---
 .../PostgresRelationSchemaChangeResolverTest.java  |  65 +++++++++++
 .../jdbc/catalog/psql/PostgresCatalog.java         |  28 ++++-
 .../dialect/kingbase/KingbaseTypeConverter.java    |   8 ++
 .../dialect/psql/PostgresTypeConverter.java        |  30 ++++-
 .../dialect/redshift/RedshiftTypeConverter.java    |   8 ++
 .../psql/PostgresCatalogBuildColumnTest.java       | 126 +++++++++++++++++++++
 .../kingbase/KingbaseTypeConverterTest.java        |  16 +++
 .../dialect/psql/PostgresTypeConverterTest.java    |  68 +++++++++++
 .../redshift/RedshiftTypeConverterTest.java        |  16 +++
 9 files changed, 360 insertions(+), 5 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresRelationSchemaChangeResolverTest.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresRelationSchemaChangeResolverTest.java
index 35ced0289b..ed1d80b24c 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresRelationSchemaChangeResolverTest.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/PostgresRelationSchemaChangeResolverTest.java
@@ -82,6 +82,59 @@ class PostgresRelationSchemaChangeResolverTest {
         Assertions.assertEquals("email", enabledEvent.getAfterColumn());
     }
 
+    @Test
+    void shouldResolveAddedEnumColumnAsString() {
+        SourceRecord record =
+                createRecord(intColumn("id", 1), varcharColumn("name", 2, 64), 
enumColumn("m", 3));
+
+        SchemaChangeEvent event =
+                resolver.resolve(record, 
Collections.singletonList(createCatalogTable()));
+
+        AlterTableColumnsEvent columnsEvent = (AlterTableColumnsEvent) event;
+        Assertions.assertEquals(1, columnsEvent.getEvents().size());
+        AlterTableAddColumnEvent enumEvent =
+                (AlterTableAddColumnEvent) columnsEvent.getEvents().get(0);
+        Assertions.assertEquals("m", enumEvent.getColumn().getName());
+        Assertions.assertEquals(BasicType.STRING_TYPE, 
enumEvent.getColumn().getDataType());
+        Assertions.assertEquals("mood", enumEvent.getColumn().getSourceType());
+    }
+
+    @Test
+    void shouldMatchBaselineWithEnumColumn() {
+        CatalogTable baseline =
+                CatalogTable.of(
+                        TABLE_IDENTIFIER,
+                        TableSchema.builder()
+                                .column(
+                                        PhysicalColumn.builder()
+                                                .name("id")
+                                                .dataType(BasicType.INT_TYPE)
+                                                .nullable(false)
+                                                .sourceType("int4")
+                                                .build())
+                                .column(
+                                        PhysicalColumn.builder()
+                                                .name("m")
+                                                
.dataType(BasicType.STRING_TYPE)
+                                                .nullable(true)
+                                                .sourceType("inv.mood")
+                                                .build())
+                                .build(),
+                        Collections.emptyMap(),
+                        Collections.emptyList(),
+                        null,
+                        null);
+        Table relation =
+                Table.editor()
+                        .tableId(new TableId(null, SCHEMA_NAME, TABLE_NAME))
+                        .setPrimaryKeyNames(Collections.singletonList("id"))
+                        .setColumns(Arrays.asList(intColumn("id", 1), 
enumColumn("m", 2)))
+                        .create();
+
+        Assertions.assertTrue(
+                
PostgresRelationSchemaChangeResolver.hasSameCatalogSchema(baseline, relation));
+    }
+
     @Test
     void shouldIgnoreUnchangedRelationRecord() {
         SourceRecord record = createRecord(intColumn("id", 1), 
varcharColumn("name", 2, 64));
@@ -213,6 +266,18 @@ class PostgresRelationSchemaChangeResolverTest {
                 .create();
     }
 
+    // Debezium reports enum columns with the enum type name and Types.VARCHAR.
+    private Column enumColumn(String name, int position) {
+        return Column.editor()
+                .name(name)
+                .jdbcType(Types.VARCHAR)
+                .nativeType(16384)
+                .type("mood", "mood")
+                .position(position)
+                .optional(true)
+                .create();
+    }
+
     private Column unsupportedArrayColumn(String name, int position) {
         return Column.editor()
                 .name(name)
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
index d352f67c8e..a3f6809617 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalog.java
@@ -19,9 +19,11 @@ package 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.psql;
 
 import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.catalog.Column;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
 import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.catalog.exception.CatalogException;
 import org.apache.seatunnel.api.table.converter.BasicTypeDefine;
+import org.apache.seatunnel.api.table.type.BasicType;
 import org.apache.seatunnel.common.utils.JdbcUrlUtil;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.AbstractJdbcCatalog;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.utils.CatalogUtils;
@@ -42,11 +44,18 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
     public static final String TABLE_OPTION_TABLESPACE = "tablespace";
     public static final String TABLE_OPTION_FILLFACTOR = "fillfactor";
 
+    // pg_type.typtype of enum types
+    private static final String PG_TYPTYPE_ENUM = "e";
+
     private static final String SELECT_COLUMNS_SQL_TEMPLATE =
             "SELECT \n"
                     + "    a.attname AS column_name, \n"
                     + "\t\tt.typname as type_name,\n"
+                    + "\t\tt.typtype as type_type,\n"
                     + "    CASE \n"
+                    // Enum types are schema objects, format_type qualifies 
them when not on the
+                    // search_path.
+                    + "        WHEN t.typtype = 'e' THEN 
format_type(a.atttypid, NULL)\n"
                     + "        WHEN a.atttypmod = -1 THEN t.typname\n"
                     + "        WHEN t.typname = 'varchar' THEN t.typname || 
'(' || (a.atttypmod - 4) || ')'\n"
                     + "        WHEN t.typname = 'bpchar' THEN 'char' || '(' || 
(a.atttypmod - 4) || ')'\n"
@@ -132,21 +141,34 @@ public class PostgresCatalog extends AbstractJdbcCatalog {
         String columnName = resultSet.getString("column_name");
         String typeName = resultSet.getString("type_name");
         String fullTypeName = resultSet.getString("full_type_name");
+        String typeType = resultSet.getString("type_type");
         long columnLength = resultSet.getLong("column_length");
         int columnScale = resultSet.getInt("column_scale");
         String columnComment = resultSet.getString("column_comment");
         Object defaultValue = resultSet.getObject("default_value");
         boolean isNullable = resultSet.getString("is_nullable").equals("YES");
 
+        if (defaultValue != null && 
defaultValue.toString().contains("regclass")) {
+            defaultValue = null;
+        }
+        if (PG_TYPTYPE_ENUM.equals(typeType)) {
+            // Enum names are user-defined and may match built-in type names, 
map them directly.
+            return PhysicalColumn.builder()
+                    .name(columnName)
+                    .dataType(BasicType.STRING_TYPE)
+                    .sourceType(fullTypeName)
+                    .nullable(isNullable)
+                    .defaultValue(defaultValue)
+                    .comment(columnComment)
+                    .build();
+        }
+
         // dealingSpecialNumeric
         if (typeName.equals(PostgresTypeConverter.PG_NUMERIC) && columnLength 
< 1) {
             fullTypeName = "numeric(38,10)";
             columnLength = 38;
             columnScale = 10;
         }
-        if (defaultValue != null && 
defaultValue.toString().contains("regclass")) {
-            defaultValue = null;
-        }
 
         BasicTypeDefine typeDefine =
                 BasicTypeDefine.builder()
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverter.java
index abf624fe8b..08133ba83d 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverter.java
@@ -36,6 +36,8 @@ import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.sqlserver
 import com.google.auto.service.AutoService;
 import lombok.extern.slf4j.Slf4j;
 
+import java.sql.Types;
+
 // reference 
https://help.kingbase.com.cn/v8/development/sql-plsql/sql/datatype.html#id2
 @Slf4j
 @AutoService(TypeConverter.class)
@@ -53,6 +55,12 @@ public class KingbaseTypeConverter extends 
PostgresTypeConverter {
         return DatabaseIdentifier.KINGBASE;
     }
 
+    @Override
+    protected boolean isUserDefinedStringType(int sqlType) {
+        // VARCHAR type names unknown to PostgreSQL keep the Kingbase handling 
below.
+        return sqlType == Types.OTHER;
+    }
+
     @Override
     public Column convert(BasicTypeDefine typeDefine) {
         try {
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverter.java
index 465a587e3f..6d4a3b13ba 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverter.java
@@ -35,6 +35,7 @@ import com.google.auto.service.AutoService;
 import lombok.extern.slf4j.Slf4j;
 
 import java.sql.Types;
+import java.util.regex.Pattern;
 
 // reference http://www.postgres.cn/docs/13/datatype.html
 @Slf4j
@@ -120,6 +121,11 @@ public class PostgresTypeConverter implements 
TypeConverter<BasicTypeDefine> {
     public static final int MAX_VARCHAR_LENGTH = 10485760;
     public static final PostgresTypeConverter INSTANCE = new 
PostgresTypeConverter();
 
+    private static final String TYPE_NAME_PART =
+            "(\"[A-Za-z_][A-Za-z0-9_$]*\"|[A-Za-z_][A-Za-z0-9_$]*)";
+    private static final Pattern USER_DEFINED_TYPE_NAME =
+            Pattern.compile(TYPE_NAME_PART + "(\\." + TYPE_NAME_PART + ")?");
+
     @Override
     public String identifier() {
         return DatabaseIdentifier.POSTGRESQL;
@@ -285,9 +291,9 @@ public class PostgresTypeConverter implements 
TypeConverter<BasicTypeDefine> {
                 }
                 break;
             default:
-                if (typeDefine.getSqlType() == Types.OTHER) {
+                if (isUserDefinedStringType(typeDefine.getSqlType())) {
                     builder.dataType(BasicType.STRING_TYPE);
-                    builder.sourceType(typeDefine.getColumnType());
+                    
builder.sourceType(userDefinedSourceType(typeDefine.getColumnType()));
                     break;
                 }
                 throw CommonError.convertToSeaTunnelTypeError(
@@ -296,6 +302,26 @@ public class PostgresTypeConverter implements 
TypeConverter<BasicTypeDefine> {
         return builder.build();
     }
 
+    /**
+     * Whether a type name not handled above is read as STRING. The PostgreSQL 
JDBC driver and
+     * Debezium report user-defined types as {@link Types#OTHER}, and enum 
types and the built-in
+     * {@code name} type as {@link Types#VARCHAR}.
+     */
+    protected boolean isUserDefinedStringType(int sqlType) {
+        return sqlType == Types.OTHER || sqlType == Types.VARCHAR;
+    }
+
+    /**
+     * The reported name of a user-defined type is not escaped and may be 
copied into sink DDL, so
+     * it is kept only when it is a plain or quoted identifier, optionally 
schema-qualified.
+     */
+    private static String userDefinedSourceType(String typeName) {
+        if (typeName != null && 
USER_DEFINED_TYPE_NAME.matcher(typeName).matches()) {
+            return typeName;
+        }
+        return PG_TEXT;
+    }
+
     @Override
     public BasicTypeDefine reconvert(Column column) {
         BasicTypeDefine.BasicTypeDefineBuilder builder =
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverter.java
index f82291198c..3b4243030f 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverter.java
@@ -34,6 +34,8 @@ import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.psql.Post
 import com.google.auto.service.AutoService;
 import lombok.extern.slf4j.Slf4j;
 
+import java.sql.Types;
+
 // reference 
https://docs.aws.amazon.com/redshift/latest/dg/c_Supported_data_types.html
 @Slf4j
 @AutoService(TypeConverter.class)
@@ -73,6 +75,12 @@ public class RedshiftTypeConverter extends 
PostgresTypeConverter {
         return DatabaseIdentifier.REDSHIFT;
     }
 
+    @Override
+    protected boolean isUserDefinedStringType(int sqlType) {
+        // VARCHAR type names unknown to PostgreSQL stay unsupported for 
Redshift.
+        return sqlType == Types.OTHER;
+    }
+
     @Override
     public Column convert(BasicTypeDefine typeDefine) {
         PhysicalColumn.PhysicalColumnBuilder builder =
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalogBuildColumnTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalogBuildColumnTest.java
new file mode 100644
index 0000000000..a247d0d681
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/psql/PostgresCatalogBuildColumnTest.java
@@ -0,0 +1,126 @@
+/*
+ * 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.seatunnel.connectors.seatunnel.jdbc.catalog.psql;
+
+import org.apache.seatunnel.api.table.catalog.Column;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
+import org.apache.seatunnel.common.utils.JdbcUrlUtil;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.sql.ResultSet;
+import java.sql.SQLException;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+class PostgresCatalogBuildColumnTest {
+
+    private final PostgresCatalog catalog =
+            new PostgresCatalog(
+                    "Postgres",
+                    "postgres",
+                    "postgres",
+                    
JdbcUrlUtil.getUrlInfo("jdbc:postgresql://localhost:5432/test"),
+                    null,
+                    null);
+
+    @Test
+    void testBuildEnumColumnAsString() throws SQLException {
+        ResultSet resultSet = mockColumn("m", "mood", "inv.mood", "e", true);
+        
when(resultSet.getObject("default_value")).thenReturn("'happy'::inv.mood");
+        when(resultSet.getString("column_comment")).thenReturn("current mood");
+
+        Column column = catalog.buildColumn(resultSet);
+
+        Assertions.assertEquals("m", column.getName());
+        Assertions.assertEquals(BasicType.STRING_TYPE, column.getDataType());
+        Assertions.assertEquals("inv.mood", column.getSourceType());
+        Assertions.assertTrue(column.isNullable());
+        Assertions.assertEquals("'happy'::inv.mood", column.getDefaultValue());
+        Assertions.assertEquals("current mood", column.getComment());
+    }
+
+    @Test
+    void testBuildEnumNamedLikeBuiltInTypeAsString() throws SQLException {
+        Column dateEnum = catalog.buildColumn(mockColumn("d", "date", 
"inv.date", "e", true));
+        Column numericEnum =
+                catalog.buildColumn(mockColumn("n", "numeric", 
"inv.\"numeric\"", "e", true));
+
+        Assertions.assertEquals(BasicType.STRING_TYPE, dateEnum.getDataType());
+        Assertions.assertEquals("inv.date", dateEnum.getSourceType());
+        Assertions.assertEquals(BasicType.STRING_TYPE, 
numericEnum.getDataType());
+        Assertions.assertEquals("inv.\"numeric\"", 
numericEnum.getSourceType());
+    }
+
+    @Test
+    void testBuildEnumKeepsQuotedFormatTypeName() throws SQLException {
+        Column column =
+                catalog.buildColumn(mockColumn("w", "Weird Name", "inv.\"Weird 
Name\"", "e", true));
+
+        Assertions.assertEquals(BasicType.STRING_TYPE, column.getDataType());
+        Assertions.assertEquals("inv.\"Weird Name\"", column.getSourceType());
+    }
+
+    @Test
+    void testSelectColumnsSqlReadsEnumTypeType() {
+        String sql = catalog.getSelectColumnsSql(TablePath.of("test", "inv", 
"products"));
+
+        Assertions.assertTrue(sql.contains("t.typtype as type_type"));
+        Assertions.assertTrue(
+                sql.contains("WHEN t.typtype = 'e' THEN 
format_type(a.atttypid, NULL)"));
+    }
+
+    @Test
+    void testBuildBaseColumnUnchanged() throws SQLException {
+        ResultSet resultSet = mockColumn("id", "int4", "int4", "b", false);
+
+        Column column = catalog.buildColumn(resultSet);
+
+        Assertions.assertEquals(BasicType.INT_TYPE, column.getDataType());
+        Assertions.assertEquals("int4", column.getSourceType());
+        Assertions.assertFalse(column.isNullable());
+    }
+
+    @Test
+    void testBuildUnsupportedNonEnumColumnStillFails() throws SQLException {
+        ResultSet resultSet = mockColumn("v", "tsvector", "tsvector", "b", 
true);
+
+        Assertions.assertThrows(
+                SeaTunnelRuntimeException.class, () -> 
catalog.buildColumn(resultSet));
+    }
+
+    private static ResultSet mockColumn(
+            String columnName,
+            String typeName,
+            String fullTypeName,
+            String typeType,
+            boolean nullable)
+            throws SQLException {
+        ResultSet resultSet = mock(ResultSet.class);
+        when(resultSet.getString("column_name")).thenReturn(columnName);
+        when(resultSet.getString("type_name")).thenReturn(typeName);
+        when(resultSet.getString("full_type_name")).thenReturn(fullTypeName);
+        when(resultSet.getString("type_type")).thenReturn(typeType);
+        when(resultSet.getString("is_nullable")).thenReturn(nullable ? "YES" : 
"NO");
+        return resultSet;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverterTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverterTest.java
index 03302b48b9..cfad668924 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverterTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/kingbase/KingbaseTypeConverterTest.java
@@ -31,7 +31,23 @@ import 
org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.sql.Types;
+
 public class KingbaseTypeConverterTest {
+    @Test
+    public void testConvertUnsupportedVarcharSqlType() {
+        BasicTypeDefine<Object> typeDefine =
+                BasicTypeDefine.builder()
+                        .name("test")
+                        .columnType("aaa")
+                        .dataType("aaa")
+                        .sqlType(Types.VARCHAR)
+                        .build();
+        Assertions.assertThrows(
+                SeaTunnelRuntimeException.class,
+                () -> KingbaseTypeConverter.INSTANCE.convert(typeDefine));
+    }
+
     @Test
     public void testConvertUnsupported() {
         BasicTypeDefine<Object> typeDefine =
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverterTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverterTest.java
index 842b2abf8e..dfe1118a67 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverterTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/psql/PostgresTypeConverterTest.java
@@ -309,6 +309,74 @@ public class PostgresTypeConverterTest {
         Assertions.assertEquals("JobStatus", column.getSourceType());
     }
 
+    @Test
+    public void testConvertEnumReportedAsVarcharAsString() {
+        BasicTypeDefine<Object> typeDefine =
+                BasicTypeDefine.builder()
+                        .name("m")
+                        .columnType("\"inv\".\"mood\"")
+                        .dataType("\"inv\".\"mood\"")
+                        .sqlType(Types.VARCHAR)
+                        .nullable(true)
+                        .build();
+
+        Column column = PostgresTypeConverter.INSTANCE.convert(typeDefine);
+
+        Assertions.assertEquals(BasicType.STRING_TYPE, column.getDataType());
+        Assertions.assertEquals("\"inv\".\"mood\"", column.getSourceType());
+        Assertions.assertNull(column.getColumnLength());
+    }
+
+    @Test
+    public void testConvertUserDefinedTypeKeepsOnlyIdentifierNames() {
+        Assertions.assertEquals("mood", convertUserDefined("mood", 
Types.VARCHAR).getSourceType());
+        Assertions.assertEquals(
+                "inv.mood", convertUserDefined("inv.mood", 
Types.VARCHAR).getSourceType());
+        Assertions.assertEquals(
+                "\"inv\".\"JobStatus\"",
+                convertUserDefined("\"inv\".\"JobStatus\"", 
Types.OTHER).getSourceType());
+    }
+
+    @Test
+    public void testConvertUserDefinedTypeWithUnsafeNameAsText() {
+        String[] unsafeNames = {
+            "text;DROP TABLE t--", "\"a\"b\"", "inv.\"Weird Name\"", "a.b.c", 
"mood NULL"
+        };
+        for (String name : unsafeNames) {
+            for (int sqlType : new int[] {Types.VARCHAR, Types.OTHER}) {
+                Column column = convertUserDefined(name, sqlType);
+                Assertions.assertEquals(BasicType.STRING_TYPE, 
column.getDataType());
+                Assertions.assertEquals(PostgresTypeConverter.PG_TEXT, 
column.getSourceType());
+            }
+        }
+    }
+
+    private static Column convertUserDefined(String typeName, int sqlType) {
+        return PostgresTypeConverter.INSTANCE.convert(
+                BasicTypeDefine.builder()
+                        .name("m")
+                        .columnType(typeName)
+                        .dataType(typeName)
+                        .sqlType(sqlType)
+                        .build());
+    }
+
+    @Test
+    public void testTypeMapperMapsEnumReportedAsVarchar() throws SQLException {
+        ResultSetMetaData metadata = mock(ResultSetMetaData.class);
+        when(metadata.getColumnLabel(1)).thenReturn("m");
+        when(metadata.getColumnTypeName(1)).thenReturn("mood");
+        when(metadata.getColumnType(1)).thenReturn(Types.VARCHAR);
+        
when(metadata.isNullable(1)).thenReturn(ResultSetMetaData.columnNullable);
+        when(metadata.getPrecision(1)).thenReturn(Integer.MAX_VALUE);
+        when(metadata.getScale(1)).thenReturn(0);
+
+        Column column = new PostgresTypeMapper().mappingColumn(metadata, 1);
+
+        Assertions.assertEquals(BasicType.STRING_TYPE, column.getDataType());
+        Assertions.assertEquals("mood", column.getSourceType());
+    }
+
     @Test
     public void testConvertBinary() {
         BasicTypeDefine<Object> typeDefine =
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverterTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverterTest.java
index 8b3e9fc1f9..97d0721ea2 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverterTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/redshift/RedshiftTypeConverterTest.java
@@ -31,8 +31,24 @@ import org.junit.jupiter.api.Test;
 
 import lombok.extern.slf4j.Slf4j;
 
+import java.sql.Types;
+
 @Slf4j
 public class RedshiftTypeConverterTest {
+    @Test
+    public void testConvertUnsupportedVarcharSqlType() {
+        BasicTypeDefine<Object> typeDefine =
+                BasicTypeDefine.builder()
+                        .name("test")
+                        .columnType("aaa")
+                        .dataType("aaa")
+                        .sqlType(Types.VARCHAR)
+                        .build();
+        Assertions.assertThrows(
+                SeaTunnelRuntimeException.class,
+                () -> RedshiftTypeConverter.INSTANCE.convert(typeDefine));
+    }
+
     @Test
     public void testConvertUnsupported() {
         BasicTypeDefine<Object> typeDefine =

Reply via email to