diff --git a/paimon-hive/paimon-hive-catalog/src/main/java/org/apache/paimon/hive/HiveCatalog.java b/paimon-hive/paimon-hive-catalog/src/main/java/org/apache/paimon/hive/HiveCatalog.java index 926fc024f27f..5811796e5a2a 100644 --- a/paimon-hive/paimon-hive-catalog/src/main/java/org/apache/paimon/hive/HiveCatalog.java +++ b/paimon-hive/paimon-hive-catalog/src/main/java/org/apache/paimon/hive/HiveCatalog.java @@ -595,47 +595,109 @@ public void markDonePartitions(Identifier identifier, List> public List listPartitions(Identifier identifier) throws TableNotExistException { FileStoreTable table = (FileStoreTable) getTable(identifier); - String tagToPartitionField = table.coreOptions().tagToPartitionField(); + CoreOptions coreOptions = table.coreOptions(); + String tagToPartitionField = coreOptions.tagToPartitionField(); + List partitions; if (tagToPartitionField != null) { try { - List partitions = listPartitionsFromHms(identifier); - return partitions.stream() - .map( - part -> { - Map parameters = part.getParameters(); - long recordCount = - Long.parseLong( - parameters.getOrDefault(NUM_ROWS_PROP, "1")); - long fileSizeInBytes = - Long.parseLong( - parameters.getOrDefault(TOTAL_SIZE_PROP, "1")); - long fileCount = - Long.parseLong( - parameters.getOrDefault(NUM_FILES_PROP, "1")); - long lastFileCreationTime = - Long.parseLong( - parameters.getOrDefault( - LAST_UPDATE_TIME_PROP, - System.currentTimeMillis() + "")); - int totalBuckets = - Integer.parseInt( - parameters.getOrDefault(TOTAL_BUCKETS, "0")); - return new org.apache.paimon.partition.Partition( - Collections.singletonMap( - tagToPartitionField, part.getValues().get(0)), - recordCount, - fileSizeInBytes, - fileCount, - lastFileCreationTime, - totalBuckets, - false); - }) - .collect(Collectors.toList()); + List hivePartitions = listPartitionsFromHms(identifier); + partitions = + hivePartitions.stream() + .map( + part -> { + Map parameters = part.getParameters(); + long recordCount = + Long.parseLong( + parameters.getOrDefault( + NUM_ROWS_PROP, "1")); + long fileSizeInBytes = + Long.parseLong( + parameters.getOrDefault( + TOTAL_SIZE_PROP, "1")); + long fileCount = + Long.parseLong( + parameters.getOrDefault( + NUM_FILES_PROP, "1")); + long lastFileCreationTime = + Long.parseLong( + parameters.getOrDefault( + LAST_UPDATE_TIME_PROP, + System.currentTimeMillis() + + "")); + int totalBuckets = + Integer.parseInt( + parameters.getOrDefault( + TOTAL_BUCKETS, "0")); + return new org.apache.paimon.partition.Partition( + Collections.singletonMap( + tagToPartitionField, + part.getValues().get(0)), + recordCount, + fileSizeInBytes, + fileCount, + lastFileCreationTime, + totalBuckets, + false); + }) + .collect(Collectors.toList()); } catch (Exception e) { throw new RuntimeException(e); } + } else { + partitions = listPartitionsFromFileSystem(table); + } + + if (coreOptions.partitionedTableInMetastore() + && coreOptions + .partitionMarkDoneActions() + .contains(CoreOptions.PartitionMarkDoneAction.MARK_EVENT)) { + return withDoneStatus(identifier, partitions); } - return listPartitionsFromFileSystem(table); + return partitions; + } + + private List withDoneStatus( + Identifier identifier, List partitions) + throws TableNotExistException { + try { + return clients() + .run( + client -> { + List result = + new ArrayList<>(partitions.size()); + for (org.apache.paimon.partition.Partition partition : partitions) { + boolean done = + client.isPartitionMarkedForEvent( + identifier.getDatabaseName(), + identifier.getTableName(), + partition.spec(), + PartitionEventType.LOAD_DONE); + result.add(copyWithDone(partition, done)); + } + return result; + }); + } catch (UnknownTableException e) { + throw new TableNotExistException(identifier); + } catch (TException | InterruptedException e) { + throw new RuntimeException(e); + } + } + + private org.apache.paimon.partition.Partition copyWithDone( + org.apache.paimon.partition.Partition partition, boolean done) { + return new org.apache.paimon.partition.Partition( + partition.spec(), + partition.recordCount(), + partition.fileSizeInBytes(), + partition.fileCount(), + partition.lastFileCreationTime(), + partition.totalBuckets(), + done, + partition.createdAt(), + partition.createdBy(), + partition.updatedAt(), + partition.updatedBy(), + partition.options()); } @VisibleForTesting diff --git a/paimon-hive/paimon-hive-catalog/src/test/java/org/apache/paimon/hive/HiveCatalogTest.java b/paimon-hive/paimon-hive-catalog/src/test/java/org/apache/paimon/hive/HiveCatalogTest.java index e1053c47152c..d158f0398c32 100644 --- a/paimon-hive/paimon-hive-catalog/src/test/java/org/apache/paimon/hive/HiveCatalogTest.java +++ b/paimon-hive/paimon-hive-catalog/src/test/java/org/apache/paimon/hive/HiveCatalogTest.java @@ -25,6 +25,8 @@ import org.apache.paimon.catalog.CatalogTestBase; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.client.ClientPool; +import org.apache.paimon.data.BinaryString; +import org.apache.paimon.data.GenericRow; import org.apache.paimon.fs.Path; import org.apache.paimon.options.CatalogOptions; import org.apache.paimon.options.Options; @@ -33,6 +35,9 @@ import org.apache.paimon.schema.Schema; import org.apache.paimon.schema.SchemaChange; import org.apache.paimon.table.object.ObjectTable; +import org.apache.paimon.table.sink.BatchTableCommit; +import org.apache.paimon.table.sink.BatchTableWrite; +import org.apache.paimon.table.sink.BatchWriteBuilder; import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.utils.CommonTestUtils; @@ -579,6 +584,48 @@ public void testTagToPartitionTable() throws Exception { Collections.singletonMap("dt", "20250101")); } + @Test + public void testListPartitionsWithDoneStatus() throws Exception { + String databaseName = "testListPartitionsWithDoneStatus"; + catalog.createDatabase(databaseName, false); + Identifier identifier = Identifier.create(databaseName, "table"); + catalog.createTable( + identifier, + Schema.newBuilder() + .option(METASTORE_PARTITIONED_TABLE.key(), "true") + .option(CoreOptions.PARTITION_MARK_DONE_ACTION.key(), "mark-event") + .column("col", DataTypes.INT()) + .column("dt", DataTypes.STRING()) + .partitionKeys("dt") + .build(), + false); + + List> partitionSpecs = + Arrays.asList( + Collections.singletonMap("dt", "20250101"), + Collections.singletonMap("dt", "20250102")); + BatchWriteBuilder writeBuilder = catalog.getTable(identifier).newBatchWriteBuilder(); + try (BatchTableWrite write = writeBuilder.newWrite(); + BatchTableCommit commit = writeBuilder.newCommit()) { + for (Map partitionSpec : partitionSpecs) { + write.write(GenericRow.of(0, BinaryString.fromString(partitionSpec.get("dt")))); + } + commit.commit(write.prepareCommit()); + } + + assertThat(catalog.listPartitions(identifier)).allMatch(partition -> !partition.done()); + + catalog.markDonePartitions(identifier, Collections.singletonList(partitionSpecs.get(0))); + + Map, Boolean> doneByPartition = new HashMap<>(); + for (Partition partition : catalog.listPartitions(identifier)) { + doneByPartition.put(partition.spec(), partition.done()); + } + assertThat(doneByPartition) + .containsEntry(partitionSpecs.get(0), true) + .containsEntry(partitionSpecs.get(1), false); + } + @Test public void testCreateTableWithBlob() throws Exception { String databaseName = "testCreateTableWithBlob";