JingsongLi commented on code in PR #9194: URL: https://github.com/apache/paimon/pull/9194#discussion_r3765535708
########## paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java: ########## @@ -0,0 +1,538 @@ +/* + * 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.manifest; + +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.format.SimpleColStats; +import org.apache.paimon.format.SimpleStatsCollector; +import org.apache.paimon.format.avro.AvroBlockWriter; +import org.apache.paimon.format.avro.AvroFileFormat; +import org.apache.paimon.format.avro.AvroRawBlock; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.PositionOutputStream; +import org.apache.paimon.io.RollingFileWriter; +import org.apache.paimon.schema.SchemaManager; +import org.apache.paimon.stats.SimpleStatsConverter; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.IOUtils; +import org.apache.paimon.utils.ObjectSerializer; +import org.apache.paimon.utils.PathFactory; + +import javax.annotation.Nullable; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.IdentityHashMap; +import java.util.List; +import java.util.Map; + +/** + * Avro writer for manifest entries. + * + * <p>The writer accepts materialized entries, encoded records and compressed Avro blocks. It is + * intentionally separate from the generic writer abstractions because encoded Avro data is an + * implementation detail of manifest run merging. + */ +public final class ManifestAvroWriter implements AutoCloseable { + + private final FileIO fileIO; + private final SchemaManager schemaManager; + private final RowType partitionType; + private final AvroFileFormat avroFileFormat; + private final ObjectSerializer<ManifestEntry> serializer; + private final String compression; + private final PathFactory pathFactory; + private final long targetFileSize; + + private final List<ManifestFileMeta> results = new ArrayList<>(); + private final List<Path> completedPaths = new ArrayList<>(); + private @Nullable FileWriter currentWriter; + private long recordCount; + private boolean closed; + + ManifestAvroWriter( + FileIO fileIO, + SchemaManager schemaManager, + RowType partitionType, + AvroFileFormat avroFileFormat, + ObjectSerializer<ManifestEntry> serializer, + String compression, + PathFactory pathFactory, + long targetFileSize) { + this.fileIO = fileIO; + this.schemaManager = schemaManager; + this.partitionType = partitionType; + this.avroFileFormat = avroFileFormat; + this.serializer = serializer; + this.compression = compression; + this.pathFactory = pathFactory; + this.targetFileSize = targetFileSize; + } + + public void write(ManifestEntry entry) throws IOException { + try { + currentWriter().write(entry); + afterWrite(1, false); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + public void write(Iterable<? extends ManifestEntry> entries) throws IOException { + for (ManifestEntry entry : entries) { + write(entry); + } + } + + public void writeEncoded(ByteBuffer encodedRecord, EncodedEntry metadata) throws IOException { + try { + currentWriter().writeEncoded(encodedRecord, metadata); + afterWrite(1, false); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + public void writeEncodedBlock(AvroRawBlock block, EncodedBlock metadata, long blockRecordCount) + throws IOException { + if (blockRecordCount != block.recordCount()) { + throw new IllegalArgumentException( + String.format( + "Manifest block record count mismatch: expected %s, actual %s.", + blockRecordCount, block.recordCount())); + } + try { + currentWriter().writeEncodedBlock(block, metadata); + afterWrite(blockRecordCount, true); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + private FileWriter currentWriter() { + if (closed) { + throw new IllegalStateException("Manifest writer has already closed."); + } + if (currentWriter == null) { + currentWriter = new FileWriter(pathFactory.newPath()); + } + return currentWriter; + } + + private void afterWrite(long addedRecords, boolean forceSizeCheck) throws IOException { + recordCount = Math.addExact(recordCount, addedRecords); + if (currentWriter.reachTargetSize( + forceSizeCheck || recordCount % RollingFileWriter.CHECK_ROLLING_RECORD_CNT == 0, + targetFileSize)) { + closeCurrentWriter(); + } + } + + private void closeCurrentWriter() throws IOException { + if (currentWriter == null) { + return; + } + currentWriter.close(); + ManifestFileMeta result = currentWriter.result(); + completedPaths.add(currentWriter.path); + results.add(result); + currentWriter = null; + } + + public long recordCount() { + return recordCount; + } + + public List<ManifestFileMeta> result() { + if (!closed) { + throw new IllegalStateException( + "Cannot access manifest results before closing the writer."); + } + return results; + } + + public void abort() { + if (currentWriter != null) { + currentWriter.abort(); + currentWriter = null; + } + for (Path path : completedPaths) { + fileIO.deleteQuietly(path); + } + completedPaths.clear(); + results.clear(); + closed = true; Review Comment: [P2] Keep aborted writes from exposing successful results abort() clears results but marks the writer as closed, while result() only checks closed. Consequently, an explicit abort or a write failure followed by result() returns an empty list instead of rejecting the failed state. The previous RollingFileWriterImpl did not treat abort as successful completion. Please track successful close separately from abort and make result() throw after an abort; an abort() -> result() regression test would cover this contract. ########## paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroWriter.java: ########## @@ -0,0 +1,538 @@ +/* + * 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.manifest; + +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.format.SimpleColStats; +import org.apache.paimon.format.SimpleStatsCollector; +import org.apache.paimon.format.avro.AvroBlockWriter; +import org.apache.paimon.format.avro.AvroFileFormat; +import org.apache.paimon.format.avro.AvroRawBlock; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.PositionOutputStream; +import org.apache.paimon.io.RollingFileWriter; +import org.apache.paimon.schema.SchemaManager; +import org.apache.paimon.stats.SimpleStatsConverter; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.IOUtils; +import org.apache.paimon.utils.ObjectSerializer; +import org.apache.paimon.utils.PathFactory; + +import javax.annotation.Nullable; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.IdentityHashMap; +import java.util.List; +import java.util.Map; + +/** + * Avro writer for manifest entries. + * + * <p>The writer accepts materialized entries, encoded records and compressed Avro blocks. It is + * intentionally separate from the generic writer abstractions because encoded Avro data is an + * implementation detail of manifest run merging. + */ +public final class ManifestAvroWriter implements AutoCloseable { + + private final FileIO fileIO; + private final SchemaManager schemaManager; + private final RowType partitionType; + private final AvroFileFormat avroFileFormat; + private final ObjectSerializer<ManifestEntry> serializer; + private final String compression; + private final PathFactory pathFactory; + private final long targetFileSize; + + private final List<ManifestFileMeta> results = new ArrayList<>(); + private final List<Path> completedPaths = new ArrayList<>(); + private @Nullable FileWriter currentWriter; + private long recordCount; + private boolean closed; + + ManifestAvroWriter( + FileIO fileIO, + SchemaManager schemaManager, + RowType partitionType, + AvroFileFormat avroFileFormat, + ObjectSerializer<ManifestEntry> serializer, + String compression, + PathFactory pathFactory, + long targetFileSize) { + this.fileIO = fileIO; + this.schemaManager = schemaManager; + this.partitionType = partitionType; + this.avroFileFormat = avroFileFormat; + this.serializer = serializer; + this.compression = compression; + this.pathFactory = pathFactory; + this.targetFileSize = targetFileSize; + } + + public void write(ManifestEntry entry) throws IOException { + try { + currentWriter().write(entry); + afterWrite(1, false); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + public void write(Iterable<? extends ManifestEntry> entries) throws IOException { + for (ManifestEntry entry : entries) { + write(entry); + } + } + + public void writeEncoded(ByteBuffer encodedRecord, EncodedEntry metadata) throws IOException { + try { + currentWriter().writeEncoded(encodedRecord, metadata); + afterWrite(1, false); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + public void writeEncodedBlock(AvroRawBlock block, EncodedBlock metadata, long blockRecordCount) + throws IOException { + if (blockRecordCount != block.recordCount()) { + throw new IllegalArgumentException( + String.format( + "Manifest block record count mismatch: expected %s, actual %s.", + blockRecordCount, block.recordCount())); + } + try { + currentWriter().writeEncodedBlock(block, metadata); + afterWrite(blockRecordCount, true); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } + } + + private FileWriter currentWriter() { + if (closed) { + throw new IllegalStateException("Manifest writer has already closed."); + } + if (currentWriter == null) { + currentWriter = new FileWriter(pathFactory.newPath()); + } + return currentWriter; + } + + private void afterWrite(long addedRecords, boolean forceSizeCheck) throws IOException { + recordCount = Math.addExact(recordCount, addedRecords); + if (currentWriter.reachTargetSize( + forceSizeCheck || recordCount % RollingFileWriter.CHECK_ROLLING_RECORD_CNT == 0, + targetFileSize)) { + closeCurrentWriter(); + } + } + + private void closeCurrentWriter() throws IOException { + if (currentWriter == null) { + return; + } + currentWriter.close(); + ManifestFileMeta result = currentWriter.result(); + completedPaths.add(currentWriter.path); + results.add(result); + currentWriter = null; + } + + public long recordCount() { + return recordCount; + } + + public List<ManifestFileMeta> result() { + if (!closed) { + throw new IllegalStateException( + "Cannot access manifest results before closing the writer."); + } + return results; + } + + public void abort() { + if (currentWriter != null) { + currentWriter.abort(); + currentWriter = null; + } + for (Path path : completedPaths) { + fileIO.deleteQuietly(path); + } + completedPaths.clear(); + results.clear(); + closed = true; + } + + @Override + public void close() throws IOException { + if (closed) { + return; + } + try { + closeCurrentWriter(); + } catch (IOException | RuntimeException | Error failure) { + abort(); + throw failure; + } finally { + closed = true; + } + } + + /** Reusable statistics needed when an encoded manifest record is copied directly. */ + public static final class EncodedEntry { + + private byte kind; + private BinaryRow partition; + private int bucket; + private int level; + private long schemaId; + private long firstRowId; + private long rowCount; + + public EncodedEntry replace( + byte kind, + BinaryRow partition, + int bucket, + int level, + long schemaId, + long firstRowId, + long rowCount) { + this.kind = kind; + this.partition = partition; + this.bucket = bucket; + this.level = level; + this.schemaId = schemaId; + this.firstRowId = firstRowId; + this.rowCount = rowCount; + return this; + } + } + + /** Aggregate statistics for an encoded Avro block copied without decompression. */ + public static final class EncodedBlock { + + private final long addedFiles; + private final long schemaId; + private final int minBucket; + private final int maxBucket; + private final int minLevel; + private final int maxLevel; + private final long minRowId; + private final long maxRowId; + private final @Nullable BinaryRow nullPartition; + private final long nullPartitionCount; + private final @Nullable BinaryRow minNonNullPartition; + private final @Nullable BinaryRow maxNonNullPartition; + + public EncodedBlock( + long addedFiles, + long schemaId, + int minBucket, + int maxBucket, + int minLevel, + int maxLevel, + long minRowId, + long maxRowId, + @Nullable BinaryRow nullPartition, + long nullPartitionCount, + @Nullable BinaryRow minNonNullPartition, + @Nullable BinaryRow maxNonNullPartition) { + this.addedFiles = addedFiles; + this.schemaId = schemaId; + this.minBucket = minBucket; + this.maxBucket = maxBucket; + this.minLevel = minLevel; + this.maxLevel = maxLevel; + this.minRowId = minRowId; + this.maxRowId = maxRowId; + this.nullPartition = nullPartition; + this.nullPartitionCount = nullPartitionCount; + this.minNonNullPartition = minNonNullPartition; + this.maxNonNullPartition = maxNonNullPartition; + } + } + + private final class FileWriter { + + private final Path path; + private final SimpleStatsCollector partitionStatsCollector; + private final SimpleStatsConverter partitionStatsSerializer; + private final Map<BinaryRow, long[]> encodedPartitionCounts = new IdentityHashMap<>(); + private final long[] repeatedNullCounts = new long[partitionType.getFieldCount()]; + private @Nullable PositionOutputStream out; + private @Nullable AvroBlockWriter writer; + private @Nullable Long outputBytes; + private long numAddedFiles; + private long numDeletedFiles; + 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 @Nullable RowIdStats rowIdStats = new RowIdStats(); + private boolean closed; + + private FileWriter(Path path) { + this.path = path; + this.partitionStatsCollector = new SimpleStatsCollector(partitionType); + this.partitionStatsSerializer = new SimpleStatsConverter(partitionType); + boolean outputCreated = false; + try { + out = fileIO.newOutputStream(path, false); + outputCreated = true; + writer = + avroFileFormat.createBlockWriter( + out, ManifestEntry.MANIFEST_ROW_TYPE, compression); + } catch (IOException failure) { + IOUtils.closeQuietly(writer); + IOUtils.closeQuietly(out); + if (outputCreated) { + fileIO.deleteQuietly(path); + } + throw new UncheckedIOException( + "Failed to create manifest Avro writer for " + path, failure); + } catch (RuntimeException | Error failure) { + IOUtils.closeQuietly(writer); + IOUtils.closeQuietly(out); + if (outputCreated) { + fileIO.deleteQuietly(path); + } + throw failure; + } + } + + private void write(ManifestEntry entry) throws IOException { + ensureOpen(); + writer.addElement( + entry instanceof ProjectedManifestEntry + ? ((ProjectedManifestEntry) entry).fullRow() + : serializer.toRow(entry)); + collectStats(entry); + } + + private void writeEncoded(ByteBuffer encodedRecord, EncodedEntry metadata) + throws IOException { + ensureOpen(); + writer.addEncoded(encodedRecord); + collectStats(metadata); + addEncodedPartition(metadata.partition, 1); + } + + private void writeEncodedBlock(AvroRawBlock block, EncodedBlock metadata) + throws IOException { + ensureOpen(); + writer.addEncodedBlock(block); + collectStats(metadata); + if (metadata.nullPartitionCount > 0) { + addEncodedPartition(metadata.nullPartition, metadata.nullPartitionCount); + } + if (metadata.minNonNullPartition != null) { + addEncodedPartition(metadata.minNonNullPartition, 1); + if (metadata.maxNonNullPartition != metadata.minNonNullPartition) { + addEncodedPartition(metadata.maxNonNullPartition, 1); + } + } + } + + private void collectStats(ManifestEntry entry) { + switch (entry.kind()) { + case ADD: + numAddedFiles++; + break; + case DELETE: + numDeletedFiles++; + break; + default: + throw new UnsupportedOperationException("Unknown entry kind: " + entry.kind()); + } + schemaId = Math.max(schemaId, entry.file().schemaId()); + minBucket = Math.min(minBucket, entry.bucket()); + maxBucket = Math.max(maxBucket, entry.bucket()); + minLevel = Math.min(minLevel, entry.level()); + maxLevel = Math.max(maxLevel, entry.level()); + if (rowIdStats != null) { + Long firstRowId = entry.file().firstRowId(); + if (firstRowId == null) { + rowIdStats = null; + } else { + rowIdStats.collect(firstRowId, entry.file().rowCount()); + } + } + partitionStatsCollector.collect(entry.partition()); + } + + private void collectStats(EncodedEntry entry) { + switch (FileKind.fromByteValue(entry.kind)) { + case ADD: + numAddedFiles++; + break; + case DELETE: + numDeletedFiles++; + break; + default: + throw new UnsupportedOperationException("Unknown entry kind: " + entry.kind); + } + schemaId = Math.max(schemaId, entry.schemaId); + minBucket = Math.min(minBucket, entry.bucket); + maxBucket = Math.max(maxBucket, entry.bucket); + minLevel = Math.min(minLevel, entry.level); + maxLevel = Math.max(maxLevel, entry.level); + if (rowIdStats != null) { + rowIdStats.collect(entry.firstRowId, entry.rowCount); + } + } + + private void collectStats(EncodedBlock block) { + numAddedFiles = Math.addExact(numAddedFiles, block.addedFiles); Review Comment: [P1] Account for DELETE entries in copied raw blocks This path increments only numAddedFiles. For a mixed or DELETE-only raw block, the Avro records are copied correctly while numDeletedFiles remains zero. Downstream, FileEntry.readDeletedEntries filters out manifests whose numDeletedFiles() == 0, so these DELETE records are never applied and logically deleted files may be resurrected. I reproduced this by copying a normal mixed block: the output contained 36 DELETE records but reported 0. Please validate 0 <= addedFiles <= blockRecordCount and add blockRecordCount - addedFiles to numDeletedFiles (or carry both counts explicitly), with mixed and DELETE-only block tests. -- 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]
