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

davidzollo 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 7c7701dd33 [Fix][Connector-V2][JDBC] Detect Hive partition keys in 
metadata (#10706)
7c7701dd33 is described below

commit 7c7701dd33b248e5c90f99d83a59244e5c749c1e
Author: zhangqingsong <[email protected]>
AuthorDate: Tue May 26 16:51:50 2026 +0800

    [Fix][Connector-V2][JDBC] Detect Hive partition keys in metadata (#10706)
    
    Co-authored-by: zhangqingsong <[email protected]>
---
 .../jdbc/catalog/AbstractJdbcCatalog.java          |   7 +-
 .../seatunnel/jdbc/catalog/utils/CatalogUtils.java |  45 ++++++-
 .../jdbc/catalog/utils/JdbcIdentifierUtils.java    |  34 +++--
 .../jdbc/internal/dialect/JdbcDialect.java         |  13 ++
 .../jdbc/internal/dialect/hive/HiveDialect.java    |  49 +++++++
 .../internal/dialect/inceptor/InceptorDialect.java |   7 +
 .../seatunnel/jdbc/utils/JdbcCatalogUtils.java     |   7 +-
 .../jdbc/catalog/utils/CatalogUtilsTest.java       | 142 +++++++++++++++++++++
 .../internal/dialect/hive/HiveDialectTest.java     | 103 +++++++++++++++
 .../dialect/inceptor/InceptorDialectTest.java}     |  18 +--
 .../connectors/seatunnel/jdbc/JdbcHiveIT.java      |  84 +++++++++++-
 .../jdbc_hive_partition_source_and_assert.conf     |  92 +++++++++++++
 12 files changed, 572 insertions(+), 29 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/AbstractJdbcCatalog.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/AbstractJdbcCatalog.java
index d2f43f05cc..83cec9eb1f 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/AbstractJdbcCatalog.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/AbstractJdbcCatalog.java
@@ -232,7 +232,7 @@ public abstract class AbstractJdbcCatalog implements 
Catalog {
                     tableIdentifier,
                     tableSchemaBuilder.build(),
                     buildConnectorOptions(tablePath),
-                    Collections.emptyList(),
+                    getPartitionKeys(conn, tablePath),
                     comment.orElse(""),
                     catalogName);
 
@@ -254,6 +254,11 @@ public abstract class AbstractJdbcCatalog implements 
Catalog {
         return columnsBuilder;
     }
 
+    protected List<String> getPartitionKeys(Connection connection, TablePath 
tablePath)
+            throws SQLException {
+        return Collections.emptyList();
+    }
+
     protected void buildColumnsWithErrorCheck(
             TablePath tablePath, ResultSet resultSet, TableSchema.Builder 
builder)
             throws SQLException {
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtils.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtils.java
index f6681c5e08..bf66767003 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtils.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtils.java
@@ -30,6 +30,7 @@ import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.common.exception.CommonError;
 import org.apache.seatunnel.common.exception.CommonErrorCode;
 import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialectTypeMapper;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.dialectenum.FieldIdeEnum;
 
@@ -223,7 +224,25 @@ public class CatalogUtils {
     public static TableSchema getTableSchema(
             DatabaseMetaData metadata, TablePath tablePath, 
JdbcDialectTypeMapper typeMapper)
             throws SQLException {
-        Optional<PrimaryKey> primaryKey = getPrimaryKey(metadata, tablePath);
+        return getTableSchema(metadata, tablePath, typeMapper, 
getPrimaryKey(metadata, tablePath));
+    }
+
+    public static TableSchema getTableSchema(
+            DatabaseMetaData metadata, TablePath tablePath, JdbcDialect 
dialect)
+            throws SQLException {
+        Optional<PrimaryKey> primaryKey =
+                dialect.supportsPrimaryKeyMetadata()
+                        ? getPrimaryKey(metadata, tablePath)
+                        : Optional.empty();
+        return getTableSchema(metadata, tablePath, 
dialect.getJdbcDialectTypeMapper(), primaryKey);
+    }
+
+    private static TableSchema getTableSchema(
+            DatabaseMetaData metadata,
+            TablePath tablePath,
+            JdbcDialectTypeMapper typeMapper,
+            Optional<PrimaryKey> primaryKey)
+            throws SQLException {
         List<ConstraintKey> constraintKeys = getConstraintKeys(metadata, 
tablePath);
         List<Column> columns;
         try {
@@ -247,8 +266,30 @@ public class CatalogUtils {
     public static CatalogTable getCatalogTable(
             Connection connection, TablePath tablePath, JdbcDialectTypeMapper 
typeMapper)
             throws SQLException {
+        return getCatalogTable(connection, tablePath, typeMapper, new 
ArrayList<>());
+    }
+
+    public static CatalogTable getCatalogTable(
+            Connection connection,
+            TablePath tablePath,
+            JdbcDialectTypeMapper typeMapper,
+            List<String> partitionKeys)
+            throws SQLException {
         DatabaseMetaData metadata = connection.getMetaData();
         TableSchema tableSchema = getTableSchema(metadata, tablePath, 
typeMapper);
+        return getCatalogTable(tablePath, tableSchema, partitionKeys);
+    }
+
+    public static CatalogTable getCatalogTable(
+            Connection connection, TablePath tablePath, JdbcDialect dialect) 
throws SQLException {
+        List<String> partitionKeys = dialect.getPartitionKeys(connection, 
tablePath);
+        DatabaseMetaData metadata = connection.getMetaData();
+        TableSchema tableSchema = getTableSchema(metadata, tablePath, dialect);
+        return getCatalogTable(tablePath, tableSchema, partitionKeys);
+    }
+
+    private static CatalogTable getCatalogTable(
+            TablePath tablePath, TableSchema tableSchema, List<String> 
partitionKeys) {
         String catalogName = "jdbc_catalog";
         return CatalogTable.of(
                 TableIdentifier.of(
@@ -258,7 +299,7 @@ public class CatalogUtils {
                         tablePath.getTableName()),
                 tableSchema,
                 new HashMap<>(),
-                new ArrayList<>(),
+                partitionKeys,
                 "",
                 catalogName);
     }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/JdbcIdentifierUtils.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/JdbcIdentifierUtils.java
index a788f985a2..f67beed9b8 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/JdbcIdentifierUtils.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/JdbcIdentifierUtils.java
@@ -19,6 +19,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.jdbc.catalog.utils;
 
 import java.sql.DatabaseMetaData;
 import java.sql.SQLException;
+import java.sql.SQLFeatureNotSupportedException;
 import java.util.Locale;
 
 public final class JdbcIdentifierUtils {
@@ -44,18 +45,35 @@ public final class JdbcIdentifierUtils {
         if (metadata == null) {
             return IdentifierCaseStrategy.CASE_INSENSITIVE;
         }
-        if (metadata.supportsMixedCaseIdentifiers()) {
-            return IdentifierCaseStrategy.CASE_SENSITIVE;
-        }
-        if (metadata.storesLowerCaseIdentifiers()) {
-            return IdentifierCaseStrategy.LOWER_CASE;
-        }
-        if (metadata.storesUpperCaseIdentifiers()) {
-            return IdentifierCaseStrategy.UPPER_CASE;
+        try {
+            if (metadata.supportsMixedCaseIdentifiers()) {
+                return IdentifierCaseStrategy.CASE_SENSITIVE;
+            }
+            if (metadata.storesLowerCaseIdentifiers()) {
+                return IdentifierCaseStrategy.LOWER_CASE;
+            }
+            if (metadata.storesUpperCaseIdentifiers()) {
+                return IdentifierCaseStrategy.UPPER_CASE;
+            }
+        } catch (SQLException e) {
+            if (isUnsupportedIdentifierCaseMetadata(e)) {
+                // Hive JDBC 3.x may report this capability gap as a plain 
SQLException.
+                return IdentifierCaseStrategy.CASE_INSENSITIVE;
+            }
+            throw e;
         }
         return IdentifierCaseStrategy.CASE_INSENSITIVE;
     }
 
+    private static boolean isUnsupportedIdentifierCaseMetadata(SQLException e) 
{
+        if (e instanceof SQLFeatureNotSupportedException) {
+            return true;
+        }
+        // Keep this narrow: other SQLExceptions from metadata calls should 
still fail fast.
+        String message = e.getMessage();
+        return message != null && "Method not 
supported".equalsIgnoreCase(message);
+    }
+
     public static boolean identifierEquals(
             IdentifierCaseStrategy caseStrategy, String expected, String 
actual) {
         if (expected == null) {
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
index 17f9d5d49e..d943e22d8a 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java
@@ -99,6 +99,19 @@ public interface JdbcDialect extends Serializable {
      */
     JdbcDialectTypeMapper getJdbcDialectTypeMapper();
 
+    /**
+     * Whether this dialect can reliably read primary-key metadata through 
{@link
+     * java.sql.DatabaseMetaData#getPrimaryKeys(String, String, String)}.
+     */
+    default boolean supportsPrimaryKeyMetadata() {
+        return true;
+    }
+
+    default List<String> getPartitionKeys(Connection connection, TablePath 
tablePath)
+            throws SQLException {
+        return new ArrayList<>();
+    }
+
     default String hashModForField(String nativeType, String fieldName, int 
mod) {
         return hashModForField(fieldName, mod);
     }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialect.java
index 081735018d..06a7f6bc26 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialect.java
@@ -17,6 +17,9 @@
 
 package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.hive;
 
+import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
+
+import org.apache.seatunnel.api.table.catalog.TablePath;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.JdbcConnectionProvider;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
@@ -29,6 +32,8 @@ import java.sql.PreparedStatement;
 import java.sql.ResultSet;
 import java.sql.ResultSetMetaData;
 import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.List;
 import java.util.Optional;
 
 public class HiveDialect implements JdbcDialect {
@@ -48,6 +53,12 @@ public class HiveDialect implements JdbcDialect {
         return new HiveTypeMapper();
     }
 
+    @Override
+    public boolean supportsPrimaryKeyMetadata() {
+        // Hive 3.x can throw from getPrimaryKeys while partition metadata is 
still available.
+        return false;
+    }
+
     @Override
     public Optional<String> getUpsertStatement(
             String database, String tableName, String[] fieldNames, String[] 
pkNames) {
@@ -68,4 +79,42 @@ public class HiveDialect implements JdbcDialect {
             JdbcConnectionConfig jdbcConnectionConfig) {
         return new HiveJdbcConnectionProvider(jdbcConnectionConfig);
     }
+
+    @Override
+    public List<String> getPartitionKeys(Connection connection, TablePath 
tablePath)
+            throws SQLException {
+        List<String> partitionKeys = new ArrayList<>();
+        boolean partitionSection = false;
+        boolean partitionHeader = false;
+        String describeSql = "DESCRIBE " + tableIdentifier(tablePath);
+        try (PreparedStatement preparedStatement = 
connection.prepareStatement(describeSql);
+                ResultSet resultSet = preparedStatement.executeQuery()) {
+            while (resultSet.next()) {
+                String columnName = 
StringUtils.trimToEmpty(resultSet.getString(1));
+                if (StringUtils.isBlank(columnName)) {
+                    continue;
+                }
+                if ("# Partition Information".equalsIgnoreCase(columnName)) {
+                    partitionSection = true;
+                    partitionHeader = true;
+                    continue;
+                }
+                if (!partitionSection) {
+                    continue;
+                }
+                if (columnName.startsWith("#")) {
+                    if (!partitionHeader) {
+                        break;
+                    }
+                    partitionHeader = false;
+                    continue;
+                }
+                if (partitionHeader) {
+                    partitionHeader = false;
+                }
+                partitionKeys.add(columnName);
+            }
+        }
+        return partitionKeys;
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
index 9770fb63bd..ece4f1e3ff 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
@@ -28,6 +28,13 @@ public class InceptorDialect extends HiveDialect {
         return DatabaseIdentifier.INCEPTOR;
     }
 
+    @Override
+    public boolean supportsPrimaryKeyMetadata() {
+        // Inceptor shares Hive syntax, but it should not inherit Hive's 
primary-key metadata
+        // workaround because its JDBC metadata support is a separate contract.
+        return true;
+    }
+
     @Override
     public JdbcRowConverter getRowConverter() {
         return new InceptorJdbcRowConverter();
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/utils/JdbcCatalogUtils.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/utils/JdbcCatalogUtils.java
index 7d7908c821..ba07f04108 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/utils/JdbcCatalogUtils.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/utils/JdbcCatalogUtils.java
@@ -382,9 +382,7 @@ public class JdbcCatalogUtils {
             TablePath tablePath = 
jdbcDialect.parse(tableConfig.getTablePath());
             CatalogTable tableOfPath = null;
             try {
-                tableOfPath =
-                        CatalogUtils.getCatalogTable(
-                                connection, tablePath, 
jdbcDialect.getJdbcDialectTypeMapper());
+                tableOfPath = CatalogUtils.getCatalogTable(connection, 
tablePath, jdbcDialect);
             } catch (Exception e) {
                 // ignore
                 log.debug("User-defined table path: {}", tablePath);
@@ -408,8 +406,7 @@ public class JdbcCatalogUtils {
         }
         if (StringUtils.isNotEmpty(tableConfig.getTablePath())) {
             TablePath tablePath = 
jdbcDialect.parse(tableConfig.getTablePath());
-            return CatalogUtils.getCatalogTable(
-                    connection, tablePath, 
jdbcDialect.getJdbcDialectTypeMapper());
+            return CatalogUtils.getCatalogTable(connection, tablePath, 
jdbcDialect);
         }
 
         return getCatalogTable(connection, tableConfig.getQuery(), 
jdbcDialect);
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtilsTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtilsTest.java
index eed75c213f..ec5650662d 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtilsTest.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/catalog/utils/CatalogUtilsTest.java
@@ -26,6 +26,7 @@ import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.converter.BasicTypeDefine;
 import org.apache.seatunnel.api.table.type.BasicType;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialectTypeMapper;
 
 import org.junit.jupiter.api.Assertions;
@@ -36,6 +37,7 @@ import java.sql.PreparedStatement;
 import java.sql.ResultSetMetaData;
 import java.sql.SQLException;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
@@ -54,6 +56,55 @@ public class CatalogUtilsTest {
         Assertions.assertEquals("testfdawe_", 
primaryKey.get().getPrimaryKey());
     }
 
+    @Test
+    void testGetTableSchemaPropagatesPrimaryKeyMetadataFailure() {
+        TestDatabaseMetaData metadata =
+                new TestDatabaseMetaData() {
+                    @Override
+                    public java.sql.ResultSet getPrimaryKeys(
+                            String catalog, String schema, String table) 
throws SQLException {
+                        throw new SQLException("getPrimaryKeys is not 
supported");
+                    }
+                };
+
+        Assertions.assertThrows(
+                SQLException.class,
+                () ->
+                        CatalogUtils.getTableSchema(
+                                metadata,
+                                TablePath.of("test.test"),
+                                new JdbcDialectTypeMapper() {}));
+    }
+
+    @Test
+    void testGetCatalogTableCanSkipPrimaryKeyMetadata() throws SQLException {
+        Connection connection =
+                new TestConnection() {
+                    @Override
+                    public java.sql.DatabaseMetaData getMetaData() {
+                        return new TestDatabaseMetaData() {
+                            @Override
+                            public java.sql.ResultSet getPrimaryKeys(
+                                    String catalog, String schema, String 
table)
+                                    throws SQLException {
+                                throw new SQLException("getPrimaryKeys is not 
supported");
+                            }
+                        };
+                    }
+                };
+
+        TablePath tablePath = TablePath.of("test.test");
+        JdbcDialect dialect = mock(JdbcDialect.class);
+        when(dialect.supportsPrimaryKeyMetadata()).thenReturn(false);
+        when(dialect.getJdbcDialectTypeMapper()).thenReturn(new 
JdbcDialectTypeMapper() {});
+        when(dialect.getPartitionKeys(connection, 
tablePath)).thenReturn(Arrays.asList("dt", "hr"));
+
+        CatalogTable catalogTable = CatalogUtils.getCatalogTable(connection, 
tablePath, dialect);
+
+        Assertions.assertNull(catalogTable.getTableSchema().getPrimaryKey());
+        Assertions.assertEquals(Arrays.asList("dt", "hr"), 
catalogTable.getPartitionKeys());
+    }
+
     @Test
     void testConstraintKeysNameWithOutSpecialChar() throws SQLException {
         List<ConstraintKey> constraintKeys =
@@ -347,6 +398,72 @@ public class CatalogUtilsTest {
         Assertions.assertEquals(Collections.emptyList(), 
tableSchema.getColumns());
     }
 
+    @Test
+    void testGetTableSchemaFallsBackWhenIdentifierCaseMetadataUnsupported() 
throws SQLException {
+        TestDatabaseMetaData metadata =
+                new TestDatabaseMetaData() {
+                    @Override
+                    public boolean supportsMixedCaseIdentifiers() throws 
SQLException {
+                        // Hive JDBC 3.1.3 throws a plain SQLException for 
this metadata API on
+                        // JDK 8, not SQLFeatureNotSupportedException.
+                        throw new SQLException("Method not supported");
+                    }
+
+                    @Override
+                    public java.sql.ResultSet getColumns(
+                            String catalog,
+                            String schemaPattern,
+                            String tableNamePattern,
+                            String columnNamePattern)
+                            throws SQLException {
+                        List<Map<String, Object>> value = new ArrayList<>();
+                        value.add(
+                                new HashMap<String, Object>() {
+                                    {
+                                        put("TABLE_NAME", "USER_INFO");
+                                        put("TABLE_SCHEM", "public");
+                                        put("COLUMN_NAME", "id");
+                                        put("DATA_TYPE", 1);
+                                        put("TYPE_NAME", "INT");
+                                        put("COLUMN_SIZE", 11);
+                                        put("DECIMAL_DIGITS", 0);
+                                        put("NULLABLE", 0);
+                                        put("REMARKS", "id comment");
+                                    }
+                                });
+                        return new TestResultSet(value);
+                    }
+                };
+
+        TablePath tablePath = TablePath.of("test_db", "public", "user_info");
+        TableSchema tableSchema =
+                CatalogUtils.getTableSchema(metadata, tablePath, new 
JdbcDialectTypeMapper() {});
+
+        Assertions.assertEquals(1, tableSchema.getColumns().size());
+        Assertions.assertEquals("id", 
tableSchema.getColumns().get(0).getName());
+    }
+
+    @Test
+    void testGetTableSchemaPropagatesIdentifierCaseMetadataFailure() {
+        TestDatabaseMetaData metadata =
+                new TestDatabaseMetaData() {
+                    @Override
+                    public boolean supportsMixedCaseIdentifiers() throws 
SQLException {
+                        throw new SQLException("connection broken");
+                    }
+                };
+
+        SQLException exception =
+                Assertions.assertThrows(
+                        SQLException.class,
+                        () ->
+                                CatalogUtils.getTableSchema(
+                                        metadata,
+                                        TablePath.of("test_db", "public", 
"user_info"),
+                                        new JdbcDialectTypeMapper() {}));
+        Assertions.assertEquals("connection broken", exception.getMessage());
+    }
+
     @Test
     void testGetTableSchemaStoresUpperCaseIdentifiersCanMatchLowerCaseInput() 
throws SQLException {
         TestDatabaseMetaData metadata =
@@ -478,4 +595,29 @@ public class CatalogUtilsTest {
 
         Assertions.assertNull(catalogTable.getTableSchema().getPrimaryKey());
     }
+
+    @Test
+    void testGetCatalogTableKeepsPartitionKeys() throws SQLException {
+        Connection connection =
+                new TestConnection() {
+                    @Override
+                    public java.sql.DatabaseMetaData getMetaData() {
+                        return new TestDatabaseMetaData();
+                    }
+                };
+
+        CatalogTable catalogTable =
+                CatalogUtils.getCatalogTable(
+                        connection,
+                        TablePath.of("test.test"),
+                        new JdbcDialectTypeMapper() {
+                            @Override
+                            public Column mappingColumn(BasicTypeDefine 
typeDefine) {
+                                return 
JdbcDialectTypeMapper.super.mappingColumn(typeDefine);
+                            }
+                        },
+                        Arrays.asList("dt", "hr"));
+
+        Assertions.assertEquals(Arrays.asList("dt", "hr"), 
catalogTable.getPartitionKeys());
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialectTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialectTest.java
new file mode 100644
index 0000000000..9ae0ca3517
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/hive/HiveDialectTest.java
@@ -0,0 +1,103 @@
+/*
+ * 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.internal.dialect.hive;
+
+import org.apache.seatunnel.api.table.catalog.TablePath;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.util.Arrays;
+import java.util.Collections;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+class HiveDialectTest {
+
+    @Test
+    void testHiveDoesNotSupportPrimaryKeyMetadata() {
+        Assertions.assertFalse(new HiveDialect().supportsPrimaryKeyMetadata());
+    }
+
+    @Test
+    void testGetPartitionKeysFromDescribeOutput() throws Exception {
+        HiveDialect hiveDialect = new HiveDialect();
+        Connection connection = mock(Connection.class);
+        PreparedStatement preparedStatement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+
+        when(connection.prepareStatement("DESCRIBE test_db.test_table"))
+                .thenReturn(preparedStatement);
+        when(preparedStatement.executeQuery()).thenReturn(resultSet);
+        when(resultSet.next()).thenReturn(true, true, true, true, true, false);
+        when(resultSet.getString(1))
+                .thenReturn("id", "# Partition Information", "# col_name", 
"dt", "hr");
+
+        Assertions.assertEquals(
+                Arrays.asList("dt", "hr"),
+                hiveDialect.getPartitionKeys(connection, 
TablePath.of("test_db.test_table")));
+    }
+
+    @Test
+    void testGetPartitionKeysStopsAtDetailedTableInformation() throws 
Exception {
+        HiveDialect hiveDialect = new HiveDialect();
+        Connection connection = mock(Connection.class);
+        PreparedStatement preparedStatement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+
+        when(connection.prepareStatement("DESCRIBE test_db.test_table"))
+                .thenReturn(preparedStatement);
+        when(preparedStatement.executeQuery()).thenReturn(resultSet);
+        when(resultSet.next()).thenReturn(true, true, true, true, true, true, 
true, false);
+        when(resultSet.getString(1))
+                .thenReturn(
+                        "id",
+                        "# Partition Information",
+                        "# col_name",
+                        "dt",
+                        "hr",
+                        "# Detailed Table Information",
+                        "Database: default");
+
+        Assertions.assertEquals(
+                Arrays.asList("dt", "hr"),
+                hiveDialect.getPartitionKeys(connection, 
TablePath.of("test_db.test_table")));
+    }
+
+    @Test
+    void testGetPartitionKeysWithoutPartitionSection() throws Exception {
+        HiveDialect hiveDialect = new HiveDialect();
+        Connection connection = mock(Connection.class);
+        PreparedStatement preparedStatement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+
+        when(connection.prepareStatement("DESCRIBE test_db.test_table"))
+                .thenReturn(preparedStatement);
+        when(preparedStatement.executeQuery()).thenReturn(resultSet);
+        when(resultSet.next()).thenReturn(true, false);
+        when(resultSet.getString(1)).thenReturn("id");
+
+        Assertions.assertEquals(
+                Collections.emptyList(),
+                hiveDialect.getPartitionKeys(connection, 
TablePath.of("test_db.test_table")));
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialectTest.java
similarity index 62%
copy from 
seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
copy to 
seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialectTest.java
index 9770fb63bd..543befd944 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialect.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/inceptor/InceptorDialectTest.java
@@ -17,19 +17,13 @@
 
 package 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.inceptor;
 
-import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
-import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.DatabaseIdentifier;
-import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.hive.HiveDialect;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
 
-public class InceptorDialect extends HiveDialect {
+class InceptorDialectTest {
 
-    @Override
-    public String dialectName() {
-        return DatabaseIdentifier.INCEPTOR;
-    }
-
-    @Override
-    public JdbcRowConverter getRowConverter() {
-        return new InceptorJdbcRowConverter();
+    @Test
+    void testSupportsPrimaryKeyMetadata() {
+        Assertions.assertTrue(new 
InceptorDialect().supportsPrimaryKeyMetadata());
     }
 }
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcHiveIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcHiveIT.java
index 08e37d3c92..9fb14c8f9a 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcHiveIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/JdbcHiveIT.java
@@ -24,7 +24,12 @@ import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
 import org.apache.seatunnel.common.utils.ExceptionUtils;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcSourceTableConfig;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.source.JdbcSourceTable;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.utils.JdbcCatalogUtils;
 
+import org.junit.jupiter.api.Assertions;
 import org.testcontainers.containers.GenericContainer;
 import org.testcontainers.containers.output.Slf4jLogConsumer;
 import org.testcontainers.utility.DockerLoggerFactory;
@@ -45,6 +50,7 @@ public class JdbcHiveIT extends AbstractJdbcIT {
     private static final String HIVE_DATABASE = "default";
 
     private static final String HIVE_SOURCE = "hive_e2e_source_table";
+    private static final String HIVE_PARTITION_TABLE = 
"hive_e2e_partition_table";
     private static final String HIVE_USERNAME = "root";
     private static final String HIVE_PASSWORD = null;
     private static final int HIVE_PORT = 10000;
@@ -53,7 +59,9 @@ public class JdbcHiveIT extends AbstractJdbcIT {
     private static final String DRIVER_CLASS = 
"org.apache.hive.jdbc.HiveDriver";
 
     private static final List<String> CONFIG_FILE =
-            Lists.newArrayList("/jdbc_hive_source_and_assert.conf");
+            Lists.newArrayList(
+                    "/jdbc_hive_source_and_assert.conf",
+                    "/jdbc_hive_partition_source_and_assert.conf");
     private static final String CREATE_SQL =
             "CREATE TABLE hive_e2e_source_table"
                     + "("
@@ -75,6 +83,18 @@ public class JdbcHiveIT extends AbstractJdbcIT {
                     + "    numeric_column          NUMERIC(10, 2)"
                     + ")";
 
+    private static final String PARTITION_TABLE_CREATE_SQL =
+            "CREATE TABLE IF NOT EXISTS "
+                    + HIVE_DATABASE
+                    + "."
+                    + HIVE_PARTITION_TABLE
+                    + "("
+                    + "    id                      INT,"
+                    + "    name                    STRING,"
+                    + "    amount                  DOUBLE"
+                    + ")"
+                    + " PARTITIONED BY (dt STRING, hr STRING)";
+
     @Override
     JdbcCase getJdbcCase() {
         Map<String, String> containerEnv = new HashMap<>();
@@ -108,6 +128,7 @@ public class JdbcHiveIT extends AbstractJdbcIT {
                             buildTableInfoWithSchema(
                                     jdbcCase.getDatabase(), 
jdbcCase.getSourceTable()));
             statement.execute(createSource);
+            statement.execute(PARTITION_TABLE_CREATE_SQL);
         } catch (Exception exception) {
             log.error(ExceptionUtils.getMessage(exception));
             throw new 
SeaTunnelRuntimeException(JdbcITErrorCode.CREATE_TABLE_FAILED, exception);
@@ -136,6 +157,26 @@ public class JdbcHiveIT extends AbstractJdbcIT {
                                 + "        42.10,"
                                 + "        42.12)");
             }
+            // Insert data into partitioned table across 3 dates and 2 hours
+            String[] dates = {"2023-09-01", "2023-09-02", "2023-09-03"};
+            String[] hours = {"00", "01"};
+            int id = 1;
+            for (String dt : dates) {
+                for (String hr : hours) {
+                    statement.execute(
+                            String.format(
+                                    "INSERT INTO %s.%s PARTITION(dt='%s', 
hr='%s') "
+                                            + "VALUES(%d, 'name_%d', %.2f)",
+                                    HIVE_DATABASE,
+                                    HIVE_PARTITION_TABLE,
+                                    dt,
+                                    hr,
+                                    id,
+                                    id,
+                                    id * 10.5));
+                    id++;
+                }
+            }
         } catch (Exception exception) {
             log.error(ExceptionUtils.getMessage(exception));
             throw new 
SeaTunnelRuntimeException(JdbcITErrorCode.INSERT_DATA_FAILED, exception);
@@ -169,4 +210,45 @@ public class JdbcHiveIT extends AbstractJdbcIT {
     public void clearTable(String schema, String table) {
         // do nothing.
     }
+
+    @Override
+    void checkResult(
+            String executeKey,
+            org.apache.seatunnel.e2e.common.container.TestContainer container,
+            org.testcontainers.containers.Container.ExecResult execResult) {
+        try {
+            TablePath partitionTablePath = TablePath.of(HIVE_DATABASE + "." + 
HIVE_PARTITION_TABLE);
+            JdbcSourceTableConfig tableConfig =
+                    JdbcSourceTableConfig.builder()
+                            .tablePath(partitionTablePath.getFullName())
+                            .useSelectCount(false)
+                            .build();
+            Map<TablePath, JdbcSourceTable> tables =
+                    JdbcCatalogUtils.getTables(
+                            JdbcConnectionConfig.builder()
+                                    .url(jdbcCase.getJdbcUrl().replace(HOST, 
dbServer.getHost()))
+                                    .driverName(jdbcCase.getDriverClass())
+                                    .username(jdbcCase.getUserName())
+                                    .password(jdbcCase.getPassword())
+                                    .build(),
+                            Lists.newArrayList(tableConfig));
+            JdbcSourceTable partitionSourceTable = 
tables.get(partitionTablePath);
+            Assertions.assertNotNull(
+                    partitionSourceTable,
+                    String.format(
+                            "JdbcCatalogUtils should load partitioned table 
%s",
+                            partitionTablePath.getFullName()));
+            List<String> partitionKeys = 
partitionSourceTable.getCatalogTable().getPartitionKeys();
+            log.info(
+                    "Detected catalog partition keys for {}: {}",
+                    partitionTablePath.getFullName(),
+                    partitionKeys);
+            Assertions.assertEquals(
+                    Lists.newArrayList("dt", "hr"),
+                    partitionKeys,
+                    "JdbcCatalogUtils should carry Hive partition keys into 
CatalogTable");
+        } catch (Exception e) {
+            Assertions.fail("Failed to load Hive partition keys from 
JdbcCatalogUtils", e);
+        }
+    }
 }
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/resources/jdbc_hive_partition_source_and_assert.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/resources/jdbc_hive_partition_source_and_assert.conf
new file mode 100644
index 0000000000..67f0148c49
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-jdbc-e2e/connector-jdbc-e2e-part-3/src/test/resources/jdbc_hive_partition_source_and_assert.conf
@@ -0,0 +1,92 @@
+#
+# 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.
+#
+######
+###### This config tests Hive partition key detection via table_path mode
+######
+
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+    Jdbc {
+        url = "jdbc:hive2://e2ehivejdbc:10000/default"
+        username = "root"
+        driver = "org.apache.hive.jdbc.HiveDriver"
+        table_path = "default.hive_e2e_partition_table"
+    }
+}
+
+transform {
+}
+
+sink{
+  assert {
+    rules =
+      {
+        row_rules = [
+          {
+            rule_type = MAX_ROW
+            rule_value = 6
+          },
+          {
+            rule_type = MIN_ROW
+            rule_value = 6
+          }
+        ],
+        # Keep field_value present for type-only checks because the assert 
executor expects a
+        # non-null rule list when field_rules are configured.
+        field_rules = [
+        {
+          field_name = id
+          field_type = int
+          field_value = [
+            {
+              rule_type = MIN
+              rule_value = 1
+            },
+            {
+              rule_type = MAX
+              rule_value = 6
+            }
+          ]
+        },
+        {
+          field_name = name
+          field_type = string
+          field_value = []
+        },
+        {
+          field_name = amount
+          field_type = "double"
+          field_value = []
+        },
+        {
+          field_name = dt
+          field_type = string
+          field_value = []
+        },
+        {
+          field_name = hr
+          field_type = string
+          field_value = []
+        }
+        ]
+      }
+  }
+}


Reply via email to