danny0405 commented on code in PR #18372:
URL: https://github.com/apache/hudi/pull/18372#discussion_r3541091148
##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java:
##########
@@ -1116,102 +1110,15 @@ public BatchMetadataConversionFunction(String
instantTime, HoodieCommitMetadata
}
@Override
- public Map<String, HoodieData<HoodieRecord>> convertMetadata() {
- Map<String, HoodieData<HoodieRecord>> partitionToRecordMap =
- HoodieMetadataWriteUtils.convertMetadataToRecords(
- engineContext, dataWriteConfig, commitMetadata, instantTime,
dataMetaClient, getTableMetadata(),
- dataWriteConfig.getMetadataConfig(),
- partitionsToUpdate, dataWriteConfig.getBloomFilterType(),
- dataWriteConfig.getBloomIndexParallelism(),
dataWriteConfig.getWritesFileIdEncoding(), getEngineType(),
- Option.of(dataWriteConfig.getRecordMerger().getRecordType()));
-
- // Updates for record index are created by parsing the WriteStatus which
is a hudi-client object. Hence, we cannot yet move this code
- // to the HoodieTableMetadataUtil class in hudi-common.
- if (partitionsToUpdate.contains(RECORD_INDEX.getPartitionPath())) {
- HoodieData<HoodieRecord> additionalUpdates =
getRecordIndexAdditionalUpserts(partitionToRecordMap.get(RECORD_INDEX.getPartitionPath()),
commitMetadata);
- partitionToRecordMap.put(RECORD_INDEX.getPartitionPath(),
partitionToRecordMap.get(RECORD_INDEX.getPartitionPath()).union(additionalUpdates));
- }
- if (partitionsToUpdate.stream().anyMatch(partition ->
partition.startsWith(EXPRESSION_INDEX.getPartitionPath()))) {
- updateExpressionIndexIfPresent(commitMetadata, instantTime,
partitionToRecordMap);
- }
- if (partitionsToUpdate.stream().anyMatch(partition ->
partition.startsWith(SECONDARY_INDEX.getPartitionPath()))) {
- updateSecondaryIndexIfPresent(commitMetadata, partitionToRecordMap,
instantTime);
- }
- return partitionToRecordMap;
- }
- }
-
- /**
- * Update expression index from {@link HoodieCommitMetadata}.
- */
- private void updateExpressionIndexIfPresent(HoodieCommitMetadata
commitMetadata, String instantTime,
- Map<String,
HoodieData<HoodieRecord>> partitionToRecordMap) {
- if
(!MetadataPartitionType.EXPRESSION_INDEX.isMetadataPartitionAvailable(dataMetaClient))
{
- return;
+ public List<IndexPartitionAndRecords> convertMetadata() {
+ return partitionsToUpdate.stream().flatMap(indexPartition -> {
Review Comment:
This now iterates over concrete metadata partition paths, but the indexer
returned by `fromPartitionPath` is type-level. For multi-instance indexes this
duplicates updates: if table config contains two secondary or expression index
partitions, the loop invokes the same `SecondaryIndexer`/`ExpressionIndexer`
once per partition, and each invocation internally scans
`getMetadataPartitions()` and returns records for all partitions with that
prefix. The old writer path only called those type-level update helpers once
when any matching partition was present. Could we either de-duplicate by
`MetadataPartitionType` before invoking `buildUpdate`, or change the indexer
API so it builds only the requested `indexPartition`?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]