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 =