diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java index 591d86759e6e..1d99c8bd918d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ProjectedManifestEntry.java @@ -151,6 +151,7 @@ private static Projection createEntryLayoutProjection() { DataFileMeta.LEVEL, DataFileMeta.SCHEMA_ID, DataFileMeta.FIRST_ROW_ID, + DataFileMeta.MAX_SEQUENCE_NUMBER, DataFileMeta.EXTRA_FILES, DataFileMeta.EMBEDDED_FILE_INDEX, DataFileMeta.EXTERNAL_PATH))))); diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryExternalSort.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryExternalSort.java index 5c9b83772e8e..c2613890938a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryExternalSort.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryExternalSort.java @@ -24,6 +24,7 @@ import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.serializer.InternalRowSerializer; import org.apache.paimon.disk.IOManager; +import org.apache.paimon.manifest.CollectedDeletes; import org.apache.paimon.manifest.CompactFileIdentifierSet; import org.apache.paimon.manifest.FileEntry.ReusableIdentifier; import org.apache.paimon.manifest.ManifestAvroWriter; @@ -84,9 +85,10 @@ static List sortAndWriteFullEntries( ExternalSortConfig config, ManifestFile manifestFile, List newFilesForAbort, - CompactFileIdentifierSet deleteEntries, + CollectedDeletes deletes, @Nullable Integer manifestReadParallelism) throws Exception { + ReusableIdentifier identifier = new ReusableIdentifier(); try (EntrySorter sorter = new EntrySorter(sortKey, config)) { scanEntries( section, @@ -94,13 +96,15 @@ static List sortAndWriteFullEntries( manifestReadParallelism, entry -> { if (entry.isAdd() - && (deleteEntries.isEmpty() || !deleteEntries.contains(entry))) { + && (deletes.isEmpty() || !deletes.isDeleted(entry, identifier))) { sorter.write(entry); } }); List files = sorter.writeToManifest(manifestFile); newFilesForAbort.addAll(files); return files; + } finally { + identifier.release(); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java index 1b9d4f8db6a1..18cd08dfb35d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMerge.java @@ -20,31 +20,36 @@ import org.apache.paimon.data.BinaryArray; import org.apache.paimon.data.BinaryRow; -import org.apache.paimon.data.GenericRow; -import org.apache.paimon.data.InternalRow; +import org.apache.paimon.data.BinaryString; import org.apache.paimon.format.SimpleStatsCollector; -import org.apache.paimon.io.DataFileMeta; -import org.apache.paimon.manifest.CompactFileIdentifierSet; -import org.apache.paimon.manifest.DeletedRowIdSet; -import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.CollectedDeletes; +import org.apache.paimon.manifest.FileEntry.ReusableIdentifier; import org.apache.paimon.manifest.ManifestAvroReader; import org.apache.paimon.manifest.ManifestAvroReader.RawBlock; import org.apache.paimon.manifest.ManifestAvroReader.RowIterator; import org.apache.paimon.manifest.ManifestAvroWriter.EncodedBlockMeta; -import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.manifest.ManifestFile; import org.apache.paimon.manifest.ManifestFileMeta; +import org.apache.paimon.manifest.PartitionDictionary; +import org.apache.paimon.manifest.ProjectedManifestEntry; +import org.apache.paimon.memory.MemorySegment; +import org.apache.paimon.memory.MemorySegmentUtils; import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.stats.SimpleStatsConverter; -import org.apache.paimon.types.DataField; import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.ByteArrayKey; +import org.apache.paimon.utils.ByteArrayLookupKey; import org.apache.paimon.utils.Pair; import javax.annotation.Nullable; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; +import java.util.Comparator; import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.function.Function; import static org.apache.paimon.utils.ManifestReadThreadPool.sequentialBatchedExecute; @@ -57,91 +62,9 @@ final class ManifestEntryRunMerge { private static final long MAX_IN_MEMORY_FRAGMENTED_ENTRIES = 25_000L; private static final int MAX_STREAM_CURSORS = 128; private static final int MAX_STREAM_READ_AMPLIFICATION = 8; - static final int KIND = 0; - static final int PARTITION = 1; - static final int BUCKET = 2; - static final int FILE = 3; - static final int FILE_NAME = 0; - static final int ROW_COUNT = 1; - static final int LEVEL = 2; - static final int SCHEMA_ID = 3; - static final int FIRST_ROW_ID = 4; - static final int MAX_SEQUENCE_NUMBER = 5; - static final int EXTRA_FILES = 6; - static final int EMBEDDED_FILE_INDEX = 7; - static final int EXTERNAL_PATH = 8; - static final int FILE_FIELD_COUNT = 9; - private static final String[] ENTRY_FILE_FIELD_NAMES = { - DataFileMeta.FILE_NAME, - DataFileMeta.ROW_COUNT, - DataFileMeta.LEVEL, - DataFileMeta.SCHEMA_ID, - DataFileMeta.FIRST_ROW_ID, - DataFileMeta.MAX_SEQUENCE_NUMBER, - DataFileMeta.EXTRA_FILES, - DataFileMeta.EMBEDDED_FILE_INDEX, - DataFileMeta.EXTERNAL_PATH - }; - private static final InternalRow.FieldGetter[] ENTRY_FILE_GETTERS = entryFileGetters(); - static final RowType ENTRY_LAYOUT = entryLayout(); - private static final int FULL_KIND = - ManifestEntry.MANIFEST_ROW_TYPE.getFieldIndex(ManifestEntry.KIND); - private static final int FULL_PARTITION = - ManifestEntry.MANIFEST_ROW_TYPE.getFieldIndex(ManifestEntry.PARTITION); - private static final int FULL_BUCKET = - ManifestEntry.MANIFEST_ROW_TYPE.getFieldIndex(ManifestEntry.BUCKET); - private static final int FULL_FILE = - ManifestEntry.MANIFEST_ROW_TYPE.getFieldIndex(ManifestEntry.FILE); private ManifestEntryRunMerge() {} - private static InternalRow.FieldGetter[] entryFileGetters() { - InternalRow.FieldGetter[] getters = - new InternalRow.FieldGetter[ENTRY_FILE_FIELD_NAMES.length]; - for (int field = 0; field < getters.length; field++) { - int position = DataFileMeta.SCHEMA.getFieldIndex(ENTRY_FILE_FIELD_NAMES[field]); - getters[field] = - InternalRow.createFieldGetter( - DataFileMeta.SCHEMA.getTypeAt(position), position); - } - return getters; - } - - private static RowType entryLayout() { - List fields = new ArrayList<>(); - fields.add(ManifestEntry.MANIFEST_ROW_TYPE.getField(ManifestEntry.KIND)); - fields.add(ManifestEntry.MANIFEST_ROW_TYPE.getField(ManifestEntry.PARTITION)); - fields.add(ManifestEntry.MANIFEST_ROW_TYPE.getField(ManifestEntry.BUCKET)); - fields.add( - ManifestEntry.MANIFEST_ROW_TYPE - .getField(ManifestEntry.FILE) - .newType( - DataFileMeta.SCHEMA.project( - DataFileMeta.FILE_NAME, - DataFileMeta.ROW_COUNT, - DataFileMeta.LEVEL, - DataFileMeta.SCHEMA_ID, - DataFileMeta.FIRST_ROW_ID, - DataFileMeta.MAX_SEQUENCE_NUMBER, - DataFileMeta.EXTRA_FILES, - DataFileMeta.EMBEDDED_FILE_INDEX, - DataFileMeta.EXTERNAL_PATH))); - return new RowType(false, fields); - } - - static GenericRow projectEntryLayout( - GenericRow fullRow, GenericRow reuse, GenericRow reuseFile) { - reuse.setField(KIND, fullRow.getByte(FULL_KIND)); - reuse.setField(PARTITION, fullRow.getBinary(FULL_PARTITION)); - reuse.setField(BUCKET, fullRow.getInt(FULL_BUCKET)); - InternalRow fullFile = fullRow.getRow(FULL_FILE, DataFileMeta.SCHEMA.getFieldCount()); - for (int field = 0; field < ENTRY_FILE_GETTERS.length; field++) { - reuseFile.setField(field, ENTRY_FILE_GETTERS[field].getFieldOrNull(fullFile)); - } - reuse.setField(FILE, reuseFile); - return reuse; - } - /** * Returns null when the input is too fragmented for a bounded streaming merge. The caller must * fall back to the spillable external sorter in that case. @@ -153,26 +76,24 @@ static List sortAndWriteFullEntries( RowType partitionType, ManifestFile manifestFile, List newFilesForAbort, - CompactFileIdentifierSet deletedIdentifiers, - DeletedRowIdSet deletedRowIds, + CollectedDeletes deletes, int maxNumFileHandles, @Nullable Integer manifestReadParallelism) throws Exception { - ManifestEntryRunMergeEntry.Filter filter = - new ManifestEntryRunMergeEntry.Filter(deletedIdentifiers, deletedRowIds, true); ManifestEntryRunMergePlan plan = discoverRuns( section, sortKey, partitionType, manifestFile, - filter, + deletes, + false, maxNumFileHandles, manifestReadParallelism); if (plan == null) { return null; } - return plan.mergeToManifest(sortKey, manifestFile, filter, newFilesForAbort); + return plan.mergeToManifest(sortKey, manifestFile, newFilesForAbort); } /** @@ -189,38 +110,24 @@ static Pair, List> sortAndWriteMinorEnt int maxNumFileHandles, @Nullable Integer manifestReadParallelism) throws Exception { - CompactFileIdentifierSet deletedIdentifiers = new CompactFileIdentifierSet(); - DeletedRowIdSet deletedRowIds = new DeletedRowIdSet(); - ManifestEntryRunMergeEntry.Filter.Minor filter = - new ManifestEntryRunMergeEntry.Filter.Minor( - deletedIdentifiers, deletedRowIds, true); + CollectedDeletes deletes = new CollectedDeletes(true); try { - ManifestEntryRunMergePlan plan; - try { - plan = - discoverRuns( - section, - sortKey, - partitionType, - manifestFile, - filter, - maxNumFileHandles, - manifestReadParallelism); - } finally { - deletedRowIds.releaseRangeIndex(); - } + ManifestEntryRunMergePlan plan = + discoverRuns( + section, + sortKey, + partitionType, + manifestFile, + deletes, + true, + maxNumFileHandles, + manifestReadParallelism); if (plan == null) { return null; } - return plan.mergeMinorToManifest( - sortKey, - manifestFile, - filter, - deletedIdentifiers, - deletedRowIds, - newFilesForAbort); + return plan.mergeMinorToManifest(sortKey, manifestFile, newFilesForAbort); } finally { - deletedIdentifiers.release(); + deletes.release(); } } @@ -230,12 +137,13 @@ private static ManifestEntryRunMergePlan discoverRuns( ManifestFileSorter.RowIdEntrySortKey sortKey, RowType partitionType, ManifestFile manifestFile, - ManifestEntryRunMergeEntry.Filter filter, + CollectedDeletes deletes, + boolean minor, int maxNumFileHandles, @Nullable Integer manifestReadParallelism) throws Exception { - ManifestEntryRunMergeEntry.PartitionDictionary partitions = - new ManifestEntryRunMergeEntry.PartitionDictionary(sortKey); + SortPartitionDictionary partitions = + new SortPartitionDictionary(sortKey::comparePartitions); List sources = new ArrayList<>(); int streamCursorCount = 0; long inMemoryEntries = 0; @@ -243,36 +151,88 @@ private static ManifestEntryRunMergePlan discoverRuns( if (section.size() <= 1 || (manifestReadParallelism != null && manifestReadParallelism <= 1)) { for (ManifestFileMeta meta : section) { - Discovery.DiscoveredManifest manifest = - discoverManifestRuns(meta, manifestFile, partitionType, partitions, filter); + CollectedDeletes discoveryDeletes = + minor ? new CollectedDeletes(deletes.useRowIdFilter()) : deletes; + Discovery.DiscoveredManifest manifest; + try { + manifest = + discoverManifestRuns( + meta, + manifestFile, + partitionType, + partitions, + discoveryDeletes, + minor); + } catch (Exception e) { + if (minor) { + discoveryDeletes.release(); + } + throw e; + } + if (minor) { + try { + deletes.combine(discoveryDeletes); + } finally { + discoveryDeletes.release(); + } + } if (manifest.requiresExternalSort) { return null; } discovered.add(manifest); } } else { - Function> reader = - meta -> { - try { - return Collections.singletonList( - discoverManifestRuns( - meta, manifestFile, partitionType, partitions, filter)); - } catch (Exception e) { - throw new RuntimeException( - "Failed to discover sorted Avro runs in " + meta.fileName(), e); - } - }; - for (Discovery.DiscoveredManifest manifest : + boolean requiresExternalSort = false; + Function>> + reader = + meta -> { + CollectedDeletes discoveryDeletes = + minor + ? new CollectedDeletes(deletes.useRowIdFilter()) + : deletes; + try { + return Collections.singletonList( + Pair.of( + discoverManifestRuns( + meta, + manifestFile, + partitionType, + partitions, + discoveryDeletes, + minor), + discoveryDeletes)); + } catch (Exception e) { + if (minor) { + discoveryDeletes.release(); + } + throw new RuntimeException( + "Failed to discover sorted Avro runs in " + + meta.fileName(), + e); + } + }; + for (Pair scan : sequentialBatchedExecute(reader, section, manifestReadParallelism)) { - discovered.add(manifest); + requiresExternalSort |= scan.getLeft().requiresExternalSort; + if (minor) { + try { + deletes.combine(scan.getRight()); + } finally { + scan.getRight().release(); + } + } + discovered.add(scan.getLeft()); + } + // Drain every task in the bounded discovery batch before falling back. Returning from + // the lazy iterator early would leave already submitted manifest scans running beside + // the external sorter and duplicate their I/O and retained memory. + if (requiresExternalSort) { + return null; } } for (int manifestIndex = 0; manifestIndex < section.size(); manifestIndex++) { ManifestFileMeta meta = section.get(manifestIndex); Discovery.DiscoveredManifest manifest = discovered.get(manifestIndex); - if (manifest.requiresExternalSort) { - return null; - } if (manifest.fragmented) { long entryCount = meta.numAddedFiles() + meta.numDeletedFiles(); inMemoryEntries += entryCount; @@ -289,26 +249,34 @@ private static ManifestEntryRunMergePlan discoverRuns( return null; } } + if (minor) { + deletes.toImmutable(); + } partitions.finish(); for (Discovery.DiscoveredManifest manifest : discovered) { - manifest.finishFiltering(filter); + manifest.finishFiltering(deletes, minor); manifest.updatePartitionRanks(partitions); } - return new ManifestEntryRunMergePlan(sources, partitions); + return new ManifestEntryRunMergePlan(sources, partitions, deletes, minor); } private static Discovery.DiscoveredManifest discoverManifestRuns( ManifestFileMeta meta, ManifestFile manifestFile, RowType partitionType, - ManifestEntryRunMergeEntry.PartitionDictionary partitions, - ManifestEntryRunMergeEntry.Filter filter) + SortPartitionDictionary partitions, + CollectedDeletes deletes, + boolean minor) throws Exception { + ReusableIdentifier identifier = new ReusableIdentifier(); try (ManifestAvroReader reader = manifestFile.scanAvroBlocks(meta.fileName(), meta.fileSize())) { - return discoverManifestRuns(meta, reader, partitionType, partitions, filter); + return discoverManifestRuns( + meta, reader, partitionType, partitions, deletes, minor, identifier); } catch (UnsupportedOperationException unsupported) { return Discovery.DiscoveredManifest.requiresExternalSort(); + } finally { + identifier.release(); } } @@ -316,26 +284,32 @@ private static Discovery.DiscoveredManifest discoverManifestRuns( ManifestFileMeta meta, ManifestAvroReader reader, RowType partitionType, - ManifestEntryRunMergeEntry.PartitionDictionary partitions, - ManifestEntryRunMergeEntry.Filter filter) + SortPartitionDictionary partitions, + CollectedDeletes deletes, + boolean minor, + ReusableIdentifier identifier) throws Exception { SimpleStatsConverter partitionStatsConverter = new SimpleStatsConverter(partitionType); List runs = new ArrayList<>(); List blocks = new ArrayList<>(); - ManifestEntryRunMergeEntry.Key previous = new ManifestEntryRunMergeEntry.Key(); - ManifestEntryRunMergeEntry.Key current = new ManifestEntryRunMergeEntry.Key(); + SortKey previous = new SortKey(); + SortKey current = new SortKey(); boolean hasPrevious = false; long runStart = 0; long position = 0; long entryCount = meta.numAddedFiles() + meta.numDeletedFiles(); boolean fragmented = false; + ProjectedManifestEntry entry = ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION.createEntry(); while (reader.hasNext()) { RawBlock rawBlock = reader.next(); - RowIterator rows = rawBlock.toRows(ENTRY_LAYOUT); + RowIterator rows = + rawBlock.toRows(ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION.projectedType()); while (rows.hasNext()) { - GenericRow row = rows.next(); - current.replace(row, partitions); - filter.observe(row, current); + entry.replace(rows.next()); + current.replace(entry, partitions); + if (minor && entry.isDelete()) { + deletes.add(entry, deletes.useRowIdFilter(), false); + } if (fragmented) { position++; continue; @@ -350,7 +324,7 @@ private static Discovery.DiscoveredManifest discoverManifestRuns( partitionType)); } Discovery.BlockInfo block = blocks.get(blocks.size() - 1); - block.collectForSort(row, current, partitions, filter); + block.collectForSort(entry, current, partitions, deletes, identifier, minor); boolean inversion = hasPrevious && compareDiscoveryKeys(previous, current, partitions) > 0; if (inversion) { @@ -374,7 +348,7 @@ private static Discovery.DiscoveredManifest discoverManifestRuns( } position++; if (rows.recordIndex() + 1 == rawBlock.recordCount()) { - ManifestEntryRunMergeEntry.Key stableLastKey = current.stableCopy(); + SortKey stableLastKey = current.stableCopy(); block.finishSort(position, stableLastKey, partitionStatsConverter); previous.copyFrom(stableLastKey); } else { @@ -413,23 +387,17 @@ private static boolean exceedsStreamingReadAmplification( } private static int compareDiscoveryKeys( - ManifestEntryRunMergeEntry.Key left, - ManifestEntryRunMergeEntry.Key right, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) { + SortKey left, SortKey right, SortPartitionDictionary partitions) { return compareRemainingKeys( left, right, partitions.compareIds(left.partitionId, right.partitionId)); } - static int compareMergeKeys( - ManifestEntryRunMergeEntry.Key left, ManifestEntryRunMergeEntry.Key right) { + static int compareMergeKeys(SortKey left, SortKey right) { return compareRemainingKeys( left, right, Integer.compare(left.partitionRank, right.partitionRank)); } - private static int compareRemainingKeys( - ManifestEntryRunMergeEntry.Key left, - ManifestEntryRunMergeEntry.Key right, - int comparison) { + private static int compareRemainingKeys(SortKey left, SortKey right, int comparison) { if (comparison == 0) { comparison = Byte.compare(left.kind, right.kind); } @@ -437,10 +405,10 @@ private static int compareRemainingKeys( comparison = Long.compare(left.firstRowId, right.firstRowId); } if (comparison == 0) { - comparison = Long.compare(left.rangeEnd, right.rangeEnd); + comparison = Long.compare(left.lastRowId, right.lastRowId); } if (comparison == 0) { - comparison = Long.compare(left.reverseSequence, right.reverseSequence); + comparison = Long.compare(left.descendingSequenceKey, right.descendingSequenceKey); } if (comparison == 0) { comparison = compareBytes(left, right); @@ -448,12 +416,15 @@ private static int compareRemainingKeys( return comparison; } - private static int compareBytes( - ManifestEntryRunMergeEntry.Key left, ManifestEntryRunMergeEntry.Key right) { + private static int compareBytes(SortKey left, SortKey right) { int minLength = Math.min(left.fileNameLength, right.fileNameLength); for (int i = 0; i < minLength; i++) { - int leftByte = left.fileNameBytes[left.fileNameOffset + i] & 0xFF; - int rightByte = right.fileNameBytes[right.fileNameOffset + i] & 0xFF; + int leftByte = + MemorySegmentUtils.getByte(left.fileNameSegments, left.fileNameOffset + i) + & 0xFF; + int rightByte = + MemorySegmentUtils.getByte(right.fileNameSegments, right.fileNameOffset + i) + & 0xFF; if (leftByte != rightByte) { return leftByte - rightByte; } @@ -461,6 +432,80 @@ private static int compareBytes( return left.fileNameLength - right.fileNameLength; } + static final class SortKey { + + int partitionId; + int partitionRank; + byte kind; + long firstRowId; + long lastRowId; + long descendingSequenceKey; + MemorySegment[] fileNameSegments; + int fileNameOffset; + int fileNameLength; + byte[] ownedFileNameBytes; + MemorySegment[] ownedFileNameSegments; + + static SortKey viewOf(ProjectedManifestEntry entry, SortPartitionDictionary partitions) { + SortKey key = new SortKey(); + key.replace(entry, partitions); + return key; + } + + void replace(ProjectedManifestEntry entry, SortPartitionDictionary partitions) { + long firstRowId = entry.file().nonNullFirstRowId(); + this.partitionId = partitions.id(entry.partitionBytes()); + this.partitionRank = partitions.rank(partitionId); + this.kind = entry.kind().toByteValue(); + this.firstRowId = firstRowId; + this.lastRowId = firstRowId + entry.file().rowCount() - 1L; + this.descendingSequenceKey = Long.MAX_VALUE - entry.file().maxSequenceNumber(); + BinaryString fileName = entry.file().fileNameBinary(); + this.fileNameSegments = fileName.getSegments(); + this.fileNameOffset = fileName.getOffset(); + this.fileNameLength = fileName.getSizeInBytes(); + } + + void copyFrom(SortKey key) { + this.partitionId = key.partitionId; + this.partitionRank = key.partitionRank; + this.kind = key.kind; + this.firstRowId = key.firstRowId; + this.lastRowId = key.lastRowId; + this.descendingSequenceKey = key.descendingSequenceKey; + ensureFileNameCapacity(key.fileNameLength); + MemorySegmentUtils.copyToBytes( + key.fileNameSegments, + key.fileNameOffset, + ownedFileNameBytes, + 0, + key.fileNameLength); + this.fileNameSegments = ownedFileNameSegments; + this.fileNameOffset = 0; + this.fileNameLength = key.fileNameLength; + } + + SortKey stableCopy() { + SortKey copy = new SortKey(); + copy.copyFrom(this); + return copy; + } + + private void ensureFileNameCapacity(int length) { + if (ownedFileNameBytes == null || ownedFileNameBytes.length < length) { + ownedFileNameBytes = new byte[length]; + ownedFileNameSegments = + new MemorySegment[] {MemorySegment.wrap(ownedFileNameBytes)}; + } + } + + void clear() { + fileNameSegments = null; + ownedFileNameBytes = null; + ownedFileNameSegments = null; + } + } + /** Results and Avro block metadata collected while discovering natural manifest runs. */ static final class Discovery { @@ -500,15 +545,15 @@ static DiscoveredManifest requiresExternalSort() { Collections.emptyList(), Collections.emptyList(), false, true); } - void updatePartitionRanks(ManifestEntryRunMergeEntry.PartitionDictionary partitions) { + void updatePartitionRanks(SortPartitionDictionary partitions) { for (BlockInfo block : blocks) { block.updatePartitionRanks(partitions); } } - void finishFiltering(ManifestEntryRunMergeEntry.Filter filter) { + void finishFiltering(CollectedDeletes deletes, boolean minor) { for (BlockInfo block : blocks) { - block.finishFiltering(filter); + block.finishFiltering(deletes, minor); } } } @@ -517,159 +562,240 @@ static final class BlockInfo { final long ordinal; final long start; - final @Nullable ManifestEntryRunMergeEntry.Key firstKey; - boolean eligible; + final @Nullable SortKey firstKey; boolean sorted = true; long end; - @Nullable ManifestEntryRunMergeEntry.Key lastKey; - long addedFiles; - long deletedFiles; - long schemaId = Long.MIN_VALUE; - int minBucket = Integer.MAX_VALUE; - int maxBucket = Integer.MIN_VALUE; - int minLevel = Integer.MAX_VALUE; - int maxLevel = Integer.MIN_VALUE; - long minRowId = Long.MAX_VALUE; - long maxRowId = Long.MIN_VALUE; - final boolean singleFieldSortedPartitionStats; - @Nullable SimpleStatsCollector partitionStats; - final RowType partitionType; - @Nullable BinaryRow nullPartition; - long nullPartitionCount; - @Nullable BinaryRow minNonNullPartition; - @Nullable BinaryRow maxNonNullPartition; - EncodedBlockMeta metadata; + @Nullable SortKey lastKey; + long minRowId; + long maxRowId; + @Nullable BlockMetadataAccumulator metadataAccumulator; + @Nullable EncodedBlockMeta metadata; BlockInfo( long ordinal, long start, - boolean eligible, - ManifestEntryRunMergeEntry.Key firstKey, + boolean rawBlockCopySupported, + SortKey firstKey, RowType partitionType) { this.ordinal = ordinal; this.start = start; - this.eligible = eligible; this.firstKey = firstKey; - this.partitionType = partitionType; - this.singleFieldSortedPartitionStats = - eligible && firstKey != null && partitionType.getFieldCount() == 1; - this.partitionStats = - eligible && firstKey != null && !singleFieldSortedPartitionStats - ? new SimpleStatsCollector(partitionType) + this.metadataAccumulator = + rawBlockCopySupported && firstKey != null + ? new BlockMetadataAccumulator(partitionType) : null; } void collectForSort( - GenericRow record, - ManifestEntryRunMergeEntry.Key key, - ManifestEntryRunMergeEntry.PartitionDictionary partitions, - ManifestEntryRunMergeEntry.Filter filter) { - if (!eligible) { + ProjectedManifestEntry entry, + SortKey key, + SortPartitionDictionary partitions, + CollectedDeletes deletes, + ReusableIdentifier identifier, + boolean minor) { + if (metadataAccumulator == null) { return; } - if (!filter.copyable(record, key)) { - eligible = false; - releasePartitionStats(); + if (!deletes.copyable(entry, identifier, minor)) { + metadataAccumulator = null; return; } - collectEntryStats(record, key); - BinaryRow partition = partitions.partition(key.partitionId); - if (singleFieldSortedPartitionStats) { - if (partition.isNullAt(0)) { - nullPartition = partition; - nullPartitionCount++; - } else { - if (minNonNullPartition == null) { - minNonNullPartition = partition; - } - maxNonNullPartition = partition; - } - } else { - partitionStats.collect(partition); + metadataAccumulator.collect(entry, key, partitions.partition(key.partitionId)); + } + + void finishSort( + long end, SortKey lastKey, SimpleStatsConverter partitionStatsConverter) { + this.end = end; + this.lastKey = lastKey; + if (metadataAccumulator != null && sorted) { + minRowId = metadataAccumulator.minRowId; + maxRowId = metadataAccumulator.maxRowId; + metadata = metadataAccumulator.finish(partitionStatsConverter); + } + metadataAccumulator = null; + } + + boolean copyable(long runStart, long runEnd) { + return metadata != null && start >= runStart && end <= runEnd; + } + + void finishFiltering(CollectedDeletes deletes, boolean minor) { + if (metadata != null + && minor + && (!deletes.useRowIdFilter() + || deletes.intersectsRowIds(minRowId, maxRowId))) { + metadata = null; } } - private void collectEntryStats(GenericRow record, ManifestEntryRunMergeEntry.Key key) { - InternalRow file = ManifestEntryRunMergeEntry.file(record); - if (key.kind == FileKind.ADD.toByteValue()) { + void updatePartitionRanks(SortPartitionDictionary partitions) { + checkState(firstKey != null && lastKey != null, "Manifest block has no sort keys."); + firstKey.partitionRank = partitions.rank(firstKey.partitionId); + lastKey.partitionRank = partitions.rank(lastKey.partitionId); + } + } + + /** Mutable statistics retained only while the current Avro block is being inspected. */ + private static final class BlockMetadataAccumulator { + + private long addedFiles; + private long deletedFiles; + private long schemaId = Long.MIN_VALUE; + private int minBucket = Integer.MAX_VALUE; + private int maxBucket = Integer.MIN_VALUE; + private int minLevel = Integer.MAX_VALUE; + private int maxLevel = Integer.MIN_VALUE; + private long minRowId = Long.MAX_VALUE; + private long maxRowId = Long.MIN_VALUE; + private final BlockPartitionStats partitionStats; + + private BlockMetadataAccumulator(RowType partitionType) { + this.partitionStats = new BlockPartitionStats(partitionType); + } + + private void collect(ProjectedManifestEntry entry, SortKey key, BinaryRow partition) { + if (entry.isAdd()) { addedFiles++; } else { deletedFiles++; } - schemaId = Math.max(schemaId, file.getLong(SCHEMA_ID)); - int bucket = record.getInt(BUCKET); + schemaId = Math.max(schemaId, entry.file().schemaId()); + int bucket = entry.bucket(); minBucket = Math.min(minBucket, bucket); maxBucket = Math.max(maxBucket, bucket); - int level = file.getInt(LEVEL); + int level = entry.file().level(); minLevel = Math.min(minLevel, level); maxLevel = Math.max(maxLevel, level); minRowId = Math.min(minRowId, key.firstRowId); - maxRowId = Math.max(maxRowId, key.rangeEnd); + maxRowId = Math.max(maxRowId, key.lastRowId); + partitionStats.collect(partition); } - void finishSort( - long end, - ManifestEntryRunMergeEntry.Key lastKey, - SimpleStatsConverter partitionStatsConverter) { - this.end = end; - this.lastKey = lastKey; - if (eligible && sorted) { - SimpleStats encodedPartitionStats; - if (singleFieldSortedPartitionStats) { - BinaryRow min = - minNonNullPartition == null ? nullPartition : minNonNullPartition; - BinaryRow max = - maxNonNullPartition == null ? nullPartition : maxNonNullPartition; - checkState(min != null && max != null, "Manifest block has no partition."); - encodedPartitionStats = - new SimpleStats( - min, - max, - BinaryArray.fromLongArray(new Long[] {nullPartitionCount})); - } else { - checkState( - partitionStats != null, "Manifest block has no partition stats."); - encodedPartitionStats = - partitionStatsConverter.toBinaryAllMode(partitionStats.extract()); + private EncodedBlockMeta finish(SimpleStatsConverter partitionStatsConverter) { + return new EncodedBlockMeta( + addedFiles, + deletedFiles, + schemaId, + minBucket, + maxBucket, + minLevel, + maxLevel, + minRowId, + maxRowId, + partitionStats.finish(partitionStatsConverter)); + } + } + + /** Partition statistics for one sorted Avro block. */ + private static final class BlockPartitionStats { + + private final boolean singleField; + private final @Nullable SimpleStatsCollector collector; + private @Nullable BinaryRow nullPartition; + private @Nullable BinaryRow minNonNullPartition; + private @Nullable BinaryRow maxNonNullPartition; + private long nullCount; + + private BlockPartitionStats(RowType partitionType) { + this.singleField = partitionType.getFieldCount() == 1; + this.collector = singleField ? null : new SimpleStatsCollector(partitionType); + } + + private void collect(BinaryRow partition) { + if (!singleField) { + checkState(collector != null, "Manifest block has no partition collector."); + collector.collect(partition); + return; + } + if (partition.isNullAt(0)) { + nullPartition = partition; + nullCount++; + } else { + if (minNonNullPartition == null) { + minNonNullPartition = partition; } - metadata = - new EncodedBlockMeta( - addedFiles, - deletedFiles, - schemaId, - minBucket, - maxBucket, - minLevel, - maxLevel, - minRowId, - maxRowId, - encodedPartitionStats); + maxNonNullPartition = partition; } - releasePartitionStats(); } - private void releasePartitionStats() { - partitionStats = null; - nullPartition = null; - minNonNullPartition = null; - maxNonNullPartition = null; + private SimpleStats finish(SimpleStatsConverter converter) { + if (!singleField) { + checkState(collector != null, "Manifest block has no partition collector."); + return converter.toBinaryAllMode(collector.extract()); + } + BinaryRow min = minNonNullPartition == null ? nullPartition : minNonNullPartition; + BinaryRow max = maxNonNullPartition == null ? nullPartition : maxNonNullPartition; + checkState(min != null && max != null, "Manifest block has no partition."); + return new SimpleStats(min, max, BinaryArray.fromLongArray(new Long[] {nullCount})); } + } + } - boolean copyable(long runStart, long runEnd) { - return metadata != null && start >= runStart && end <= runEnd; - } + /** Concurrent partition dictionary and ordering used only by manifest run merge. */ + static final class SortPartitionDictionary { - void finishFiltering(ManifestEntryRunMergeEntry.Filter filter) { - if (metadata != null && !filter.copyableAfterDiscovery(minRowId, maxRowId)) { - metadata = null; + private final Comparator comparator; + private final PartitionDictionary partitions = new PartitionDictionary(); + private final Map ids = new ConcurrentHashMap<>(); + private final ThreadLocal lookup = + ThreadLocal.withInitial(ByteArrayLookupKey::new); + private int[] ranks; + + SortPartitionDictionary(Comparator comparator) { + this.comparator = comparator; + } + + int id(byte[] bytes) { + ByteArrayLookupKey lookupKey = lookup.get(); + lookupKey.reset(bytes); + try { + Integer existing = ids.get(lookupKey); + if (existing != null) { + return existing; } + synchronized (this) { + existing = ids.get(lookupKey); + if (existing != null) { + return existing; + } + checkState(ranks == null, "Manifest scan found an unknown partition."); + byte[] canonical = Arrays.copyOf(bytes, bytes.length); + int id = partitions.id(canonical); + ids.put(new ByteArrayKey(canonical), id); + return id; + } + } finally { + lookupKey.clear(); } + } - void updatePartitionRanks(ManifestEntryRunMergeEntry.PartitionDictionary partitions) { - checkState(firstKey != null && lastKey != null, "Manifest block has no sort keys."); - firstKey.partitionRank = partitions.rank(firstKey.partitionId); - lastKey.partitionRank = partitions.rank(lastKey.partitionId); + void finish() { + int partitionCount = ids.size(); + List order = new ArrayList<>(partitionCount); + for (int id = 0; id < partitionCount; id++) { + order.add(id); + } + order.sort((left, right) -> compareIds(left, right)); + ranks = new int[partitionCount]; + int rank = 0; + for (int position = 0; position < order.size(); position++) { + if (position > 0 && compareIds(order.get(position - 1), order.get(position)) != 0) { + rank++; + } + ranks[order.get(position)] = rank; } } + + int compareIds(int left, int right) { + return comparator.compare(partitions.partition(left), partitions.partition(right)); + } + + int rank(int id) { + return ranks == null ? 0 : ranks[id]; + } + + BinaryRow partition(int id) { + return partitions.partition(id); + } } } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergeEntry.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergeEntry.java deleted file mode 100644 index 5e8c318e9fda..000000000000 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergeEntry.java +++ /dev/null @@ -1,339 +0,0 @@ -/* - * 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.paimon.operation; - -import org.apache.paimon.data.BinaryRow; -import org.apache.paimon.data.BinaryString; -import org.apache.paimon.data.GenericRow; -import org.apache.paimon.data.InternalRow; -import org.apache.paimon.manifest.CompactFileIdentifierSet; -import org.apache.paimon.manifest.DeletedRowIdSet; -import org.apache.paimon.manifest.FileEntry.ReusableIdentifier; -import org.apache.paimon.manifest.FileKind; -import org.apache.paimon.manifest.ProjectedManifestEntry; -import org.apache.paimon.memory.MemorySegmentUtils; -import org.apache.paimon.utils.ByteArrayKey; -import org.apache.paimon.utils.ByteArrayLookupKey; -import org.apache.paimon.utils.SerializationUtils; - -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; - -import static org.apache.paimon.utils.Preconditions.checkState; - -/** Entry-level state shared by manifest run discovery and merge execution. */ -final class ManifestEntryRunMergeEntry { - - private ManifestEntryRunMergeEntry() {} - - static final class Key { - - int partitionId; - int partitionRank; - byte kind; - boolean hasRowId; - long firstRowId; - long rangeEnd; - long reverseSequence; - byte[] fileNameBytes; - int fileNameOffset; - int fileNameLength; - - static Key viewOf(ProjectedManifestEntry entry, PartitionDictionary partitions) { - Key key = new Key(); - key.replace(entry, partitions); - return key; - } - - void replace(ProjectedManifestEntry entry, PartitionDictionary partitions) { - long firstRowId = entry.file().nonNullFirstRowId(); - this.partitionId = partitions.id(entry.partitionBytes()); - this.partitionRank = partitions.rank(partitionId); - this.kind = entry.kind().toByteValue(); - this.hasRowId = true; - this.firstRowId = firstRowId; - this.rangeEnd = firstRowId + entry.file().rowCount() - 1L; - this.reverseSequence = Long.MAX_VALUE - entry.file().maxSequenceNumber(); - this.fileNameBytes = entry.file().fileNameBinary().toBytes(); - this.fileNameOffset = 0; - this.fileNameLength = fileNameBytes.length; - } - - void replace(GenericRow record, PartitionDictionary partitions) { - InternalRow file = file(record); - checkState( - !file.isNullAt(ManifestEntryRunMerge.FIRST_ROW_ID), - "First row id should not be null."); - this.partitionId = partitions.id(record.getBinary(ManifestEntryRunMerge.PARTITION)); - this.partitionRank = partitions.rank(partitionId); - this.kind = record.getByte(ManifestEntryRunMerge.KIND); - this.hasRowId = true; - this.firstRowId = file.getLong(ManifestEntryRunMerge.FIRST_ROW_ID); - this.rangeEnd = firstRowId + file.getLong(ManifestEntryRunMerge.ROW_COUNT) - 1L; - this.reverseSequence = - Long.MAX_VALUE - file.getLong(ManifestEntryRunMerge.MAX_SEQUENCE_NUMBER); - BinaryString fileName = file.getString(ManifestEntryRunMerge.FILE_NAME); - this.fileNameBytes = - MemorySegmentUtils.copyToBytes( - fileName.getSegments(), - fileName.getOffset(), - fileName.getSizeInBytes()); - this.fileNameOffset = 0; - this.fileNameLength = fileNameBytes.length; - } - - void replaceForCompaction(GenericRow record) { - InternalRow file = file(record); - this.kind = record.getByte(ManifestEntryRunMerge.KIND); - this.hasRowId = !file.isNullAt(ManifestEntryRunMerge.FIRST_ROW_ID); - if (hasRowId) { - this.firstRowId = file.getLong(ManifestEntryRunMerge.FIRST_ROW_ID); - this.rangeEnd = firstRowId + file.getLong(ManifestEntryRunMerge.ROW_COUNT) - 1L; - } - } - - void copyFrom(Key key) { - this.partitionId = key.partitionId; - this.partitionRank = key.partitionRank; - this.kind = key.kind; - this.hasRowId = key.hasRowId; - this.firstRowId = key.firstRowId; - this.rangeEnd = key.rangeEnd; - this.reverseSequence = key.reverseSequence; - this.fileNameBytes = key.fileNameBytes; - this.fileNameOffset = key.fileNameOffset; - this.fileNameLength = key.fileNameLength; - } - - Key stableCopy() { - Key copy = new Key(); - copy.copyFrom(this); - copy.fileNameBytes = - Arrays.copyOfRange( - fileNameBytes, fileNameOffset, fileNameOffset + fileNameLength); - copy.fileNameOffset = 0; - return copy; - } - - void clear() { - fileNameBytes = null; - } - } - - /** Interns variable-width partition bytes once and assigns comparator-compatible ranks. */ - static final class PartitionDictionary { - - final ManifestFileSorter.RowIdEntrySortKey sortKey; - final Map ids = new ConcurrentHashMap<>(); - final ThreadLocal lookup = - ThreadLocal.withInitial(ByteArrayLookupKey::new); - volatile BinaryRow[] partitions = new BinaryRow[16]; - int partitionCount; - int[] ranks; - - PartitionDictionary(ManifestFileSorter.RowIdEntrySortKey sortKey) { - this.sortKey = sortKey; - } - - PartitionDictionary() { - this.sortKey = null; - } - - int id(byte[] bytes) { - return id(bytes, 0, bytes.length); - } - - int id(byte[] bytes, int offset, int length) { - ByteArrayLookupKey lookupKey = lookup.get(); - lookupKey.reset(bytes, offset, length); - try { - Integer existing = ids.get(lookupKey); - if (existing != null) { - return existing; - } - synchronized (this) { - existing = ids.get(lookupKey); - if (existing != null) { - return existing; - } - checkState(ranks == null, "Full manifest scan found an unknown partition."); - byte[] canonical = Arrays.copyOfRange(bytes, offset, offset + length); - int id = partitionCount; - if (id == partitions.length) { - partitions = Arrays.copyOf(partitions, partitions.length << 1); - } - partitions[id] = SerializationUtils.deserializeBinaryRow(canonical); - ids.put(new ByteArrayKey(canonical), id); - partitionCount = id + 1; - return id; - } - } finally { - lookupKey.clear(); - } - } - - int compareIds(int left, int right) { - checkState(sortKey != null, "Partition dictionary has no sort key."); - return sortKey.comparePartitions(partitions[left], partitions[right]); - } - - void finish() { - List order = new ArrayList<>(partitionCount); - for (int id = 0; id < partitionCount; id++) { - order.add(id); - } - order.sort((left, right) -> compareIds(left, right)); - ranks = new int[partitionCount]; - int rank = 0; - for (int position = 0; position < order.size(); position++) { - if (position > 0 && compareIds(order.get(position - 1), order.get(position)) != 0) { - rank++; - } - ranks[order.get(position)] = rank; - } - } - - int rank(int id) { - return ranks == null ? 0 : ranks[id]; - } - - BinaryRow partition(int id) { - return partitions[id]; - } - } - - static class Filter { - - final CompactFileIdentifierSet deletedIdentifiers; - final DeletedRowIdSet deletedRowIds; - final boolean useRowIdFilter; - final ThreadLocal identifier = - ThreadLocal.withInitial(IdentifierEncoder::new); - - Filter( - CompactFileIdentifierSet deletedIdentifiers, - DeletedRowIdSet deletedRowIds, - boolean useRowIdFilter) { - this.deletedIdentifiers = deletedIdentifiers; - this.deletedRowIds = deletedRowIds; - this.useRowIdFilter = useRowIdFilter; - } - - boolean include(ProjectedManifestEntry entry) { - return entry.isAdd() && !deletedIdentifiers.contains(entry); - } - - boolean include(GenericRow record, Key key) { - return key.kind == FileKind.ADD.toByteValue() && !isDeleted(record, key); - } - - boolean copyable(GenericRow record, Key key) { - return include(record, key); - } - - void observe(GenericRow record, Key key) {} - - boolean copyableAfterDiscovery(long minRowId, long maxRowId) { - return true; - } - - ReusableIdentifier identifier(GenericRow record) { - return identifier.get().replace(record); - } - - boolean isDeleted(GenericRow record, Key key) { - // RowID is only a cheap negative filter. The complete identifier remains the - // authoritative match, and is also sufficient for manifests which predate RowID. - if (useRowIdFilter) { - checkState(key.hasRowId, "First row id should not be null."); - if (!deletedRowIds.contains(key.firstRowId)) { - return false; - } - } - return deletedIdentifiers.contains(identifier(record)); - } - - static final class Minor extends Filter { - - Minor( - CompactFileIdentifierSet deletedIdentifiers, - DeletedRowIdSet deletedRowIds, - boolean useRowIdFilter) { - super(deletedIdentifiers, deletedRowIds, useRowIdFilter); - } - - @Override - boolean include(ProjectedManifestEntry entry) { - return true; - } - - @Override - boolean include(GenericRow record, Key key) { - return true; - } - - @Override - boolean copyable(GenericRow record, Key key) { - return key.kind == FileKind.ADD.toByteValue(); - } - - @Override - void observe(GenericRow record, Key key) { - if (key.kind != FileKind.DELETE.toByteValue()) { - return; - } - ReusableIdentifier reusable = identifier(record); - synchronized (this) { - deletedIdentifiers.add(reusable); - if (useRowIdFilter) { - checkState(key.hasRowId, "First row id should not be null."); - deletedRowIds.add(key.firstRowId); - } - } - } - - @Override - boolean copyableAfterDiscovery(long minRowId, long maxRowId) { - // A DELETE preserves the deleted ADD's globally unique first RowID. A range hit may - // be a false positive and only disables block copying; a miss proves the block has - // no deleted ADD. - return useRowIdFilter && !deletedRowIds.intersects(minRowId, maxRowId); - } - } - - private static final class IdentifierEncoder { - - final ProjectedManifestEntry entry = - ProjectedManifestEntry.Projection.create(ManifestEntryRunMerge.ENTRY_LAYOUT) - .createEntry(); - final ReusableIdentifier identifier = new ReusableIdentifier(); - - ReusableIdentifier replace(GenericRow record) { - return identifier.replaceWithPartition(entry.replace(record)); - } - } - } - - static InternalRow file(GenericRow record) { - return record.getRow(ManifestEntryRunMerge.FILE, ManifestEntryRunMerge.FILE_FIELD_COUNT); - } -} diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java index eb7702b8aa5f..997bfb362b7b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestEntryRunMergePlan.java @@ -23,8 +23,8 @@ import org.apache.paimon.data.InternalRow; import org.apache.paimon.data.serializer.InternalRowSerializer; import org.apache.paimon.format.avro.AvroRawBlock; +import org.apache.paimon.manifest.CollectedDeletes; import org.apache.paimon.manifest.CompactFileIdentifierSet; -import org.apache.paimon.manifest.DeletedRowIdSet; import org.apache.paimon.manifest.FileEntry.ReusableIdentifier; import org.apache.paimon.manifest.FileKind; import org.apache.paimon.manifest.ManifestAvroReader; @@ -53,25 +53,31 @@ final class ManifestEntryRunMergePlan { final List sources; - final ManifestEntryRunMergeEntry.PartitionDictionary partitions; + final ManifestEntryRunMerge.SortPartitionDictionary partitions; + final CollectedDeletes deletes; + final boolean minor; ManifestEntryRunMergePlan( - List sources, ManifestEntryRunMergeEntry.PartitionDictionary partitions) { + List sources, + ManifestEntryRunMerge.SortPartitionDictionary partitions, + CollectedDeletes deletes, + boolean minor) { this.sources = sources; this.partitions = partitions; + this.deletes = deletes; + this.minor = minor; } List mergeToManifest( ManifestFileSorter.RowIdEntrySortKey sortKey, ManifestFile manifestFile, - ManifestEntryRunMergeEntry.Filter filter, List newFilesForAbort) throws Exception { List cursors = new ArrayList<>(sources.size()); Exception failure = null; try { for (Source.Spec source : sources) { - Cursor cursor = source.open(manifestFile, sortKey, filter, partitions); + Cursor cursor = source.open(manifestFile, sortKey, deletes, minor, partitions); cursors.add(cursor); cursor.advance(); } @@ -100,16 +106,13 @@ List mergeToManifest( Pair, List> mergeMinorToManifest( ManifestFileSorter.RowIdEntrySortKey sortKey, ManifestFile manifestFile, - ManifestEntryRunMergeEntry.Filter filter, - CompactFileIdentifierSet deletedIdentifiers, - DeletedRowIdSet deletedRowIds, List newFilesForAbort) throws Exception { List cursors = new ArrayList<>(sources.size()); Exception failure = null; try { for (Source.Spec source : sources) { - Cursor cursor = source.open(manifestFile, sortKey, filter, partitions); + Cursor cursor = source.open(manifestFile, sortKey, deletes, minor, partitions); cursors.add(cursor); cursor.advance(); } @@ -118,8 +121,7 @@ Pair, List> mergeMinorToManifest( return Pair.of(Collections.emptyList(), Collections.emptyList()); } Pair, List> files = - writeMinorSelected( - selectionTree, manifestFile, deletedIdentifiers, deletedRowIds); + writeMinorSelected(selectionTree, manifestFile, deletes); newFilesForAbort.addAll(files.getLeft()); newFilesForAbort.addAll(files.getRight()); return files; @@ -169,10 +171,7 @@ static List writeSelected( } private static Pair, List> writeMinorSelected( - SelectionTree selectionTree, - ManifestFile manifestFile, - CompactFileIdentifierSet deletedIdentifiers, - DeletedRowIdSet deletedRowIds) + SelectionTree selectionTree, ManifestFile manifestFile, CollectedDeletes deletes) throws Exception { ManifestAvroWriter addWriter = manifestFile.createAvroWriter(); ManifestAvroWriter deleteWriter = manifestFile.createAvroWriter(); @@ -192,15 +191,11 @@ private static Pair, List> writeMinorSe cursor.materializeCurrent(); if (cursor.key().kind == FileKind.ADD.toByteValue()) { - if (!deletedRowIds.contains(cursor.key().firstRowId)) { - writeCurrent(addWriter, cursor); + ReusableIdentifier identifier = cursor.identifier(); + if (deletes.isDeleted(cursor.current(), identifier)) { + matchedEntries.add(identifier); } else { - ReusableIdentifier identifier = cursor.identifier(); - if (deletedIdentifiers.contains(identifier)) { - matchedEntries.add(identifier); - } else { - writeCurrent(addWriter, cursor); - } + writeCurrent(addWriter, cursor); } } else { ReusableIdentifier identifier = cursor.identifier(); @@ -270,8 +265,9 @@ interface Spec { Cursor open( ManifestFile manifestFile, ManifestFileSorter.RowIdEntrySortKey sortKey, - ManifestEntryRunMergeEntry.Filter filter, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) + CollectedDeletes deletes, + boolean minor, + ManifestEntryRunMerge.SortPartitionDictionary partitions) throws Exception; } @@ -311,11 +307,12 @@ long prefixBlockCount() { public Cursor open( ManifestFile manifestFile, ManifestFileSorter.RowIdEntrySortKey sortKey, - ManifestEntryRunMergeEntry.Filter filter, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) + CollectedDeletes deletes, + boolean minor, + ManifestEntryRunMerge.SortPartitionDictionary partitions) throws Exception { return new PrimitiveManifestRunCursor( - manifestFile, meta, start, end, blocks, filter, partitions); + manifestFile, meta, start, end, blocks, deletes, minor, partitions); } } @@ -331,10 +328,12 @@ static final class FragmentedManifestSpec implements Spec { public Cursor open( ManifestFile manifestFile, ManifestFileSorter.RowIdEntrySortKey sortKey, - ManifestEntryRunMergeEntry.Filter filter, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) + CollectedDeletes deletes, + boolean minor, + ManifestEntryRunMerge.SortPartitionDictionary partitions) throws Exception { - return new InMemoryManifestCursor(manifestFile, meta, sortKey, filter, partitions); + return new InMemoryManifestCursor( + manifestFile, meta, sortKey, deletes, minor, partitions); } } } @@ -351,7 +350,7 @@ interface Cursor extends AutoCloseable { @Nullable EncodedEntry metadata(); - ManifestEntryRunMergeEntry.Key key(); + ManifestEntryRunMerge.SortKey key(); @Nullable ByteBuffer encodedRecord(); @@ -366,7 +365,7 @@ default boolean hasCopyableBlock() { return false; } - default ManifestEntryRunMergeEntry.Key blockLastKey() { + default ManifestEntryRunMerge.SortKey blockLastKey() { throw new UnsupportedOperationException(); } @@ -392,10 +391,16 @@ static final class PrimitiveManifestRunCursor implements Cursor { final ManifestAvroReader reader; final boolean encodedRecordsCompatible; - final ManifestEntryRunMergeEntry.Filter filter; - final ManifestEntryRunMergeEntry.PartitionDictionary partitions; - final ManifestEntryRunMergeEntry.Key key = new ManifestEntryRunMergeEntry.Key(); + final CollectedDeletes deletes; + final boolean minor; + final ManifestEntryRunMerge.SortPartitionDictionary partitions; + final ManifestEntryRunMerge.SortKey key = new ManifestEntryRunMerge.SortKey(); final EncodedEntry metadata = new EncodedEntry(); + final ProjectedManifestEntry projectedEntry = + ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION.createEntry(); + final ProjectedManifestEntry fullEntry = + ProjectedManifestEntry.fullProjection().createEntry(); + final ReusableIdentifier identifier = new ReusableIdentifier(); final List blocks; final long runStart; final long runEnd; @@ -406,10 +411,8 @@ static final class PrimitiveManifestRunCursor implements Cursor { boolean current; @Nullable RawBlock currentRawBlock; @Nullable RowIterator currentRows; - @Nullable GenericRow currentRow; @Nullable GenericRow currentSourceRow; - @Nullable GenericRow compactRow; - @Nullable GenericRow compactFile; + @Nullable ProjectedManifestEntry currentEntry; @Nullable ManifestEntryRunMerge.Discovery.BlockInfo currentBlock; boolean closed; @@ -419,17 +422,14 @@ static final class PrimitiveManifestRunCursor implements Cursor { long start, long end, List blocks, - ManifestEntryRunMergeEntry.Filter filter, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) + CollectedDeletes deletes, + boolean minor, + ManifestEntryRunMerge.SortPartitionDictionary partitions) throws Exception { this.reader = manifestFile.scanAvroBlocks(meta.fileName(), meta.fileSize()); this.encodedRecordsCompatible = reader.rawBlockCopySupported(); - if (!encodedRecordsCompatible) { - this.compactRow = - new GenericRow(ManifestEntryRunMerge.ENTRY_LAYOUT.getFieldCount()); - this.compactFile = new GenericRow(ManifestEntryRunMerge.FILE_FIELD_COUNT); - } - this.filter = filter; + this.deletes = deletes; + this.minor = minor; this.partitions = partitions; this.blocks = blocks; this.runStart = start; @@ -469,24 +469,22 @@ public boolean advance() throws Exception { currentRows != null && currentRows.hasNext(), "Manifest block ends before its discovered boundary."); currentSourceRow = currentRows.next(); - currentRow = + currentEntry = encodedRecordsCompatible - ? currentSourceRow - : ManifestEntryRunMerge.projectEntryLayout( - currentSourceRow, compactRow, compactFile); + ? projectedEntry.replace(currentSourceRow) + : fullEntry.replace(currentSourceRow); decodedRemaining--; - key.replace(currentRow, partitions); - if (filter.include(currentRow, key)) { + key.replace(currentEntry, partitions); + if (minor || deletes.copyable(currentEntry, identifier, false)) { current = true; - InternalRow file = ManifestEntryRunMergeEntry.file(currentRow); metadata.replace( key.kind, partitions.partition(key.partitionId), - currentRow.getInt(ManifestEntryRunMerge.BUCKET), - file.getInt(ManifestEntryRunMerge.LEVEL), - file.getLong(ManifestEntryRunMerge.SCHEMA_ID), + currentEntry.bucket(), + currentEntry.file().level(), + currentEntry.file().schemaId(), key.firstRowId, - file.getLong(ManifestEntryRunMerge.ROW_COUNT)); + currentEntry.file().rowCount()); return true; } } @@ -496,8 +494,8 @@ boolean prepareNextBlock() throws Exception { rawBlock = false; current = false; currentRows = null; - currentRow = null; currentSourceRow = null; + currentEntry = null; while (blockIndex < blocks.size()) { ManifestEntryRunMerge.Discovery.BlockInfo info = blocks.get(blockIndex); if (info.start >= runEnd) { @@ -524,7 +522,8 @@ boolean prepareNextBlock() throws Exception { currentRows = currentRawBlock.toRows( encodedRecordsCompatible - ? ManifestEntryRunMerge.ENTRY_LAYOUT + ? ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION + .projectedType() : ManifestEntry.MANIFEST_ROW_TYPE); for (long i = 0; i < prefix; i++) { checkState( @@ -548,7 +547,7 @@ public boolean hasCurrent() { @Override public ProjectedManifestEntry current() { - return null; + return current ? currentEntry : null; } @Override @@ -557,7 +556,7 @@ public EncodedEntry metadata() { } @Override - public ManifestEntryRunMergeEntry.Key key() { + public ManifestEntryRunMerge.SortKey key() { return key; } @@ -574,7 +573,7 @@ public InternalRow decodedRow() { @Override public ReusableIdentifier identifier() { checkState(current, "Manifest entry has not been materialized."); - return filter.identifier(currentRow); + return identifier.replaceWithPartition(currentEntry); } @Override @@ -583,7 +582,7 @@ public boolean hasCopyableBlock() { } @Override - public ManifestEntryRunMergeEntry.Key blockLastKey() { + public ManifestEntryRunMerge.SortKey blockLastKey() { return currentBlock.lastKey; } @@ -617,30 +616,28 @@ public void materializeCurrent() throws Exception { currentRows = currentRawBlock.toRows( encodedRecordsCompatible - ? ManifestEntryRunMerge.ENTRY_LAYOUT + ? ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION.projectedType() : ManifestEntry.MANIFEST_ROW_TYPE); checkState(currentRows.hasNext(), "Manifest block cannot be decompressed."); currentSourceRow = currentRows.next(); - currentRow = + currentEntry = encodedRecordsCompatible - ? currentSourceRow - : ManifestEntryRunMerge.projectEntryLayout( - currentSourceRow, compactRow, compactFile); + ? projectedEntry.replace(currentSourceRow) + : fullEntry.replace(currentSourceRow); decodedRemaining--; - key.replace(currentRow, partitions); + key.replace(currentEntry, partitions); checkState( - filter.include(currentRow, key), + minor || deletes.copyable(currentEntry, identifier, false), "Copyable manifest block contains a filtered entry."); current = true; - InternalRow file = ManifestEntryRunMergeEntry.file(currentRow); metadata.replace( key.kind, partitions.partition(key.partitionId), - currentRow.getInt(ManifestEntryRunMerge.BUCKET), - file.getInt(ManifestEntryRunMerge.LEVEL), - file.getLong(ManifestEntryRunMerge.SCHEMA_ID), + currentEntry.bucket(), + currentEntry.file().level(), + currentEntry.file().schemaId(), key.firstRowId, - file.getLong(ManifestEntryRunMerge.ROW_COUNT)); + currentEntry.file().rowCount()); blockIndex++; } @@ -653,10 +650,11 @@ public void close() throws Exception { current = false; currentRawBlock = null; currentRows = null; - currentRow = null; currentSourceRow = null; - compactRow = null; - compactFile = null; + currentEntry = null; + projectedEntry.clear(); + fullEntry.clear(); + identifier.release(); currentBlock = null; rawBlock = false; key.clear(); @@ -676,8 +674,9 @@ static final class InMemoryManifestCursor implements Cursor { ManifestFile manifestFile, ManifestFileMeta meta, ManifestFileSorter.RowIdEntrySortKey sortKey, - ManifestEntryRunMergeEntry.Filter filter, - ManifestEntryRunMergeEntry.PartitionDictionary partitions) + CollectedDeletes deletes, + boolean minor, + ManifestEntryRunMerge.SortPartitionDictionary partitions) throws Exception { long entryCount = meta.numAddedFiles() + meta.numDeletedFiles(); this.entries = new ArrayList<>((int) entryCount); @@ -688,14 +687,14 @@ static final class InMemoryManifestCursor implements Cursor { manifestFile.scan(meta.fileName(), ProjectedManifestEntry.fullProjection())) { while (iterator.hasNext()) { ProjectedManifestEntry entry = iterator.next(); - if (!filter.include(entry)) { + if (!minor && !deletes.copyable(entry, identifier, false)) { continue; } BinaryRow row = serializer.toBinaryRow(entry.fullRow()).copy(); entries.add( new StoredEntry( row, - ManifestEntryRunMergeEntry.Key.viewOf( + ManifestEntryRunMerge.SortKey.viewOf( view.replace(row), partitions))); } } @@ -732,7 +731,7 @@ public EncodedEntry metadata() { } @Override - public ManifestEntryRunMergeEntry.Key key() { + public ManifestEntryRunMerge.SortKey key() { return entries.get(position).key; } @@ -758,9 +757,9 @@ public void close() { private static final class StoredEntry { final BinaryRow row; - final ManifestEntryRunMergeEntry.Key key; + final ManifestEntryRunMerge.SortKey key; - StoredEntry(BinaryRow row, ManifestEntryRunMergeEntry.Key key) { + StoredEntry(BinaryRow row, ManifestEntryRunMerge.SortKey key) { this.row = row; this.key = key; } @@ -821,7 +820,7 @@ int select(int left, int right) { return comparison < 0 || (comparison == 0 && left < right) ? left : right; } - boolean blockPrecedesOthers(int cursor, ManifestEntryRunMergeEntry.Key blockLastKey) { + boolean blockPrecedesOthers(int cursor, ManifestEntryRunMerge.SortKey blockLastKey) { for (int other = 0; other < cursors.size(); other++) { if (other == cursor || !cursors.get(other).hasCurrent()) { continue; diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileSorter.java b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileSorter.java index eda776751df8..76c5b0ef5ca0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileSorter.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ManifestFileSorter.java @@ -26,8 +26,7 @@ import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.disk.IOManager; -import org.apache.paimon.manifest.CompactFileIdentifierSet; -import org.apache.paimon.manifest.DeletedRowIdSet; +import org.apache.paimon.manifest.CollectedDeletes; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.manifest.ManifestFile; import org.apache.paimon.manifest.ManifestFileMeta; @@ -75,8 +74,7 @@ static class CompactionContext { final ManifestSortKey sortKey; final RowType partitionType; final ManifestEntryExternalSort.ExternalSortConfig externalSortConfig; - final CompactFileIdentifierSet deleteEntries; - final DeletedRowIdSet deletedRowIds; + final CollectedDeletes deletes; /** * Manifest files that need unsorted compaction. * @@ -96,8 +94,7 @@ static class CompactionContext { ManifestSortKey sortKey, RowType partitionType, ManifestEntryExternalSort.ExternalSortConfig externalSortConfig, - CompactFileIdentifierSet deleteEntries, - DeletedRowIdSet deletedRowIds, + CollectedDeletes deletes, Map defaultCompactFiles, List levelRuns, List pickedRuns) { @@ -106,8 +103,7 @@ static class CompactionContext { this.sortKey = sortKey; this.partitionType = partitionType; this.externalSortConfig = externalSortConfig; - this.deleteEntries = deleteEntries; - this.deletedRowIds = deletedRowIds; + this.deletes = deletes; this.defaultCompactFiles = defaultCompactFiles; this.levelRuns = levelRuns; this.pickedRuns = pickedRuns; @@ -122,8 +118,7 @@ boolean isMarkedForDefaultCompaction(ManifestFileMeta file) { /** Result of classifying manifest files. */ static class ClassifyResult { final List lsmFiles; - final CompactFileIdentifierSet deleteEntries; - final DeletedRowIdSet deletedRowIds; + final CollectedDeletes deletes; /** * Manifest files that need unsorted compaction. * @@ -136,32 +131,14 @@ static class ClassifyResult { ClassifyResult( List lsmFiles, - CompactFileIdentifierSet deleteEntries, - DeletedRowIdSet deletedRowIds, + CollectedDeletes deletes, Map compactWithoutSort) { this.lsmFiles = lsmFiles; - this.deleteEntries = deleteEntries; - this.deletedRowIds = deletedRowIds; + this.deletes = deletes; this.compactWithoutSort = compactWithoutSort; } } - /** Binary identifiers and partition values collected from DELETE entries. */ - private static class DeletedEntryInfo { - final CompactFileIdentifierSet identifiers; - final DeletedRowIdSet rowIds; - final Set partitions; - - private DeletedEntryInfo( - CompactFileIdentifierSet identifiers, - DeletedRowIdSet rowIds, - Set partitions) { - this.identifiers = identifiers; - this.rowIds = rowIds; - this.partitions = partitions; - } - } - /** * Try to sort-rewrite the merged manifest list by a configured partition field. If the sort * field cannot be resolved, the input is returned as-is. @@ -268,64 +245,68 @@ private static Optional> tryFullCompaction( sortedRunSizeRatio, externalSortConfig, manifestReadParallelism); - List levelRuns = ctx.levelRuns; - List pickedRuns = ctx.pickedRuns; - - if (pickedRuns.isEmpty() && ctx.defaultCompactFiles.isEmpty()) { - LOG.debug( - "Manifest sort full compact skipped: no runs picked and no defaultCompactFiles."); - return Optional.empty(); - } + try { + List levelRuns = ctx.levelRuns; + List pickedRuns = ctx.pickedRuns; + + if (pickedRuns.isEmpty() && ctx.defaultCompactFiles.isEmpty()) { + LOG.debug( + "Manifest sort full compact skipped: no runs picked and no defaultCompactFiles."); + return Optional.empty(); + } - LOG.info( - "Manifest sort full compact: input={} files, lsm={} runs, picked={} runs, " - + "defaultCompactFiles={}.", - input.size(), - levelRuns.size(), - pickedRuns.size(), - ctx.defaultCompactFiles.size()); - - // Step 3: Collect reused files (not picked) and picked files - Set pickedSet = new HashSet<>(pickedRuns); - List result = new ArrayList<>(); - for (ManifestAdjacentSortedRun run : levelRuns) { - if (!pickedSet.contains(run)) { - result.addAll(run.files()); + LOG.info( + "Manifest sort full compact: input={} files, lsm={} runs, picked={} runs, " + + "defaultCompactFiles={}.", + input.size(), + levelRuns.size(), + pickedRuns.size(), + ctx.defaultCompactFiles.size()); + + // Step 3: Collect reused files (not picked) and picked files + Set pickedSet = new HashSet<>(pickedRuns); + List result = new ArrayList<>(); + for (ManifestAdjacentSortedRun run : levelRuns) { + if (!pickedSet.contains(run)) { + result.addAll(run.files()); + } } - } - List pickedFiles = new ArrayList<>(); - for (ManifestAdjacentSortedRun run : pickedRuns) { - pickedFiles.addAll(run.files()); - } - pickedFiles.addAll(ctx.defaultCompactFiles.keySet()); + List pickedFiles = new ArrayList<>(); + for (ManifestAdjacentSortedRun run : pickedRuns) { + pickedFiles.addAll(run.files()); + } + pickedFiles.addAll(ctx.defaultCompactFiles.keySet()); - // Step 4: Split into sections and merge small adjacent sections - List
sections = splitIntoSections(pickedFiles, ctx); - sections = mergeSmallAdjacentSections(sections, suggestedMetaSize); + // Step 4: Split into sections and merge small adjacent sections + List
sections = splitIntoSections(pickedFiles, ctx); + sections = mergeSmallAdjacentSections(sections, suggestedMetaSize); - LOG.info( - "Manifest sort full compact: pickedFiles={}, sections={}.", - pickedFiles.size(), - sections.size()); + LOG.info( + "Manifest sort full compact: pickedFiles={}, sections={}.", + pickedFiles.size(), + sections.size()); - // Step 5: Rewrite sections - FullCompactOutput output = new FullCompactOutput(result); - rewriteSections( - sections, - output, - newFilesForAbort, - ctx, - manifestFile, - suggestedMetaSize, - suggestedMinMetaCount, - maxRewriteSize, - manifestReadParallelism); + // Step 5: Rewrite sections + FullCompactOutput output = new FullCompactOutput(result); + rewriteSections( + sections, + output, + newFilesForAbort, + ctx, + manifestFile, + suggestedMetaSize, + suggestedMinMetaCount, + maxRewriteSize, + manifestReadParallelism); - LOG.info( - "Manifest sort full compact completed: input={}, resultFiles={}.", - input.size(), - result.size()); - return Optional.of(result); + LOG.info( + "Manifest sort full compact completed: input={}, resultFiles={}.", + input.size(), + result.size()); + return Optional.of(result); + } finally { + ctx.deletes.release(); + } } /** @@ -365,95 +346,99 @@ private static List tryMinorCompaction( sortedRunSizeRatio, externalSortConfig, manifestReadParallelism); - List levelRuns = ctx.levelRuns; - List pickedRuns = ctx.pickedRuns; - - if (pickedRuns.isEmpty() && ctx.defaultCompactFiles.isEmpty()) { - LOG.debug( - "Manifest sort minor compact skipped: no runs picked and no defaultCompactFiles."); - return input; - } - - LOG.info( - "Manifest sort minor compact: input={} files, lsm={} runs, picked={} runs, " - + "defaultCompactFiles={}.", - input.size(), - levelRuns.size(), - pickedRuns.size(), - ctx.defaultCompactFiles.size()); - - // Step 2: Build fileName -> index mapping and initialize 2D result - Map fileNameToIndex = new HashMap<>(); - List> result = new ArrayList<>(input.size()); - for (int i = 0; i < input.size(); i++) { - fileNameToIndex.put(input.get(i).fileName(), i); - result.add(new ArrayList<>()); - } - - // Step 3: Collect reused files and picked files - Set pickedSet = new HashSet<>(pickedRuns); - for (ManifestAdjacentSortedRun run : levelRuns) { - if (!pickedSet.contains(run)) { - for (ManifestFileMeta file : run.files()) { - Integer idx = fileNameToIndex.get(file.fileName()); - if (idx != null) { - result.get(idx).add(file); + try { + List levelRuns = ctx.levelRuns; + List pickedRuns = ctx.pickedRuns; + + if (pickedRuns.isEmpty() && ctx.defaultCompactFiles.isEmpty()) { + LOG.debug( + "Manifest sort minor compact skipped: no runs picked and no defaultCompactFiles."); + return input; + } + + LOG.info( + "Manifest sort minor compact: input={} files, lsm={} runs, picked={} runs, " + + "defaultCompactFiles={}.", + input.size(), + levelRuns.size(), + pickedRuns.size(), + ctx.defaultCompactFiles.size()); + + // Step 2: Build fileName -> index mapping and initialize 2D result + Map fileNameToIndex = new HashMap<>(); + List> result = new ArrayList<>(input.size()); + for (int i = 0; i < input.size(); i++) { + fileNameToIndex.put(input.get(i).fileName(), i); + result.add(new ArrayList<>()); + } + + // Step 3: Collect reused files and picked files + Set pickedSet = new HashSet<>(pickedRuns); + for (ManifestAdjacentSortedRun run : levelRuns) { + if (!pickedSet.contains(run)) { + for (ManifestFileMeta file : run.files()) { + Integer idx = fileNameToIndex.get(file.fileName()); + if (idx != null) { + result.get(idx).add(file); + } } } } - } - List pickedFiles = new ArrayList<>(); - for (ManifestAdjacentSortedRun run : pickedRuns) { - pickedFiles.addAll(run.files()); - } - pickedFiles.addAll(ctx.defaultCompactFiles.keySet()); - - // Step 4: Compute index range - int minIdx = Integer.MAX_VALUE; - int maxIdx = Integer.MIN_VALUE; - for (ManifestFileMeta meta : pickedFiles) { - Integer idx = fileNameToIndex.get(meta.fileName()); - if (idx != null) { - minIdx = Math.min(minIdx, idx); - maxIdx = Math.max(maxIdx, idx); + List pickedFiles = new ArrayList<>(); + for (ManifestAdjacentSortedRun run : pickedRuns) { + pickedFiles.addAll(run.files()); } - } - Pair indexRange = Pair.of(minIdx, maxIdx); + pickedFiles.addAll(ctx.defaultCompactFiles.keySet()); + + // Step 4: Compute index range + int minIdx = Integer.MAX_VALUE; + int maxIdx = Integer.MIN_VALUE; + for (ManifestFileMeta meta : pickedFiles) { + Integer idx = fileNameToIndex.get(meta.fileName()); + if (idx != null) { + minIdx = Math.min(minIdx, idx); + maxIdx = Math.max(maxIdx, idx); + } + } + Pair indexRange = Pair.of(minIdx, maxIdx); - // Step 5: Split into sections and merge small adjacent sections - List
sections = splitIntoSections(pickedFiles, ctx); - sections = mergeSmallAdjacentSections(sections, suggestedMetaSize); + // Step 5: Split into sections and merge small adjacent sections + List
sections = splitIntoSections(pickedFiles, ctx); + sections = mergeSmallAdjacentSections(sections, suggestedMetaSize); - LOG.info( - "Manifest sort minor compact: pickedFiles={}, sections={}.", - pickedFiles.size(), - sections.size()); + LOG.info( + "Manifest sort minor compact: pickedFiles={}, sections={}.", + pickedFiles.size(), + sections.size()); - // Step 6: Rewrite sections - MinorCompactOutput output = new MinorCompactOutput(result, indexRange, fileNameToIndex); - rewriteSections( - sections, - output, - newFilesForAbort, - ctx, - manifestFile, - suggestedMetaSize, - suggestedMinMetaCount, - maxRewriteSize, - manifestReadParallelism); + // Step 6: Rewrite sections + MinorCompactOutput output = new MinorCompactOutput(result, indexRange, fileNameToIndex); + rewriteSections( + sections, + output, + newFilesForAbort, + ctx, + manifestFile, + suggestedMetaSize, + suggestedMinMetaCount, + maxRewriteSize, + manifestReadParallelism); - // Step 7: Flatten 2D result into a single list - List flatResult = new ArrayList<>(); - for (List subList : result) { - flatResult.addAll(subList); - } + // Step 7: Flatten 2D result into a single list + List flatResult = new ArrayList<>(); + for (List subList : result) { + flatResult.addAll(subList); + } - LOG.info( - "Manifest sort minor compact completed: input={}, resultFiles={}.", - input.size(), - flatResult.size()); - return flatResult; + LOG.info( + "Manifest sort minor compact completed: input={}, resultFiles={}.", + input.size(), + flatResult.size()); + return flatResult; + } finally { + ctx.deletes.release(); + } } /** @@ -508,8 +493,7 @@ private static CompactionContext prepareCompaction( sortKey, partitionType, externalSortConfig, - classification.deleteEntries, - classification.deletedRowIds, + classification.deletes, classification.compactWithoutSort, levelRuns, pickedRuns); @@ -535,7 +519,7 @@ static boolean reachesFullCompactionThreshold( *

Non-full compaction: small files go to defaultCompactFiles for minor-style merge; the rest * are returned as lsmFiles. * - * @return classification containing lsmFiles, deleteEntries, and defaultCompactFiles + * @return classification containing lsmFiles, collected DELETEs, and defaultCompactFiles */ static ClassifyResult classifyManifests( List input, @@ -565,28 +549,26 @@ private static ClassifyResult classifyManifests( // Initialize classification containers and read delete entries Map defaultCompactFiles = new LinkedHashMap<>(); List lsmFiles = new LinkedList<>(input); - CompactFileIdentifierSet classifiedDeleteEntries = new CompactFileIdentifierSet(); - DeletedRowIdSet deletedRowIds = new DeletedRowIdSet(); - Set deletePartitions = Collections.emptySet(); + CollectedDeletes deletes; PartitionPredicate predicate = null; if (fullCompaction) { - DeletedEntryInfo deletedEntries = + deletes = readDeletedEntries( manifestFile, input, runMergeOptimizeEnabled, manifestReadParallelism); - classifiedDeleteEntries = deletedEntries.identifiers; - deletedRowIds = deletedEntries.rowIds; - deletePartitions = deletedEntries.partitions; // Build partition predicate from delete entries for overlap detection. - if (classifiedDeleteEntries.isEmpty()) { + if (deletes.isEmpty()) { predicate = PartitionPredicate.ALWAYS_FALSE; } else { if (partitionType.getFieldCount() > 0) { - predicate = PartitionPredicate.fromMultiple(partitionType, deletePartitions); + predicate = + PartitionPredicate.fromMultiple(partitionType, deletes.partitions()); } else { predicate = PartitionPredicate.ALWAYS_TRUE; } } + } else { + deletes = new CollectedDeletes(runMergeOptimizeEnabled); } // Classify each file based on size and delete-partition overlap @@ -607,18 +589,15 @@ private static ClassifyResult classifyManifests( } } - return new ClassifyResult( - lsmFiles, classifiedDeleteEntries, deletedRowIds, defaultCompactFiles); + return new ClassifyResult(lsmFiles, deletes.toImmutable(), defaultCompactFiles); } - private static DeletedEntryInfo readDeletedEntries( + private static CollectedDeletes readDeletedEntries( ManifestFile manifestFile, List manifestFiles, boolean runMergeOptimizeEnabled, @Nullable Integer manifestReadParallelism) { - CompactFileIdentifierSet identifiers = new CompactFileIdentifierSet(); - DeletedRowIdSet rowIds = new DeletedRowIdSet(); - Set partitions = new HashSet<>(); + CollectedDeletes deletes = new CollectedDeletes(runMergeOptimizeEnabled); List filesWithDeletes = new ArrayList<>(); for (ManifestFileMeta meta : manifestFiles) { if (meta.numDeletedFiles() > 0) { @@ -629,44 +608,29 @@ private static DeletedEntryInfo readDeletedEntries( if (filesWithDeletes.size() <= 1 || (manifestReadParallelism != null && manifestReadParallelism <= 1)) { for (ManifestFileMeta meta : filesWithDeletes) { - collectDeletedEntries( - meta, - manifestFile, - identifiers, - rowIds, - partitions, - runMergeOptimizeEnabled, - false); + CollectedDeletes local = + collectDeletedEntries(meta, manifestFile, runMergeOptimizeEnabled); + deletes.combine(local); + local.release(); } } else { - Function> reader = - meta -> { - collectDeletedEntries( - meta, - manifestFile, - identifiers, - rowIds, - partitions, - runMergeOptimizeEnabled, - true); - return Collections.singletonList(Boolean.TRUE); - }; - for (Boolean ignored : + Function> reader = + meta -> + Collections.singletonList( + collectDeletedEntries( + meta, manifestFile, runMergeOptimizeEnabled)); + for (CollectedDeletes local : sequentialBatchedExecute(reader, filesWithDeletes, manifestReadParallelism)) { - // Iteration waits for each bounded batch of parallel reads. + deletes.combine(local); + local.release(); } } - return new DeletedEntryInfo(identifiers, rowIds, partitions); + return deletes; } - private static void collectDeletedEntries( - ManifestFileMeta meta, - ManifestFile manifestFile, - CompactFileIdentifierSet identifiers, - DeletedRowIdSet rowIds, - Set partitions, - boolean runMergeOptimizeEnabled, - boolean synchronize) { + private static CollectedDeletes collectDeletedEntries( + ManifestFileMeta meta, ManifestFile manifestFile, boolean runMergeOptimizeEnabled) { + CollectedDeletes deletes = new CollectedDeletes(runMergeOptimizeEnabled); try (CloseableIterator entries = manifestFile.scan( meta.fileName(), ProjectedManifestEntry.DELETE_ENTRY_PROJECTION)) { @@ -675,27 +639,14 @@ private static void collectDeletedEntries( if (!entry.isDelete()) { continue; } - BinaryRow partition = entry.partition().copy(); - if (synchronize) { - synchronized (identifiers) { - identifiers.add(entry); - if (runMergeOptimizeEnabled) { - rowIds.add(entry.file().nonNullFirstRowId()); - } - partitions.add(partition); - } - } else { - identifiers.add(entry); - if (runMergeOptimizeEnabled) { - rowIds.add(entry.file().nonNullFirstRowId()); - } - partitions.add(partition); - } + deletes.add(entry, runMergeOptimizeEnabled, true); } } catch (Exception e) { + deletes.release(); throw new RuntimeException( String.format("Failed to scan manifest file '%s'.", meta.fileName()), e); } + return deletes; } /** @@ -1089,7 +1040,7 @@ private static void unsortedCompactSection( } // Flush tail only if delete entries exist or file count >= minCount. if (!candidates.isEmpty()) { - if (!ctx.deleteEntries.isEmpty() || candidates.size() >= suggestedMinMetaCount) { + if (!ctx.deletes.isEmpty() || candidates.size() >= suggestedMinMetaCount) { rewriteSection( candidates, output, @@ -1153,8 +1104,7 @@ private static void rewriteFull( ctx.partitionType, manifestFile, sortNewFiles, - ctx.deleteEntries, - ctx.deletedRowIds, + ctx.deletes, ctx.externalSortConfig.maxNumFileHandles, manifestReadParallelism); } @@ -1166,7 +1116,7 @@ private static void rewriteFull( ctx.externalSortConfig, manifestFile, sortNewFiles, - ctx.deleteEntries, + ctx.deletes, manifestReadParallelism); } if (!sorted.isEmpty()) { diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ProjectedManifestEntryTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ProjectedManifestEntryTest.java index c78e7823671a..ede3a07402e1 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ProjectedManifestEntryTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ProjectedManifestEntryTest.java @@ -266,6 +266,25 @@ void testDoesNotValidateFileKindOnReplace() { assertThat(entry.isDelete()).isFalse(); } + @Test + void testEntryLayoutProjectionContainsRunMergeFields() { + RowType projectedType = ProjectedManifestEntry.ENTRY_LAYOUT_PROJECTION.projectedType(); + RowType projectedFileType = + (RowType) projectedType.getTypeAt(projectedType.getFieldIndex(ManifestEntry.FILE)); + + assertThat(projectedFileType.getFieldNames()) + .containsExactly( + DataFileMeta.FILE_NAME, + DataFileMeta.ROW_COUNT, + DataFileMeta.LEVEL, + DataFileMeta.SCHEMA_ID, + DataFileMeta.FIRST_ROW_ID, + DataFileMeta.MAX_SEQUENCE_NUMBER, + DataFileMeta.EXTRA_FILES, + DataFileMeta.EMBEDDED_FILE_INDEX, + DataFileMeta.EXTERNAL_PATH); + } + private static ProjectedManifestEntry.Projection projection( boolean includeBucket, String... projectedFileFields) { RowType manifestType = ManifestEntry.MANIFEST_ROW_TYPE;