From f89385246fe830c78f01d0d1312fab9832b8ad42 Mon Sep 17 00:00:00 2001 From: yuqi Date: Tue, 21 Jul 2026 16:08:15 +0800 Subject: [PATCH 1/3] fix(doris): load partition assignments --- .../doris/operation/DorisTableOperations.java | 32 ++++++++++++++++++- .../DorisTablePartitionOperations.java | 17 +++++----- .../operation/TestDorisTableOperations.java | 13 ++++++++ 3 files changed, 53 insertions(+), 9 deletions(-) diff --git a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java index e82e44240bd..b0f88229886 100644 --- a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java +++ b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java @@ -67,7 +67,9 @@ import org.apache.gravitino.rel.indexes.Index; import org.apache.gravitino.rel.indexes.Indexes; import org.apache.gravitino.rel.partitions.ListPartition; +import org.apache.gravitino.rel.partitions.Partition; import org.apache.gravitino.rel.partitions.RangePartition; +import org.apache.gravitino.rel.types.Type; /** Table operations for Apache Doris. */ public class DorisTableOperations extends JdbcTableOperations { @@ -633,7 +635,35 @@ protected Transform[] getTablePartitioning( } Optional transform = DorisUtils.extractPartitionInfoFromSql(createTableSql.toString()); - return transform.map(t -> new Transform[] {t}).orElse(Transforms.EMPTY_TRANSFORM); + if (!transform.isPresent()) { + return Transforms.EMPTY_TRANSFORM; + } + + Transform partitioning = transform.get(); + Map columnTypes = + DorisTablePartitionOperations.getColumnTypes(connection, tableName, typeConverter); + List assignments = new ArrayList<>(); + String showPartitionsSql = String.format("SHOW PARTITIONS FROM `%s`", tableName); + try (Statement partitionStatement = connection.createStatement(); + ResultSet partitions = partitionStatement.executeQuery(showPartitionsSql)) { + while (partitions.next()) { + assignments.add( + DorisTablePartitionOperations.fromDorisPartition( + partitions, partitioning, columnTypes)); + } + } + + if (partitioning instanceof Transforms.ListTransform) { + String[][] fieldNames = ((Transforms.ListTransform) partitioning).fieldNames(); + return new Transform[] { + Transforms.list(fieldNames, assignments.toArray(new ListPartition[0])) + }; + } + + String[] fieldName = ((Transforms.RangeTransform) partitioning).fieldName(); + return new Transform[] { + Transforms.range(fieldName, assignments.toArray(new RangePartition[0])) + }; } } diff --git a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java index 21bb654d86e..fdc55c5c934 100644 --- a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java +++ b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java @@ -104,7 +104,7 @@ public String[] listPartitionNames() { public Partition[] listPartitions() { try (Connection connection = getConnection(loadedTable.databaseName())) { Transform partitionInfo = loadedTable.partitioning()[0]; - Map columnTypes = getColumnType(connection); + Map columnTypes = getColumnTypes(connection, loadedTable.name(), typeConverter); String showPartitionsSql = String.format("SHOW PARTITIONS FROM `%s`", loadedTable.name()); try (Statement statement = connection.createStatement(); ResultSet result = statement.executeQuery(showPartitionsSql)) { @@ -123,7 +123,7 @@ public Partition[] listPartitions() { public Partition getPartition(String partitionName) throws NoSuchPartitionException { try (Connection connection = getConnection(loadedTable.databaseName())) { Transform partitionInfo = loadedTable.partitioning()[0]; - Map columnTypes = getColumnType(connection); + Map columnTypes = getColumnTypes(connection, loadedTable.name(), typeConverter); String showPartitionsSql = String.format( "SHOW PARTITIONS FROM `%s` WHERE PartitionName = \"%s\"", @@ -220,7 +220,7 @@ public boolean dropPartition(String partitionName) { } } - private Partition fromDorisPartition( + static Partition fromDorisPartition( ResultSet resultSet, Transform partitionInfo, Map columnTypes) throws SQLException { String partitionName = resultSet.getString(NAME); @@ -271,18 +271,19 @@ private Partition fromDorisPartition( partitionName, lists.build().toArray(new Literal[0][0]), properties); } else { throw new UnsupportedOperationException( - String.format("%s is not a partitioned table", loadedTable.name())); + String.format("%s is not a supported partition transform", partitionInfo)); } } - private Map getColumnType(Connection connection) throws SQLException { + static Map getColumnTypes( + Connection connection, String tableName, JdbcTypeConverter typeConverter) + throws SQLException { DatabaseMetaData metaData = connection.getMetaData(); try (ResultSet result = - metaData.getColumns( - connection.getCatalog(), connection.getSchema(), loadedTable.name(), null)) { + metaData.getColumns(connection.getCatalog(), connection.getSchema(), tableName, null)) { ImmutableMap.Builder columnTypes = ImmutableMap.builder(); while (result.next()) { - if (Objects.equals(result.getString("TABLE_NAME"), loadedTable.name())) { + if (Objects.equals(result.getString("TABLE_NAME"), tableName)) { JdbcTypeConverter.JdbcTypeBean typeBean = new JdbcTypeConverter.JdbcTypeBean(result.getString("TYPE_NAME")); typeBean.setColumnSize(result.getInt("COLUMN_SIZE")); diff --git a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperations.java b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperations.java index a503eb3c5eb..87b04bc757b 100644 --- a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperations.java +++ b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperations.java @@ -29,6 +29,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import org.apache.gravitino.catalog.doris.converter.DorisTypeConverter; @@ -559,6 +560,12 @@ public void testCreatePartitionedTable() { distribution, indexes); JdbcTable rangePartitionTable = TABLE_OPERATIONS.load(databaseName, rangePartitionTableName); + RangePartition[] loadedRangeAssignments = + ((Transforms.RangeTransform) rangePartitionTable.partitioning()[0]).assignments(); + assertEquals(3, loadedRangeAssignments.length); + assertEquals( + Set.of("p1", "p2", "p3"), + Arrays.stream(loadedRangeAssignments).map(Partition::name).collect(Collectors.toSet())); assertionsTableInfo( rangePartitionTableName, tableComment, @@ -618,6 +625,12 @@ public void testCreatePartitionedTable() { distribution, indexes); JdbcTable listPartitionTable = TABLE_OPERATIONS.load(databaseName, listPartitionTableName); + ListPartition[] loadedListAssignments = + ((Transforms.ListTransform) listPartitionTable.partitioning()[0]).assignments(); + assertEquals(2, loadedListAssignments.length); + assertEquals( + Set.of("p1", "p2"), + Arrays.stream(loadedListAssignments).map(Partition::name).collect(Collectors.toSet())); assertionsTableInfo( listPartitionTableName, tableComment, From 42c7d1fbb4197e9f01ef6a4c29289ed701ffbed6 Mon Sep 17 00:00:00 2001 From: yuqi Date: Tue, 21 Jul 2026 16:20:15 +0800 Subject: [PATCH 2/3] test(doris): verify loaded partition assignments --- .../doris/integration/test/CatalogDorisIT.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java index e474c10fa76..8f804d43c38 100644 --- a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java +++ b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java @@ -812,6 +812,12 @@ void testCreatePartitionedTable() { null, new Transform[] {Transforms.range(new String[] {DORIS_COL_NAME4})}, loadedTable); + RangePartition[] rangeAssignments = + ((Transforms.RangeTransform) loadedTable.partitioning()[0]).assignments(); + assertEquals(3, rangeAssignments.length); + assertEquals( + Set.of("p1", "p2", "p3"), + Arrays.stream(rangeAssignments).map(Partition::name).collect(Collectors.toSet())); // assert partition info SupportsPartitions tablePartitionOperations = loadedTable.supportPartitions(); @@ -876,6 +882,12 @@ void testCreatePartitionedTable() { null, new Transform[] {Transforms.list(new String[][] {{DORIS_COL_NAME1}, {DORIS_COL_NAME4}})}, loadedTable); + ListPartition[] listAssignments = + ((Transforms.ListTransform) loadedTable.partitioning()[0]).assignments(); + assertEquals(2, listAssignments.length); + assertEquals( + Set.of("p4", "p5"), + Arrays.stream(listAssignments).map(Partition::name).collect(Collectors.toSet())); // assert partition info tablePartitionOperations = loadedTable.supportPartitions(); From ce16b7b3bcae80c303a9c65ca6233774e0741a78 Mon Sep 17 00:00:00 2001 From: yuqi Date: Wed, 22 Jul 2026 14:37:55 +0800 Subject: [PATCH 3/3] [#12093] refactor(doris): share a loadPartitions helper and harden the transform switch --- .../doris/operation/DorisTableOperations.java | 23 ++++----- .../DorisTablePartitionOperations.java | 47 ++++++++++++++----- 2 files changed, 44 insertions(+), 26 deletions(-) diff --git a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java index b0f88229886..8a56a4e52cf 100644 --- a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java +++ b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java @@ -642,28 +642,25 @@ protected Transform[] getTablePartitioning( Transform partitioning = transform.get(); Map columnTypes = DorisTablePartitionOperations.getColumnTypes(connection, tableName, typeConverter); - List assignments = new ArrayList<>(); String showPartitionsSql = String.format("SHOW PARTITIONS FROM `%s`", tableName); - try (Statement partitionStatement = connection.createStatement(); - ResultSet partitions = partitionStatement.executeQuery(showPartitionsSql)) { - while (partitions.next()) { - assignments.add( - DorisTablePartitionOperations.fromDorisPartition( - partitions, partitioning, columnTypes)); - } - } + List assignments = + DorisTablePartitionOperations.loadPartitions( + connection, showPartitionsSql, partitioning, columnTypes); if (partitioning instanceof Transforms.ListTransform) { String[][] fieldNames = ((Transforms.ListTransform) partitioning).fieldNames(); return new Transform[] { Transforms.list(fieldNames, assignments.toArray(new ListPartition[0])) }; + } else if (partitioning instanceof Transforms.RangeTransform) { + String[] fieldName = ((Transforms.RangeTransform) partitioning).fieldName(); + return new Transform[] { + Transforms.range(fieldName, assignments.toArray(new RangePartition[0])) + }; } - String[] fieldName = ((Transforms.RangeTransform) partitioning).fieldName(); - return new Transform[] { - Transforms.range(fieldName, assignments.toArray(new RangePartition[0])) - }; + throw new UnsupportedOperationException( + String.format("%s is not a supported partition transform", partitioning)); } } diff --git a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java index fdc55c5c934..a62154b6a87 100644 --- a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java +++ b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTablePartitionOperations.java @@ -40,6 +40,7 @@ import java.sql.Statement; import java.util.Arrays; import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.Objects; import java.util.regex.Matcher; @@ -106,14 +107,8 @@ public Partition[] listPartitions() { Transform partitionInfo = loadedTable.partitioning()[0]; Map columnTypes = getColumnTypes(connection, loadedTable.name(), typeConverter); String showPartitionsSql = String.format("SHOW PARTITIONS FROM `%s`", loadedTable.name()); - try (Statement statement = connection.createStatement(); - ResultSet result = statement.executeQuery(showPartitionsSql)) { - ImmutableList.Builder partitions = ImmutableList.builder(); - while (result.next()) { - partitions.add(fromDorisPartition(result, partitionInfo, columnTypes)); - } - return partitions.build().toArray(new Partition[0]); - } + return loadPartitions(connection, showPartitionsSql, partitionInfo, columnTypes) + .toArray(new Partition[0]); } catch (SQLException e) { throw exceptionConverter.toGravitinoException(e); } @@ -128,11 +123,10 @@ public Partition getPartition(String partitionName) throws NoSuchPartitionExcept String.format( "SHOW PARTITIONS FROM `%s` WHERE PartitionName = \"%s\"", loadedTable.name(), partitionName); - try (Statement statement = connection.createStatement(); - ResultSet result = statement.executeQuery(showPartitionsSql)) { - if (result.next()) { - return fromDorisPartition(result, partitionInfo, columnTypes); - } + List partitions = + loadPartitions(connection, showPartitionsSql, partitionInfo, columnTypes); + if (!partitions.isEmpty()) { + return partitions.get(0); } } catch (SQLException e) { throw exceptionConverter.toGravitinoException(e); @@ -220,6 +214,33 @@ public boolean dropPartition(String partitionName) { } } + /** + * Runs a {@code SHOW PARTITIONS} query and converts every returned row into a {@link Partition}, + * shared by table loading and the partition read operations. + * + * @param connection an open connection to the table's database + * @param showPartitionsSql the {@code SHOW PARTITIONS} statement to execute + * @param partitionInfo the table's partition transform (list or range) + * @param columnTypes the partition columns' types, keyed by column name + * @return the partitions returned by the query, in result-set order + * @throws SQLException if the query fails + */ + static List loadPartitions( + Connection connection, + String showPartitionsSql, + Transform partitionInfo, + Map columnTypes) + throws SQLException { + try (Statement statement = connection.createStatement(); + ResultSet result = statement.executeQuery(showPartitionsSql)) { + ImmutableList.Builder partitions = ImmutableList.builder(); + while (result.next()) { + partitions.add(fromDorisPartition(result, partitionInfo, columnTypes)); + } + return partitions.build(); + } + } + static Partition fromDorisPartition( ResultSet resultSet, Transform partitionInfo, Map columnTypes) throws SQLException {