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..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 @@ -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,32 @@ 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); + String showPartitionsSql = String.format("SHOW PARTITIONS FROM `%s`", tableName); + 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])) + }; + } + + 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 21bb654d86e..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; @@ -104,16 +105,10 @@ 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)) { - 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); } @@ -123,16 +118,15 @@ 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\"", 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,7 +214,34 @@ public boolean dropPartition(String partitionName) { } } - private Partition fromDorisPartition( + /** + * 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 { String partitionName = resultSet.getString(NAME); @@ -271,18 +292,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/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(); 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,