This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new f2a59c20ea [core] Unify manifest Avro readers (#9182)
f2a59c20ea is described below

commit f2a59c20eadb6d0cb6efec3016d6165479463b07
Author: YeJunHao <[email protected]>
AuthorDate: Wed Aug 12 13:17:50 2026 +0800

    [core] Unify manifest Avro readers (#9182)
---
 .../apache/paimon/manifest/ManifestAvroReader.java | 409 ++++++++++++++++-----
 .../org/apache/paimon/manifest/ManifestFile.java   |   4 +-
 .../paimon/manifest/ProjectedManifestEntry.java    |   3 +-
 .../apache/paimon/manifest/ManifestFileTest.java   | 321 +++++++++++++---
 .../main/java/org/apache/avro/file/RawBlock.java   | 168 +++++++++
 .../java/org/apache/avro/file/RawBlockReader.java  |  56 +++
 .../apache/paimon/format/avro/AvroBlockReader.java | 110 ++----
 .../apache/paimon/format/avro/AvroRawBlock.java    |  64 ++++
 .../paimon/format/avro/AvroRecordDecoder.java      |  37 ++
 .../paimon/format/avro/AvroFileFormatTest.java     | 115 +++---
 10 files changed, 982 insertions(+), 305 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java 
b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java
index 83af2833eb..3fbb000871 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java
@@ -22,7 +22,7 @@ import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.format.avro.AvroBlockReader;
-import org.apache.paimon.format.avro.AvroBlockReader.BorrowedBlock;
+import org.apache.paimon.format.avro.AvroRawBlock;
 import org.apache.paimon.format.avro.AvroRecordDecoder;
 import org.apache.paimon.format.avro.AvroRecordDecoder.FieldDecoder;
 import org.apache.paimon.format.avro.AvroRecordDecoder.FieldType;
@@ -36,106 +36,113 @@ import javax.annotation.Nullable;
 import java.io.IOException;
 import java.io.InputStream;
 import java.io.UncheckedIOException;
+import java.nio.ByteBuffer;
+import java.util.Iterator;
 import java.util.NoSuchElementException;
 
 import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow;
 
 /**
  * Schema-aware Avro reader which projects fields and filters before decoding 
data file metadata.
+ *
+ * <p>This reader also exposes reusable raw Avro blocks.
  */
-final class ManifestAvroReader implements CloseableIterator<InternalRow> {
+public final class ManifestAvroReader implements AutoCloseable {
 
     private final AvroBlockReader blockReader;
-    private final AvroRecordDecoder decoder;
-    private final ManifestRecordDecoder recordDecoder;
+    private final DecoderContext decoderContext;
+    private final boolean rawBlockCopySupported;
 
-    private long recordsRemaining;
-    private @Nullable InternalRow next;
-    private boolean nextReady;
-    private boolean finished;
+    private long blockOrdinal = -1;
 
-    ManifestAvroReader(
-            InputStream input,
-            RowType projectedType,
-            @Nullable PartitionPredicate partitionFilter,
-            @Nullable BucketFilter bucketFilter)
-            throws IOException {
-        AvroBlockReader blockReader = new AvroBlockReader(input);
+    ManifestAvroReader(InputStream input) throws IOException {
+        AvroBlockReader blockReader = null;
         try {
-            AvroRecordDecoder decoder = blockReader.createRecordDecoder();
-            this.recordDecoder =
-                    new ManifestRecordDecoder(
-                            decoder, projectedType, partitionFilter, 
bucketFilter);
-            this.decoder = decoder;
+            blockReader = new AvroBlockReader(input);
             this.blockReader = blockReader;
-        } catch (RuntimeException | Error e) {
-            IOUtils.closeQuietly(blockReader);
-            throw e;
+            this.decoderContext = new 
DecoderContext(blockReader.createRecordDecoder());
+            this.rawBlockCopySupported =
+                    
blockReader.supportsRawBlockCopy(ManifestEntry.MANIFEST_ROW_TYPE);
+        } catch (IOException | RuntimeException | Error failure) {
+            IOUtils.closeQuietly(blockReader == null ? input : blockReader);
+            throw failure;
         }
     }
 
-    @Override
-    public boolean hasNext() {
-        if (nextReady) {
-            return true;
-        }
-        if (finished) {
-            return false;
-        }
-
-        try {
-            while (true) {
-                if (recordsRemaining == 0) {
-                    if (decoder.isInitialized() && !decoder.isEnd()) {
-                        throw new IOException(
-                                "Manifest Avro block contains trailing 
undecoded bytes.");
-                    }
-                    if (!blockReader.hasNextBlock()) {
-                        finished = true;
-                        return false;
-                    }
-                    BorrowedBlock block = blockReader.nextBorrowedBlock();
-                    decoder.reset(block.bytes(), block.offset(), 
block.length());
-                    recordsRemaining = block.recordCount();
-                }
-
-                InternalRow candidate = recordDecoder.read(decoder);
-                recordsRemaining--;
-                if (candidate != null) {
-                    next = candidate;
-                    nextReady = true;
-                    return true;
-                }
-            }
-        } catch (IOException e) {
-            throw new UncheckedIOException("Failed to decode Manifest Avro 
block.", e);
-        }
+    /** Returns whether another raw Avro block is available. */
+    public boolean hasNext() throws IOException {
+        return blockReader.hasNextBlock();
     }
 
-    @Override
-    public InternalRow next() {
+    /** Returns the next raw block without decompressing it. */
+    public RawBlock next() throws IOException {
         if (!hasNext()) {
             throw new NoSuchElementException();
         }
-        InternalRow result = next;
-        next = null;
-        nextReady = false;
-        return result;
+        return new RawBlock(
+                decoderContext,
+                rawBlockCopySupported,
+                blockReader.nextBorrowedRawBlock(),
+                ++blockOrdinal);
     }
 
-    long decodedDataFiles() {
-        return recordDecoder.decodedDataFiles;
-    }
+    /**
+     * Returns an iterator over projected rows from all remaining blocks.
+     *
+     * <p>Every returned row has independent backing data. Closing the 
iterator closes this reader.
+     */
+    public CloseableIterator<InternalRow> read(
+            RowType projectedType,
+            @Nullable PartitionPredicate partitionFilter,
+            @Nullable BucketFilter bucketFilter) {
+        return new CloseableIterator<InternalRow>() {
+
+            private @Nullable RowIterator rows;
+            private boolean closed;
+
+            @Override
+            public boolean hasNext() {
+                if (closed) {
+                    return false;
+                }
+                try {
+                    while (rows == null || !rows.hasNext()) {
+                        if (!ManifestAvroReader.this.hasNext()) {
+                            return false;
+                        }
+                        rows =
+                                ManifestAvroReader.this
+                                        .next()
+                                        .toRows(
+                                                projectedType,
+                                                partitionFilter,
+                                                bucketFilter,
+                                                false);
+                    }
+                    return true;
+                } catch (IOException e) {
+                    throw new UncheckedIOException("Failed to decode Manifest 
Avro block.", e);
+                }
+            }
 
-    long skippedDataFiles() {
-        return recordDecoder.skippedDataFiles;
+            @Override
+            public InternalRow next() {
+                if (!hasNext()) {
+                    throw new NoSuchElementException();
+                }
+                return rows.next();
+            }
+
+            @Override
+            public void close() throws IOException {
+                closed = true;
+                ManifestAvroReader.this.close();
+            }
+        };
     }
 
     @Override
     public void close() throws IOException {
-        next = null;
-        nextReady = false;
-        finished = true;
         blockReader.close();
     }
 
@@ -157,22 +164,9 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
         private final int bucketPosition;
         private final int totalBucketsPosition;
         private final int filePosition;
-
         private final FieldDecoder fileReader;
 
-        private final @Nullable PartitionPredicate partitionFilter;
-        private final @Nullable BucketFilter bucketFilter;
-        private final boolean partitionNeededForFilter;
-        private final boolean bucketNeededForFilter;
-
-        private long decodedDataFiles;
-        private long skippedDataFiles;
-
-        private ManifestRecordDecoder(
-                AvroRecordDecoder decoder,
-                RowType projectedType,
-                @Nullable PartitionPredicate partitionFilter,
-                @Nullable BucketFilter bucketFilter) {
+        private ManifestRecordDecoder(AvroRecordDecoder decoder, RowType 
projectedType) {
             this.projectedFieldCount = projectedType.getFieldCount();
             this.versionPosition =
                     
projectedType.getFieldIndex(ManifestSchemaUtils.FORMAT_IDENTIFIER_FIELD);
@@ -181,10 +175,6 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
             this.bucketPosition = 
projectedType.getFieldIndex(ManifestEntry.BUCKET);
             this.totalBucketsPosition = 
projectedType.getFieldIndex(ManifestEntry.TOTAL_BUCKETS);
             this.filePosition = 
projectedType.getFieldIndex(ManifestEntry.FILE);
-            this.partitionFilter = partitionFilter;
-            this.bucketFilter = bucketFilter;
-            this.partitionNeededForFilter = partitionFilter != null || 
bucketFilter != null;
-            this.bucketNeededForFilter = bucketFilter != null;
 
             validateTopLevelFields(decoder);
             // Manifest v2 has a fixed top-level layout, but DataFileMeta has 
gained nullable
@@ -195,7 +185,12 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
                             5, filePosition >= 0 ? 
projectedType.getTypeAt(filePosition) : null);
         }
 
-        private @Nullable InternalRow read(AvroRecordDecoder decoder) throws 
IOException {
+        private boolean read(
+                AvroRecordDecoder decoder,
+                GenericRow row,
+                @Nullable PartitionPredicate partitionFilter,
+                @Nullable BucketFilter bucketFilter)
+                throws IOException {
             if (!decoder.readRecordStart()) {
                 throw new IOException("Unexpected null or non-record Manifest 
Avro value.");
             }
@@ -203,6 +198,7 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
             int version = decoder.readInt();
             ManifestEntrySerializer.checkFormatIdentifier(version);
             int kind = decoder.readInt();
+            boolean partitionNeededForFilter = partitionFilter != null || 
bucketFilter != null;
             byte[] partitionBytes;
             if (partitionPosition >= 0 || partitionNeededForFilter) {
                 partitionBytes = decoder.readBytes();
@@ -215,9 +211,10 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
                     partitionNeededForFilter ? 
deserializeBinaryRow(partitionBytes) : null;
             if (partitionFilter != null && !partitionFilter.test(partition)) {
                 skipBucketAndFile(decoder);
-                return null;
+                return false;
             }
 
+            boolean bucketNeededForFilter = bucketFilter != null;
             int bucket;
             if (bucketPosition >= 0 || bucketNeededForFilter) {
                 bucket = decoder.readInt();
@@ -236,27 +233,24 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
 
             if (bucketFilter != null && !bucketFilter.test(partition, bucket, 
totalBuckets)) {
                 skipFile(decoder);
-                return null;
+                return false;
             }
 
             Object file;
             if (filePosition >= 0) {
-                file = fileReader.read(decoder, null);
-                decodedDataFiles++;
+                file = fileReader.read(decoder, row.getField(filePosition));
             } else {
                 fileReader.skip(decoder);
-                skippedDataFiles++;
                 file = null;
             }
 
-            GenericRow row = new GenericRow(projectedFieldCount);
             setProjected(row, versionPosition, version);
             setProjected(row, kindPosition, (byte) kind);
             setProjected(row, partitionPosition, partitionBytes);
             setProjected(row, bucketPosition, bucket);
             setProjected(row, totalBucketsPosition, totalBuckets);
             setProjected(row, filePosition, file);
-            return row;
+            return true;
         }
 
         private void skipBucketAndFile(AvroRecordDecoder decoder) throws 
IOException {
@@ -267,7 +261,6 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
 
         private void skipFile(AvroRecordDecoder decoder) throws IOException {
             fileReader.skip(decoder);
-            skippedDataFiles++;
         }
 
         private static void setProjected(
@@ -314,4 +307,220 @@ final class ManifestAvroReader implements 
CloseableIterator<InternalRow> {
             }
         }
     }
+
+    /** Borrowed raw block which must be consumed before the enclosing reader 
advances. */
+    public static final class RawBlock {
+
+        private final DecoderContext decoderContext;
+        private final boolean rawBlockCopySupported;
+        private final AvroRawBlock block;
+        private final long blockOrdinal;
+        private final long blockRecordCount;
+
+        private RawBlock(
+                DecoderContext decoderContext,
+                boolean rawBlockCopySupported,
+                AvroRawBlock block,
+                long blockOrdinal) {
+            this.decoderContext = decoderContext;
+            this.rawBlockCopySupported = rawBlockCopySupported;
+            this.block = block;
+            this.blockOrdinal = blockOrdinal;
+            this.blockRecordCount = block.recordCount();
+        }
+
+        /** Lazily decompresses this block and returns an iterator over one 
reusable row. */
+        public RowIterator toRows(RowType projectedType) throws IOException {
+            return toRows(projectedType, null, null, true);
+        }
+
+        private RowIterator toRows(
+                RowType projectedType,
+                @Nullable PartitionPredicate partitionFilter,
+                @Nullable BucketFilter bucketFilter,
+                boolean reuseRow)
+                throws IOException {
+            ManifestRecordDecoder recordDecoder = 
decoderContext.recordDecoder(projectedType);
+
+            ByteBuffer decompressed = decoderContext.decompress(block);
+            decoderContext.decoder.reset(
+                    decompressed.array(),
+                    decompressed.arrayOffset() + decompressed.position(),
+                    decompressed.remaining());
+            GenericRow reuse = reuseRow ? new 
GenericRow(recordDecoder.projectedFieldCount) : null;
+            return new RowIterator(
+                    blockRecordCount,
+                    decoderContext.decoder,
+                    recordDecoder,
+                    reuse,
+                    partitionFilter,
+                    bucketFilter);
+        }
+
+        public long blockOrdinal() {
+            return blockOrdinal;
+        }
+
+        public long recordCount() {
+            return blockRecordCount;
+        }
+
+        public boolean rawBlockCopySupported() {
+            return rawBlockCopySupported;
+        }
+
+        public AvroRawBlock encodedBlock() {
+            return block;
+        }
+    }
+
+    /** Decoder state shared by the borrowed blocks produced by one reader. */
+    private static final class DecoderContext {
+
+        private final AvroRecordDecoder decoder;
+
+        private @Nullable ByteBuffer decompressionBuffer;
+        private RowType projectedRowType;
+        private ManifestRecordDecoder recordDecoder;
+
+        private DecoderContext(AvroRecordDecoder decoder) {
+            this.decoder = decoder;
+        }
+
+        private ByteBuffer decompress(AvroRawBlock block) throws IOException {
+            decompressionBuffer = block.decompress(decompressionBuffer);
+            return decompressionBuffer;
+        }
+
+        private ManifestRecordDecoder recordDecoder(RowType rowType) {
+            if (!rowType.equals(projectedRowType)) {
+                recordDecoder = new ManifestRecordDecoder(decoder, rowType);
+                projectedRowType = rowType;
+            }
+            return recordDecoder;
+        }
+    }
+
+    /** Iterator over the reusable row decoded from one borrowed block. */
+    public static final class RowIterator implements Iterator<GenericRow> {
+
+        private final AvroRecordDecoder decoder;
+        private final ManifestRecordDecoder recordDecoder;
+        private final @Nullable GenericRow reuseRow;
+        private final @Nullable PartitionPredicate partitionFilter;
+        private final @Nullable BucketFilter bucketFilter;
+        private final boolean filtered;
+
+        private long blockRemaining;
+        private long blockRecordIndex = -1;
+        private @Nullable GenericRow next;
+        private @Nullable ByteBuffer encodedRecord;
+
+        private RowIterator(
+                long recordCount,
+                AvroRecordDecoder decoder,
+                ManifestRecordDecoder recordDecoder,
+                @Nullable GenericRow reuseRow,
+                @Nullable PartitionPredicate partitionFilter,
+                @Nullable BucketFilter bucketFilter) {
+            blockRemaining = recordCount;
+            this.decoder = decoder;
+            this.recordDecoder = recordDecoder;
+            this.reuseRow = reuseRow;
+            this.partitionFilter = partitionFilter;
+            this.bucketFilter = bucketFilter;
+            this.filtered = partitionFilter != null || bucketFilter != null;
+        }
+
+        @Override
+        public boolean hasNext() {
+            try {
+                if (!filtered) {
+                    ensureBlockFullyConsumed();
+                    return blockRemaining > 0;
+                }
+                if (next != null) {
+                    return true;
+                }
+                while (blockRemaining > 0) {
+                    blockRecordIndex++;
+                    GenericRow row =
+                            reuseRow == null
+                                    ? new 
GenericRow(recordDecoder.projectedFieldCount)
+                                    : reuseRow;
+                    int recordStart = decoder.absolutePosition();
+                    boolean selected =
+                            recordDecoder.read(decoder, row, partitionFilter, 
bucketFilter);
+                    blockRemaining--;
+                    if (selected) {
+                        encodedRecord =
+                                decoder.borrowedView(recordStart, 
decoder.absolutePosition());
+                        next = row;
+                        return true;
+                    }
+                }
+                ensureBlockFullyConsumed();
+                return false;
+            } catch (IOException e) {
+                throw new UncheckedIOException(
+                        "Failed to decode projected Manifest Avro record.", e);
+            }
+        }
+
+        @Override
+        public GenericRow next() {
+            if (filtered) {
+                if (!hasNext()) {
+                    throw new NoSuchElementException();
+                }
+                GenericRow result = next;
+                next = null;
+                return result;
+            }
+            try {
+                if (blockRemaining == 0) {
+                    throw new NoSuchElementException();
+                }
+                blockRecordIndex++;
+                GenericRow row =
+                        reuseRow == null
+                                ? new 
GenericRow(recordDecoder.projectedFieldCount)
+                                : reuseRow;
+                int recordStart = decoder.absolutePosition();
+                recordDecoder.read(decoder, row, null, null);
+                encodedRecord = decoder.borrowedView(recordStart, 
decoder.absolutePosition());
+                blockRemaining--;
+                return row;
+            } catch (IOException e) {
+                throw new UncheckedIOException(
+                        "Failed to decode projected Manifest Avro record.", e);
+            }
+        }
+
+        public long recordIndex() {
+            checkCurrentRecord();
+            return blockRecordIndex;
+        }
+
+        /** Returns a borrowed encoded view of the current complete Avro 
record. */
+        public ByteBuffer encodedRecord() {
+            checkCurrentRecord();
+            if (encodedRecord == null) {
+                throw new IllegalStateException("No current Manifest Avro 
record.");
+            }
+            return encodedRecord;
+        }
+
+        private void ensureBlockFullyConsumed() throws IOException {
+            if (blockRemaining == 0 && !decoder.isEnd()) {
+                throw new IOException("Manifest Avro block contains trailing 
undecoded bytes.");
+            }
+        }
+
+        private void checkCurrentRecord() {
+            if (blockRecordIndex < 0) {
+                throw new IllegalStateException("No current Manifest Avro 
record.");
+            }
+        }
+    }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java 
b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java
index 2400f8cc82..0564ff198e 100644
--- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java
+++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java
@@ -200,8 +200,8 @@ public class ManifestFile extends 
ObjectsFile<ManifestEntry> {
             @Nullable BucketFilter bucketFilter)
             throws IOException {
         try {
-            return new ManifestAvroReader(
-                    fileIO.newInputStream(path), projectedType, 
partitionFilter, bucketFilter);
+            ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path));
+            return reader.read(projectedType, partitionFilter, bucketFilter);
         } catch (IOException e) {
             FileUtils.checkExists(fileIO, path);
             throw e;
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 feff0e1420..e7a279f0f0 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
@@ -112,7 +112,8 @@ public final class ProjectedManifestEntry implements 
ManifestEntry {
                                                         DataFileMeta.LEVEL,
                                                         
DataFileMeta.EXTRA_FILES,
                                                         
DataFileMeta.EMBEDDED_FILE_INDEX,
-                                                        
DataFileMeta.EXTERNAL_PATH)))));
+                                                        
DataFileMeta.EXTERNAL_PATH,
+                                                        
DataFileMeta.FIRST_ROW_ID)))));
     }
 
     private static Projection createRowRangeProjection() {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java 
b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java
index c6376b7407..0846cb2697 100644
--- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java
@@ -45,6 +45,8 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
 import java.io.IOException;
+import java.io.UncheckedIOException;
+import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
@@ -122,18 +124,14 @@ public class ManifestFileTest {
         ManifestEntrySerializer serializer = new ManifestEntrySerializer();
         List<ManifestEntry> actual = new ArrayList<>();
 
-        try (ManifestAvroReader reader =
-                new ManifestAvroReader(
-                        fileIO.newInputStream(path),
-                        ManifestEntry.MANIFEST_ROW_TYPE,
-                        partitionFilter,
-                        bucketFilter)) {
-            while (reader.hasNext()) {
-                InternalRow row = reader.next();
+        try (ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path));
+                CloseableIterator<InternalRow> rows =
+                        reader.read(
+                                ManifestEntry.MANIFEST_ROW_TYPE, 
partitionFilter, bucketFilter)) {
+            while (rows.hasNext()) {
+                InternalRow row = rows.next();
                 actual.add(serializer.fromRow(row));
             }
-            assertThat(reader.decodedDataFiles()).isEqualTo(expected.size());
-            assertThat(reader.skippedDataFiles()).isEqualTo(entries.size() - 
expected.size());
         }
 
         assertThat(actual).containsExactlyElementsOf(expected);
@@ -169,11 +167,11 @@ public class ManifestFileTest {
         ProjectedManifestEntry projectedEntry =
                 
ProjectedManifestEntry.Projection.create(projectedType).createEntry();
 
-        try (ManifestAvroReader reader =
-                new ManifestAvroReader(fileIO.newInputStream(path), 
projectedType, null, null)) {
+        try (ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path));
+                CloseableIterator<InternalRow> rows = 
reader.read(projectedType, null, null)) {
             for (ManifestEntry expected : entries) {
-                assertThat(reader.hasNext()).isTrue();
-                InternalRow row = reader.next();
+                assertThat(rows.hasNext()).isTrue();
+                InternalRow row = rows.next();
                 assertThat(row.getFieldCount()).isEqualTo(3);
                 assertThat(row.getRow(0, 
projectedFileType.getFieldCount()).getFieldCount())
                         .isEqualTo(2);
@@ -184,9 +182,7 @@ public class ManifestFileTest {
                 
assertThat(projectedEntry.partition()).isEqualTo(expected.partition());
                 assertThat(projectedEntry.kind()).isEqualTo(expected.kind());
             }
-            assertThat(reader.hasNext()).isFalse();
-            assertThat(reader.decodedDataFiles()).isEqualTo(entries.size());
-            assertThat(reader.skippedDataFiles()).isZero();
+            assertThat(rows.hasNext()).isFalse();
         }
     }
 
@@ -200,11 +196,11 @@ public class ManifestFileTest {
 
         List<DataField> fields = ManifestEntry.MANIFEST_ROW_TYPE.getFields();
         RowType projectedType = new RowType(false, 
Arrays.asList(fields.get(2), fields.get(1)));
-        try (ManifestAvroReader reader =
-                new ManifestAvroReader(fileIO.newInputStream(path), 
projectedType, null, null)) {
+        try (ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path));
+                CloseableIterator<InternalRow> rows = 
reader.read(projectedType, null, null)) {
             for (ManifestEntry expected : entries) {
-                assertThat(reader.hasNext()).isTrue();
-                InternalRow row = reader.next();
+                assertThat(rows.hasNext()).isTrue();
+                InternalRow row = rows.next();
                 assertThat(row.getFieldCount()).isEqualTo(2);
                 assertThat(row.getBinary(0))
                         .containsExactly(
@@ -212,9 +208,28 @@ public class ManifestFileTest {
                                         expected.partition()));
                 
assertThat(FileKind.fromByteValue(row.getByte(1))).isEqualTo(expected.kind());
             }
-            assertThat(reader.hasNext()).isFalse();
-            assertThat(reader.decodedDataFiles()).isZero();
-            assertThat(reader.skippedDataFiles()).isEqualTo(entries.size());
+            assertThat(rows.hasNext()).isFalse();
+        }
+    }
+
+    @Test
+    void testAvroReaderRejectsTrailingUndecodedRecords() throws Exception {
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest =
+                writeSingleManifest(manifestFile, Arrays.asList(gen.next(), 
gen.next()));
+        Path path = new Path(new Path(tempDir.toUri()), "manifest/" + 
manifest.fileName());
+        lowerFirstBlockRecordCount(path);
+
+        try (ManifestAvroReader reader =
+                        new 
ManifestAvroReader(LocalFileIO.create().newInputStream(path));
+                CloseableIterator<InternalRow> rows =
+                        reader.read(ManifestEntry.MANIFEST_ROW_TYPE, null, 
null)) {
+            assertThat(rows.hasNext()).isTrue();
+            rows.next();
+            assertThatThrownBy(rows::hasNext)
+                    .isInstanceOf(UncheckedIOException.class)
+                    .hasRootCauseInstanceOf(IOException.class)
+                    .hasStackTraceContaining("trailing undecoded bytes");
         }
     }
 
@@ -246,7 +261,8 @@ public class ManifestFileTest {
         try (CloseableIterator<ProjectedManifestEntry> entries =
                 manifestFile.scan(
                         manifest.fileName(), 
ProjectedManifestEntry.DELETE_ENTRY_PROJECTION)) {
-            assertThatThrownBy(entries::hasNext)
+            assertThat(entries.hasNext()).isTrue();
+            assertThatThrownBy(entries::next)
                     .isInstanceOf(IllegalArgumentException.class)
                     .hasMessageContaining("not compatible");
         }
@@ -286,17 +302,32 @@ public class ManifestFileTest {
         }
 
         ManifestEntry actual;
-        try (ManifestAvroReader reader =
-                new ManifestAvroReader(
-                        fileIO.newInputStream(path), 
ManifestEntry.MANIFEST_ROW_TYPE, null, null)) {
-            assertThat(reader.hasNext()).isTrue();
-            actual = serializer.fromRow(reader.next());
-            assertThat(reader.hasNext()).isFalse();
+        try (ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path));
+                CloseableIterator<InternalRow> rows =
+                        reader.read(ManifestEntry.MANIFEST_ROW_TYPE, null, 
null)) {
+            assertThat(rows.hasNext()).isTrue();
+            actual = serializer.fromRow(rows.next());
+            assertThat(rows.hasNext()).isFalse();
         }
 
         assertThat(actual.fileName()).isEqualTo(source.fileName());
         assertThat(actual.file().firstRowId()).isNull();
         assertThat(actual.file().writeCols()).isNull();
+
+        ProjectedManifestEntry.Projection projection = 
ProjectedManifestEntry.ROW_RANGE_PROJECTION;
+        ProjectedManifestEntry binaryEntry = projection.createEntry();
+        try (ManifestAvroReader reader = new 
ManifestAvroReader(fileIO.newInputStream(path))) {
+            assertThat(reader.hasNext()).isTrue();
+            ManifestAvroReader.RawBlock block = reader.next();
+            assertThat(block.rawBlockCopySupported()).isFalse();
+            ManifestAvroReader.RowIterator rows = 
block.toRows(projection.projectedType());
+            assertThat(rows.hasNext()).isTrue();
+            binaryEntry.replace(rows.next());
+            assertThat(binaryEntry.rowCount()).isEqualTo(source.rowCount());
+            assertThat(binaryEntry.firstRowId()).isNull();
+            assertThat(rows.hasNext()).isFalse();
+            assertThat(reader.hasNext()).isFalse();
+        }
     }
 
     @Test
@@ -331,12 +362,15 @@ public class ManifestFileTest {
         }
 
         assertThatThrownBy(
-                        () ->
-                                new ManifestAvroReader(
-                                        fileIO.newInputStream(path),
-                                        ManifestEntry.MANIFEST_ROW_TYPE,
-                                        null,
-                                        null))
+                        () -> {
+                            try (ManifestAvroReader reader =
+                                            new 
ManifestAvroReader(fileIO.newInputStream(path));
+                                    CloseableIterator<InternalRow> rows =
+                                            reader.read(
+                                                    
ManifestEntry.MANIFEST_ROW_TYPE, null, null)) {
+                                rows.hasNext();
+                            }
+                        })
                 .isInstanceOf(IllegalArgumentException.class)
                 .hasMessageContaining("expected _KIND but found _FILE");
     }
@@ -392,7 +426,7 @@ public class ManifestFileTest {
     }
 
     @Test
-    void testReadDeletedEntriesWithProjectedScan() {
+    void testReadDeletedEntriesWithProjectedScan() throws Exception {
         ManifestEntry first = gen.next();
         ManifestEntry second = gen.next();
         DataFileMeta firstFile =
@@ -441,6 +475,14 @@ public class ManifestFileTest {
         ManifestFileMeta secondManifest =
                 writeSingleManifest(manifestFile, Arrays.asList(secondAdd, 
secondDelete));
 
+        try (CloseableIterator<ProjectedManifestEntry> entries =
+                manifestFile.scan(
+                        firstManifest.fileName(), 
ProjectedManifestEntry.DELETE_ENTRY_PROJECTION)) {
+            
assertThat(entries.next().file().nonNullFirstRowId()).isEqualTo(10L);
+            
assertThat(entries.next().file().nonNullFirstRowId()).isEqualTo(10L);
+            assertThat(entries.hasNext()).isFalse();
+        }
+
         Set<FileEntry.Identifier> deleted =
                 FileEntry.readDeletedEntries(
                         manifestFile, Arrays.asList(firstManifest, 
secondManifest), 2);
@@ -488,33 +530,172 @@ public class ManifestFileTest {
     }
 
     @Test
-    void testScanProjectedManifestEntriesCanBeRetained() throws Exception {
+    void testScanProjectedManifestCreatesDistinctEntryWrappers() throws 
Exception {
         List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());
         ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
         ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
         ProjectedManifestEntry.Projection projection =
                 projection(DataFileMeta.FILE_NAME, DataFileMeta.ROW_COUNT);
-        List<ProjectedManifestEntry> actual = new ArrayList<>();
+        List<String> fileNames = new ArrayList<>();
+        List<Long> rowCounts = new ArrayList<>();
+        ProjectedManifestEntry previous = null;
 
         try (CloseableIterator<ProjectedManifestEntry> iterator =
                 manifestFile.scan(manifest.fileName(), projection)) {
             while (iterator.hasNext()) {
-                actual.add(iterator.next());
+                ProjectedManifestEntry current = iterator.next();
+                assertThat(current).isNotSameAs(previous);
+                fileNames.add(current.fileName());
+                rowCounts.add(current.rowCount());
+                previous = current;
             }
         }
 
-        assertThat(actual).hasSize(entries.size());
-        for (int i = 1; i < actual.size(); i++) {
-            assertThat(actual.get(i)).isNotSameAs(actual.get(i - 1));
-        }
-        
assertThat(actual.stream().map(ManifestEntry::fileName).collect(Collectors.toList()))
+        assertThat(fileNames)
                 .containsExactlyElementsOf(
                         
entries.stream().map(ManifestEntry::fileName).collect(Collectors.toList()));
-        
assertThat(actual.stream().map(ManifestEntry::rowCount).collect(Collectors.toList()))
+        assertThat(rowCounts)
                 .containsExactlyElementsOf(
                         
entries.stream().map(ManifestEntry::rowCount).collect(Collectors.toList()));
     }
 
+    @Test
+    void testBlockReaderConvertsRawBlocksToProjectedRows() throws Exception {
+        List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
+        ProjectedManifestEntry.Projection projection =
+                projection(DataFileMeta.FILE_NAME, DataFileMeta.ROW_COUNT);
+        ProjectedManifestEntry actual = projection.createEntry();
+        InternalRow reusedRow = null;
+        InternalRow reusedFileRow = null;
+        int position = 0;
+
+        try (ManifestAvroReader reader = openManifestReader(manifest)) {
+            while (reader.hasNext()) {
+                ManifestAvroReader.RawBlock block = reader.next();
+                assertThat(block.rawBlockCopySupported()).isTrue();
+                ManifestAvroReader.RowIterator rows = 
block.toRows(projection.projectedType());
+                ByteBuffer reusedEncodedRecord = null;
+                while (rows.hasNext()) {
+                    GenericRow row = rows.next();
+                    ByteBuffer encodedRecord = rows.encodedRecord();
+                    assertThat(encodedRecord.remaining()).isPositive();
+                    if (reusedEncodedRecord != null) {
+                        
assertThat(encodedRecord).isSameAs(reusedEncodedRecord);
+                    }
+                    reusedEncodedRecord = encodedRecord;
+                    InternalRow fileRow = row.getRow(2, 2);
+                    if (reusedRow != null) {
+                        assertThat(row).isSameAs(reusedRow);
+                        assertThat(fileRow).isSameAs(reusedFileRow);
+                    }
+                    reusedRow = row;
+                    reusedFileRow = fileRow;
+                    actual.replace(row);
+                    ManifestEntry expected = entries.get(position++);
+                    assertThat(actual.kind()).isEqualTo(expected.kind());
+                    
assertThat(actual.partition()).isEqualTo(expected.partition());
+                    
assertThat(actual.fileName()).isEqualTo(expected.fileName());
+                    
assertThat(actual.rowCount()).isEqualTo(expected.rowCount());
+                }
+            }
+        }
+
+        assertThat(position).isEqualTo(entries.size());
+    }
+
+    @Test
+    void testBlockReaderReadsAcrossMultipleBlocks() throws Exception {
+        List<ManifestEntry> entries = Collections.nCopies(1_000, gen.next());
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
+        ProjectedManifestEntry.Projection projection = 
projection(DataFileMeta.FILE_NAME);
+        int blockCount = 0;
+        int rowCount = 0;
+
+        try (ManifestAvroReader reader = openManifestReader(manifest)) {
+            while (reader.hasNext()) {
+                ManifestAvroReader.RowIterator rows =
+                        reader.next().toRows(projection.projectedType());
+                assertThat(rows.hasNext()).isTrue();
+                while (rows.hasNext()) {
+                    rows.next();
+                    rowCount++;
+                }
+                blockCount++;
+            }
+        }
+
+        assertThat(blockCount).isGreaterThan(1);
+        assertThat(rowCount).isEqualTo(entries.size());
+    }
+
+    @Test
+    void testBlockReaderSupportsReorderedProjection() throws Exception {
+        List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
+        RowType manifestType = ManifestEntry.MANIFEST_ROW_TYPE;
+        RowType fileType =
+                DataFileMeta.SCHEMA.project(DataFileMeta.ROW_COUNT, 
DataFileMeta.FILE_NAME);
+        RowType projectedType =
+                new RowType(
+                        false,
+                        Arrays.asList(
+                                
manifestType.getField(ManifestEntry.FILE).newType(fileType),
+                                manifestType.getField(ManifestEntry.PARTITION),
+                                manifestType.getField(ManifestEntry.KIND)));
+        ProjectedManifestEntry.Projection projection =
+                ProjectedManifestEntry.Projection.create(projectedType);
+        ProjectedManifestEntry actual = projection.createEntry();
+        int position = 0;
+
+        try (ManifestAvroReader reader = openManifestReader(manifest)) {
+            while (reader.hasNext()) {
+                ManifestAvroReader.RowIterator rows =
+                        reader.next().toRows(projection.projectedType());
+                while (rows.hasNext()) {
+                    InternalRow row = rows.next();
+                    assertThat(row.getRow(0, 2).getLong(0))
+                            .isEqualTo(entries.get(position).rowCount());
+                    actual.replace(row);
+                    
assertThat(actual.fileName()).isEqualTo(entries.get(position).fileName());
+                    
assertThat(actual.partition()).isEqualTo(entries.get(position).partition());
+                    
assertThat(actual.kind()).isEqualTo(entries.get(position).kind());
+                    position++;
+                }
+            }
+        }
+
+        assertThat(position).isEqualTo(entries.size());
+    }
+
+    @Test
+    void testBlockReaderSupportsFullManifestProjection() throws Exception {
+        List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
+        ProjectedManifestEntry.Projection projection = 
ProjectedManifestEntry.fullProjection();
+        ProjectedManifestEntry binaryEntry = projection.createEntry();
+        ManifestEntrySerializer serializer = new ManifestEntrySerializer();
+        int position = 0;
+
+        try (ManifestAvroReader reader = openManifestReader(manifest)) {
+            while (reader.hasNext()) {
+                ManifestAvroReader.RowIterator rows =
+                        reader.next().toRows(projection.projectedType());
+                while (rows.hasNext()) {
+                    binaryEntry.replace(rows.next());
+                    assertThat(serializer.fromRow(binaryEntry.fullRow()))
+                            .isEqualTo(entries.get(position++));
+                }
+            }
+        }
+
+        assertThat(position).isEqualTo(entries.size());
+    }
+
     @Test
     void testScanProjectedManifestKeepsEntryValidWhenAdvancing() throws 
Exception {
         List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next());
@@ -540,19 +721,18 @@ public class ManifestFileTest {
         List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());
         ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
         ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
-        List<ProjectedManifestEntry> retained = new ArrayList<>();
+        List<String> processedFileNames = new ArrayList<>();
 
         try (CloseableIterator<ProjectedManifestEntry> iterator =
                 manifestFile.scan(manifest.fileName(), 
projection(DataFileMeta.FILE_NAME))) {
             while (iterator.hasNext()) {
                 ProjectedManifestEntry entry = iterator.next();
-                retained.add(entry);
+                processedFileNames.add(entry.fileName());
                 break;
             }
         }
 
-        assertThat(retained).hasSize(1);
-        
assertThat(retained.get(0).fileName()).isEqualTo(entries.get(0).fileName());
+        
assertThat(processedFileNames).containsExactly(entries.get(0).fileName());
     }
 
     @Test
@@ -561,7 +741,7 @@ public class ManifestFileTest {
         ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
         ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
         RuntimeException failure = new RuntimeException("Expected processing 
failure.");
-        List<ProjectedManifestEntry> retained = new ArrayList<>();
+        List<String> processedFileNames = new ArrayList<>();
 
         assertThatThrownBy(
                         () -> {
@@ -571,13 +751,12 @@ public class ManifestFileTest {
                                             
projection(DataFileMeta.FILE_NAME))) {
                                 assertThat(iterator.hasNext()).isTrue();
                                 ProjectedManifestEntry entry = iterator.next();
-                                retained.add(entry);
+                                processedFileNames.add(entry.fileName());
                                 throw failure;
                             }
                         })
                 .isSameAs(failure);
-        assertThat(retained).hasSize(1);
-        
assertThat(retained.get(0).fileName()).isEqualTo(entries.get(0).fileName());
+        
assertThat(processedFileNames).containsExactly(entries.get(0).fileName());
     }
 
     private List<ManifestEntry> generateData() {
@@ -588,6 +767,38 @@ public class ManifestFileTest {
         return entries;
     }
 
+    private ManifestAvroReader openManifestReader(ManifestFileMeta manifest) 
throws IOException {
+        FileIO fileIO = LocalFileIO.create();
+        Path path = new Path(new Path(tempDir.toUri()), "manifest/" + 
manifest.fileName());
+        return new ManifestAvroReader(fileIO.newInputStream(path));
+    }
+
+    private void lowerFirstBlockRecordCount(Path path) throws IOException {
+        java.nio.file.Path localPath = java.nio.file.Paths.get(path.toUri());
+        byte[] bytes = java.nio.file.Files.readAllBytes(localPath);
+        byte[] syncMarker = Arrays.copyOfRange(bytes, bytes.length - 16, 
bytes.length);
+        int headerSyncPosition = indexOf(bytes, syncMarker, 4, bytes.length - 
syncMarker.length);
+        assertThat(headerSyncPosition).isGreaterThanOrEqualTo(0);
+
+        int blockCountPosition = headerSyncPosition + syncMarker.length;
+        assertThat(bytes[blockCountPosition]).isEqualTo((byte) 4);
+        bytes[blockCountPosition] = 2;
+        java.nio.file.Files.write(localPath, bytes);
+    }
+
+    private static int indexOf(byte[] bytes, byte[] target, int from, int 
limit) {
+        for (int position = from; position + target.length <= limit; 
position++) {
+            int index = 0;
+            while (index < target.length && bytes[position + index] == 
target[index]) {
+                index++;
+            }
+            if (index == target.length) {
+                return position;
+            }
+        }
+        return -1;
+    }
+
     private ManifestFile createManifestFile(String pathStr) {
         return createManifestFile(pathStr, 
ThreadLocalRandom.current().nextInt(8192) + 1024);
     }
diff --git a/paimon-format/src/main/java/org/apache/avro/file/RawBlock.java 
b/paimon-format/src/main/java/org/apache/avro/file/RawBlock.java
new file mode 100644
index 0000000000..3386074b9d
--- /dev/null
+++ b/paimon-format/src/main/java/org/apache/avro/file/RawBlock.java
@@ -0,0 +1,168 @@
+/*
+ * 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.avro.file;
+
+import org.apache.avro.Schema;
+import org.apache.avro.io.DatumReader;
+import org.apache.avro.io.Decoder;
+
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.nio.ByteBuffer;
+
+/** Reusable compressed Avro block. */
+public final class RawBlock {
+
+    private DataFileStream.DataBlock block;
+    private Codec codec;
+    private Schema schema;
+    private boolean decompressed;
+    private ByteBuffer decompressedBuffer;
+
+    RawBlock(DataFileStream.DataBlock block, Codec codec, Schema schema) {
+        replace(block, codec, schema);
+    }
+
+    RawBlock replace(DataFileStream.DataBlock block, Codec codec, Schema 
schema) {
+        this.block = block;
+        this.codec = codec;
+        this.schema = schema;
+        this.decompressed = false;
+        this.decompressedBuffer = null;
+        return this;
+    }
+
+    DataFileStream.DataBlock dataBlock() {
+        return block;
+    }
+
+    public long recordCount() {
+        return block.getNumEntries();
+    }
+
+    public ByteBuffer decompress(ByteBuffer reuse) throws IOException {
+        if (!decompressed) {
+            if (codec instanceof ZstandardCodec) {
+                ByteBuffer source = block.getAsByteBuffer();
+                ByteBuffer target =
+                        reuse != null && reuse.hasArray() ? reuse : 
ByteBuffer.allocate(256 * 1024);
+                int size = 0;
+                try (InputStream compressed =
+                                new ByteArrayInputStream(
+                                        source.array(),
+                                        source.arrayOffset() + 
source.position(),
+                                        source.remaining());
+                        InputStream input = ZstandardLoader.input(compressed, 
true)) {
+                    while (true) {
+                        if (size == target.capacity()) {
+                            int grownCapacity =
+                                    target.capacity() == 0
+                                            ? 256 * 1024
+                                            : 
Math.multiplyExact(target.capacity(), 2);
+                            ByteBuffer grown = 
ByteBuffer.allocate(grownCapacity);
+                            System.arraycopy(
+                                    target.array(),
+                                    target.arrayOffset(),
+                                    grown.array(),
+                                    grown.arrayOffset(),
+                                    size);
+                            target = grown;
+                        }
+                        int read =
+                                input.read(
+                                        target.array(),
+                                        target.arrayOffset() + size,
+                                        target.capacity() - size);
+                        if (read < 0) {
+                            break;
+                        }
+                        size += read;
+                    }
+                }
+                target.position(0);
+                target.limit(size);
+                decompressedBuffer = target.duplicate();
+            } else {
+                block.decompressUsing(codec);
+                decompressedBuffer = block.getAsByteBuffer().duplicate();
+            }
+            decompressed = true;
+        }
+        return decompressedBuffer.duplicate();
+    }
+
+    /** Returns a single-block stream for appending this compressed block to 
an Avro writer. */
+    public <D> DataFileStream<D> asStream() throws IOException {
+        if (decompressed) {
+            throw new IllegalStateException("A decompressed Avro block cannot 
be copied raw.");
+        }
+        return new SingleBlockStream<D>(schema, codec, block);
+    }
+
+    private static final class SingleBlockStream<D> extends DataFileStream<D> {
+
+        private final Schema schema;
+        private final Codec codec;
+        private DataBlock block;
+
+        private SingleBlockStream(Schema schema, Codec codec, DataBlock block) 
throws IOException {
+            super(new NoOpDatumReader<D>());
+            this.schema = schema;
+            this.codec = codec;
+            this.block = block;
+        }
+
+        @Override
+        public Schema getSchema() {
+            return schema;
+        }
+
+        @Override
+        Codec resolveCodec() {
+            return codec;
+        }
+
+        @Override
+        boolean hasNextBlock() {
+            return block != null;
+        }
+
+        @Override
+        DataBlock nextRawBlock(DataBlock reuse) {
+            DataBlock result = block;
+            block = null;
+            return result;
+        }
+
+        @Override
+        public void close() {}
+    }
+
+    private static final class NoOpDatumReader<D> implements DatumReader<D> {
+
+        @Override
+        public void setSchema(Schema schema) {}
+
+        @Override
+        public D read(D reuse, Decoder decoder) {
+            throw new UnsupportedOperationException();
+        }
+    }
+}
diff --git 
a/paimon-format/src/main/java/org/apache/avro/file/RawBlockReader.java 
b/paimon-format/src/main/java/org/apache/avro/file/RawBlockReader.java
new file mode 100644
index 0000000000..43a68c54a4
--- /dev/null
+++ b/paimon-format/src/main/java/org/apache/avro/file/RawBlockReader.java
@@ -0,0 +1,56 @@
+/*
+ * 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.avro.file;
+
+import org.apache.avro.Schema;
+import org.apache.avro.io.DatumReader;
+import org.apache.avro.io.Decoder;
+
+import java.io.IOException;
+import java.io.InputStream;
+
+/** Package bridge exposing Avro's compressed blocks without reflection. */
+public final class RawBlockReader extends DataFileStream<Void> {
+
+    public RawBlockReader(InputStream input) throws IOException {
+        super(input, new NoOpDatumReader<Void>());
+    }
+
+    public boolean hasNextRawBlock() {
+        return super.hasNextBlock();
+    }
+
+    public RawBlock nextRawBlock(RawBlock reuse) throws IOException {
+        DataBlock raw = super.nextRawBlock(reuse == null ? null : 
reuse.dataBlock());
+        return reuse == null
+                ? new RawBlock(raw, resolveCodec(), getSchema())
+                : reuse.replace(raw, resolveCodec(), getSchema());
+    }
+
+    private static final class NoOpDatumReader<D> implements DatumReader<D> {
+
+        @Override
+        public void setSchema(Schema schema) {}
+
+        @Override
+        public D read(D reuse, Decoder decoder) {
+            throw new UnsupportedOperationException();
+        }
+    }
+}
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroBlockReader.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroBlockReader.java
index ef98b15fa0..c79b0e1985 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroBlockReader.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroBlockReader.java
@@ -18,88 +18,76 @@
 
 package org.apache.paimon.format.avro;
 
+import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.IOUtils;
 
 import org.apache.avro.AvroRuntimeException;
-import org.apache.avro.Schema;
-import org.apache.avro.file.DataFileStream;
-import org.apache.avro.generic.GenericDatumReader;
+import org.apache.avro.file.RawBlock;
+import org.apache.avro.file.RawBlockReader;
+
+import javax.annotation.Nullable;
 
 import java.io.Closeable;
 import java.io.IOException;
 import java.io.InputStream;
-import java.nio.ByteBuffer;
+import java.util.Collections;
 
 /**
- * Reader which exposes decompressed blocks from an Avro object container file.
+ * Reader which exposes raw blocks from an Avro object container file.
  *
  * <p>This reader owns the input stream and closes it when construction fails 
or {@link #close()} is
  * called.
  */
 public final class AvroBlockReader implements Closeable {
 
-    private final DataFileStream<Object> reader;
+    private final RawBlockReader reader;
 
-    private long currentBlockRecordCount = -1;
+    private @Nullable AvroRawBlock borrowedRawBlock;
 
     public AvroBlockReader(InputStream input) throws IOException {
         try {
-            this.reader = new DataFileStream<>(input, new 
GenericDatumReader<>());
+            this.reader = new RawBlockReader(input);
         } catch (IOException | RuntimeException | Error e) {
             IOUtils.closeQuietly(input);
             throw e;
         }
     }
 
-    Schema schema() {
-        return reader.getSchema();
-    }
-
     /** Creates a record decoder from the writer schema stored in the Avro 
file header. */
     public AvroRecordDecoder createRecordDecoder() {
         return new AvroRecordDecoder(reader.getSchema());
     }
 
-    /** Returns whether another block is available. */
-    public boolean hasNextBlock() throws IOException {
-        return replaceAvroRuntimeException(reader::hasNext);
+    /** Returns whether blocks use the default Avro schema for the given row 
type. */
+    public boolean supportsRawBlockCopy(RowType rowType) {
+        return AvroSchemaConverter.convertToSchema(rowType, 
Collections.emptyMap())
+                .equals(reader.getSchema());
     }
 
-    /**
-     * Returns a decompressed copy of the next block.
-     *
-     * <p>The returned array is owned by the caller and remains valid after 
this reader advances or
-     * closes.
-     */
-    public byte[] nextBlock() throws IOException {
-        BorrowedBlock block = nextBorrowedBlock();
-        byte[] bytes = new byte[block.length];
-        System.arraycopy(block.bytes, block.offset, bytes, 0, block.length);
-        return bytes;
+    /** Returns whether another block is available. */
+    public boolean hasNextBlock() throws IOException {
+        return replaceAvroRuntimeException(reader::hasNextRawBlock);
     }
 
     /**
-     * Returns a borrowed view of the next decompressed block.
+     * Returns a borrowed view of the next compressed block.
      *
-     * <p>The returned bytes are owned by this reader and may be overwritten 
by the next call to
-     * {@link #hasNextBlock()}, {@link #nextBlock()}, or this method, or when 
this reader is closed.
+     * <p>The returned holder and its storage are owned by this reader and 
reused by the next call
+     * to this method. Consume the block before advancing this reader.
      */
-    public BorrowedBlock nextBorrowedBlock() throws IOException {
-        ByteBuffer block = replaceAvroRuntimeException(reader::nextBlock);
-        currentBlockRecordCount = reader.getBlockCount();
-        return new BorrowedBlock(
-                block.array(),
-                block.arrayOffset() + block.position(),
-                block.remaining(),
-                currentBlockRecordCount);
-    }
-
-    /** Returns the record count of the last block returned by a block-reading 
method. */
-    public long currentBlockRecordCount() {
-        if (currentBlockRecordCount < 0) {
-            throw new IllegalStateException("No block has been read.");
-        }
-        return currentBlockRecordCount;
+    public AvroRawBlock nextBorrowedRawBlock() throws IOException {
+        RawBlock block =
+                replaceAvroRuntimeException(
+                        () ->
+                                reader.nextRawBlock(
+                                        borrowedRawBlock == null
+                                                ? null
+                                                : 
borrowedRawBlock.rawBlock()));
+        borrowedRawBlock =
+                borrowedRawBlock == null
+                        ? new AvroRawBlock(block)
+                        : borrowedRawBlock.replace(block);
+        return borrowedRawBlock;
     }
 
     @Override
@@ -118,38 +106,6 @@ public final class AvroBlockReader implements Closeable {
         }
     }
 
-    /** Borrowed decompressed block contents and its record count. */
-    public static final class BorrowedBlock {
-
-        private final byte[] bytes;
-        private final int offset;
-        private final int length;
-        private final long recordCount;
-
-        private BorrowedBlock(byte[] bytes, int offset, int length, long 
recordCount) {
-            this.bytes = bytes;
-            this.offset = offset;
-            this.length = length;
-            this.recordCount = recordCount;
-        }
-
-        public byte[] bytes() {
-            return bytes;
-        }
-
-        public int offset() {
-            return offset;
-        }
-
-        public int length() {
-            return length;
-        }
-
-        public long recordCount() {
-            return recordCount;
-        }
-    }
-
     @FunctionalInterface
     private interface IOSupplier<T> {
 
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRawBlock.java 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRawBlock.java
new file mode 100644
index 0000000000..c7d21d4cc8
--- /dev/null
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRawBlock.java
@@ -0,0 +1,64 @@
+/*
+ * 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.format.avro;
+
+import org.apache.avro.file.DataFileStream;
+import org.apache.avro.file.RawBlock;
+
+import javax.annotation.Nullable;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+
+/** Reusable compressed block from an Avro object container file. */
+public final class AvroRawBlock {
+
+    private RawBlock block;
+
+    AvroRawBlock(RawBlock block) {
+        this.block = block;
+    }
+
+    AvroRawBlock replace(RawBlock block) {
+        this.block = block;
+        return this;
+    }
+
+    RawBlock rawBlock() {
+        return block;
+    }
+
+    public long recordCount() {
+        return block.recordCount();
+    }
+
+    /**
+     * Lazily decompresses this block, reusing the supplied heap buffer when 
possible.
+     *
+     * <p>The returned view remains owned by this block and is invalidated 
when this holder is
+     * reused for another block.
+     */
+    public ByteBuffer decompress(@Nullable ByteBuffer reuse) throws 
IOException {
+        return block.decompress(reuse);
+    }
+
+    <D> DataFileStream<D> asStream() throws IOException {
+        return block.asStream();
+    }
+}
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRecordDecoder.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRecordDecoder.java
index 78b4941897..a4dafc1296 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRecordDecoder.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/avro/AvroRecordDecoder.java
@@ -27,6 +27,7 @@ import org.apache.avro.io.DecoderFactory;
 import javax.annotation.Nullable;
 
 import java.io.IOException;
+import java.nio.ByteBuffer;
 
 /**
  * Decoder for sequentially reading records from Avro blocks without exposing 
Avro classes to
@@ -38,6 +39,9 @@ public final class AvroRecordDecoder {
     private final int recordBranch;
 
     private @Nullable BinaryDecoder decoder;
+    private @Nullable ByteBuffer borrowedView;
+    private int blockOffset;
+    private int blockLength;
 
     AvroRecordDecoder(Schema writerSchema) {
         if (writerSchema.getType() == Schema.Type.UNION) {
@@ -93,6 +97,11 @@ public final class AvroRecordDecoder {
     /** Reuses this decoder for another block. */
     public void reset(byte[] bytes, int offset, int length) {
         decoder = DecoderFactory.get().binaryDecoder(bytes, offset, length, 
decoder);
+        blockOffset = offset;
+        blockLength = length;
+        if (borrowedView == null || borrowedView.array() != bytes) {
+            borrowedView = ByteBuffer.wrap(bytes);
+        }
     }
 
     /** Returns whether a block has been supplied through {@link 
#reset(byte[], int, int)}. */
@@ -114,6 +123,34 @@ public final class AvroRecordDecoder {
         return decoder().readInt();
     }
 
+    /** Returns the byte position relative to the beginning of the current 
block. */
+    private int position() throws IOException {
+        return blockLength - decoder().inputStream().available();
+    }
+
+    /** Returns the absolute byte position in the current block's backing 
array. */
+    public int absolutePosition() throws IOException {
+        return blockOffset + position();
+    }
+
+    /** Returns a borrowed view of an absolute range in the current block's 
backing array. */
+    public ByteBuffer borrowedView(int start, int end) {
+        if (borrowedView == null) {
+            throw new IllegalStateException("No Avro block has been 
supplied.");
+        }
+        int blockEnd = blockOffset + blockLength;
+        if (start < blockOffset || end < start || end > blockEnd) {
+            throw new IllegalArgumentException(
+                    String.format(
+                            "Borrowed Avro byte range [%s, %s) is outside 
block range [%s, %s).",
+                            start, end, blockOffset, blockEnd));
+        }
+        borrowedView.clear();
+        borrowedView.position(start);
+        borrowedView.limit(end);
+        return borrowedView;
+    }
+
     public byte[] readBytes() throws IOException {
         return decoder().readBytes(null).array();
     }
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/avro/AvroFileFormatTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/avro/AvroFileFormatTest.java
index b8cb87f1c4..8c05c5b4bf 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/avro/AvroFileFormatTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/avro/AvroFileFormatTest.java
@@ -24,7 +24,6 @@ import org.apache.paimon.format.FileFormat;
 import org.apache.paimon.format.FileFormatFactory.FormatContext;
 import org.apache.paimon.format.FormatReaderContext;
 import org.apache.paimon.format.FormatWriter;
-import org.apache.paimon.format.avro.AvroBlockReader.BorrowedBlock;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.fs.PositionOutputStream;
@@ -50,6 +49,7 @@ import java.io.ByteArrayOutputStream;
 import java.io.File;
 import java.io.FileNotFoundException;
 import java.io.IOException;
+import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.NoSuchElementException;
@@ -144,59 +144,6 @@ public class AvroFileFormatTest {
         }
     }
 
-    @Test
-    void testReadDecompressedBlocks() throws IOException {
-        RowType rowType = DataTypes.ROW(DataTypes.INT().notNull()).notNull();
-        LocalFileIO fileIO = LocalFileIO.create();
-        Path file = new Path(new Path(tempPath.toUri()), 
UUID.randomUUID().toString());
-        int numRecords = 100_000;
-
-        try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
-            FormatWriter writer = 
fileFormat.createWriterFactory(rowType).create(out, "zstd");
-            for (int i = 0; i < numRecords; i++) {
-                writer.addElement(GenericRow.of(i));
-            }
-            writer.close();
-        }
-
-        int nextValue = 0;
-        int numBlocks = 0;
-        byte[] firstBlock = null;
-        byte[] firstBlockCopy = null;
-        try (AvroBlockReader reader = new 
AvroBlockReader(fileIO.newInputStream(file))) {
-            assertThat(reader.schema()).isNotNull();
-            assertThatThrownBy(reader::currentBlockRecordCount)
-                    .isInstanceOf(IllegalStateException.class)
-                    .hasMessageContaining("No block");
-
-            AvroRowDatumReader datumReader = new AvroRowDatumReader(rowType);
-            datumReader.setSchema(reader.schema());
-            BinaryDecoder decoder = null;
-            while (reader.hasNextBlock()) {
-                byte[] block = reader.nextBlock();
-                if (firstBlock == null) {
-                    firstBlock = block;
-                    firstBlockCopy = block.clone();
-                }
-                numBlocks++;
-                decoder = DecoderFactory.get().binaryDecoder(block, decoder);
-                long blockRecordCount = reader.currentBlockRecordCount();
-                assertThat(blockRecordCount).isPositive();
-                for (long i = 0; i < blockRecordCount; i++) {
-                    assertThat(datumReader.read(null, 
decoder).getInt(0)).isEqualTo(nextValue++);
-                }
-                assertThat(decoder.isEnd()).isTrue();
-            }
-
-            assertThat(reader.hasNextBlock()).isFalse();
-            
assertThatThrownBy(reader::nextBlock).isInstanceOf(NoSuchElementException.class);
-        }
-
-        assertThat(numBlocks).isGreaterThan(1);
-        assertThat(nextValue).isEqualTo(numRecords);
-        assertThat(firstBlock).containsExactly(firstBlockCopy);
-    }
-
     @Test
     void testReadBlocksFromEmptyFile() throws IOException {
         RowType rowType = DataTypes.ROW(DataTypes.INT().notNull()).notNull();
@@ -209,12 +156,13 @@ public class AvroFileFormatTest {
 
         try (AvroBlockReader reader = new 
AvroBlockReader(fileIO.newInputStream(file))) {
             assertThat(reader.hasNextBlock()).isFalse();
-            
assertThatThrownBy(reader::nextBlock).isInstanceOf(NoSuchElementException.class);
+            assertThatThrownBy(reader::nextBorrowedRawBlock)
+                    .isInstanceOf(NoSuchElementException.class);
         }
     }
 
     @Test
-    void testReadBorrowedBlocks() throws IOException {
+    void testReadBorrowedRawBlocks() throws IOException {
         RowType rowType = DataTypes.ROW(DataTypes.INT().notNull()).notNull();
         LocalFileIO fileIO = LocalFileIO.create();
         Path file = new Path(new Path(tempPath.toUri()), 
UUID.randomUUID().toString());
@@ -228,27 +176,25 @@ public class AvroFileFormatTest {
             writer.close();
         }
 
-        int nextValue = 0;
+        long records = 0;
+        int blocks = 0;
+        AvroRawBlock previous = null;
         try (AvroBlockReader reader = new 
AvroBlockReader(fileIO.newInputStream(file))) {
-            AvroRecordDecoder decoder = reader.createRecordDecoder();
-            AvroRecordDecoder.FieldDecoder fieldDecoder =
-                    decoder.createFieldDecoder(0, rowType.getTypeAt(0));
             while (reader.hasNextBlock()) {
-                BorrowedBlock block = reader.nextBorrowedBlock();
-                assertThat(block.recordCount()).isPositive();
-                
assertThat(reader.currentBlockRecordCount()).isEqualTo(block.recordCount());
-                decoder.reset(block.bytes(), block.offset(), block.length());
-                for (long i = 0; i < block.recordCount(); i++) {
-                    assertThat(decoder.readRecordStart()).isTrue();
-                    assertThat(fieldDecoder.read(decoder, 
null)).isEqualTo(nextValue++);
+                AvroRawBlock block = reader.nextBorrowedRawBlock();
+                if (previous != null) {
+                    assertThat(block).isSameAs(previous);
                 }
-                assertThat(decoder.isEnd()).isTrue();
+                previous = block;
+                records += block.recordCount();
+                blocks++;
             }
-            assertThatThrownBy(reader::nextBorrowedBlock)
+            assertThatThrownBy(reader::nextBorrowedRawBlock)
                     .isInstanceOf(NoSuchElementException.class);
         }
 
-        assertThat(nextValue).isEqualTo(numRecords);
+        assertThat(blocks).isGreaterThan(1);
+        assertThat(records).isEqualTo(numRecords);
     }
 
     @Test
@@ -283,6 +229,35 @@ public class AvroFileFormatTest {
         assertThat(decoder.isEnd()).isTrue();
     }
 
+    @Test
+    void testReadsLargeZstdBlock() throws IOException {
+        RowType rowType =
+                RowType.builder()
+                        .field("payload", 
DataTypes.VARBINARY(500_000).notNull())
+                        .field("id", DataTypes.INT().notNull())
+                        .build();
+        AvroFileFormat format = new AvroFileFormat(new FormatContext(new 
Options(), 1024, 1024));
+        LocalFileIO fileIO = LocalFileIO.create();
+        Path file = new Path(new Path(tempPath.toUri()), 
UUID.randomUUID().toString());
+        byte[] payload = new byte[400_000];
+        Arrays.fill(payload, (byte) 7);
+
+        try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
+            FormatWriter writer = 
format.createWriterFactory(rowType).create(out, "zstd");
+            writer.addElement(GenericRow.of(payload, 42));
+            writer.close();
+        }
+
+        try (AvroBlockReader blockReader = new 
AvroBlockReader(fileIO.newInputStream(file))) {
+            AvroRawBlock block = blockReader.nextBorrowedRawBlock();
+            assertThat(block.recordCount()).isEqualTo(1);
+            ByteBuffer decoded = block.decompress(null);
+            assertThat(decoded.remaining()).isGreaterThan(payload.length);
+            assertThat(decoded.get(10)).isEqualTo((byte) 7);
+            
assertThat(block.decompress(ByteBuffer.allocate(1))).isEqualTo(decoded);
+        }
+    }
+
     @Test
     void testGetRealIOException() throws IOException {
         RowType rowType = DataTypes.ROW(DataTypes.INT().notNull());

Reply via email to