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());