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 a5eb81e126 [format] Support writing and reading Arrow schema metadata
for file formats (#8321)
a5eb81e126 is described below
commit a5eb81e1269af1f6fa24a4ac07a5a3d21c89530a
Author: lxy <[email protected]>
AuthorDate: Wed Jun 24 17:49:21 2026 +0800
[format] Support writing and reading Arrow schema metadata for file formats
(#8321)
This PR is a sub-PR for shared-shredding.
Shared-shredding needs to attach dictionary and other field-level
metadata before closing data files. To support that flow, this PR adds a
generic metadata write/read path for file formats, so upper layers can
provide raw key-value metadata during writing, and readers can parse the
stored metadata back when opening files.
The metadata representation follows Arrow Parquet's existing schema
metadata convention: Arrow stores the original serialized schema under
the `ARROW:schema` file metadata key, base64-decodes it on read, and
deserializes it with Arrow IPC schema reading. See Apache Arrow's
Parquet schema implementation, where `kArrowSchemaKey` is `ARROW:schema`
and the value is base64-decoded before `ReadSchema`.
Related design:
https://cwiki.apache.org/confluence/display/PAIMON/PIP-43%3A+Columnar+Storage+Optimization+for+MAP+Type+in+Paimon
---
.../paimon/format/SupportsWriterMetadata.java | 25 ++--
paimon-format/pom.xml | 12 ++
.../apache/paimon/format/FormatMetadataUtils.java | 137 +++++++++++++++++++
...Builder.java => SupportsReaderArrowSchema.java} | 19 +--
.../apache/paimon/format/orc/OrcReaderFactory.java | 58 ++++++--
.../paimon/format/orc/writer/OrcBulkWriter.java | 29 +++-
.../format/parquet/ParquetWriterFactory.java | 18 ++-
.../reader/VectorizedParquetRecordReader.java | 16 ++-
.../format/parquet/writer/ParquetBuilder.java | 9 ++
.../format/parquet/writer/ParquetBulkWriter.java | 27 +++-
.../parquet/writer/ParquetRowDataBuilder.java | 18 +++
.../parquet/writer/RowDataParquetBuilder.java | 11 ++
.../paimon/format/FormatMetadataUtilsTest.java | 152 +++++++++++++++++++++
.../paimon/format/orc/OrcFormatReadWriteTest.java | 85 ++++++++++++
.../format/parquet/ParquetFormatReadWriteTest.java | 109 +++++++++++++++
15 files changed, 681 insertions(+), 44 deletions(-)
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
b/paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
similarity index 57%
copy from
paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
copy to
paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
index 808febd9a2..89677a86e7 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
+++
b/paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
@@ -16,22 +16,17 @@
* limitations under the License.
*/
-package org.apache.paimon.format.parquet.writer;
+package org.apache.paimon.format;
-import org.apache.parquet.hadoop.ParquetWriter;
-import org.apache.parquet.io.OutputFile;
+import java.util.Map;
-import java.io.IOException;
-import java.io.Serializable;
+/** Writer capability for adding format-specific file metadata before closing
the file. */
+public interface SupportsWriterMetadata {
-/**
- * A builder to create a {@link ParquetWriter} from a Parquet {@link
OutputFile}.
- *
- * @param <T> The type of elements written by the writer.
- */
-@FunctionalInterface
-public interface ParquetBuilder<T> extends Serializable {
-
- /** Creates and configures a parquet writer to the given output file. */
- ParquetWriter<T> createWriter(OutputFile out, String compression) throws
IOException;
+ /**
+ * Adds raw metadata entries to the file footer.
+ *
+ * <p>This method must be called before {@link FormatWriter#close()}.
+ */
+ void addMetadata(Map<String, byte[]> metadata);
}
diff --git a/paimon-format/pom.xml b/paimon-format/pom.xml
index b4136797b8..037e488168 100644
--- a/paimon-format/pom.xml
+++ b/paimon-format/pom.xml
@@ -47,6 +47,18 @@ under the License.
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.paimon</groupId>
+ <artifactId>paimon-arrow</artifactId>
+ <version>${project.version}</version>
+ </dependency>
+
+ <dependency>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-vector</artifactId>
+ <version>${arrow.version}</version>
+ </dependency>
+
<dependency>
<groupId>org.xerial.snappy</groupId>
<artifactId>snappy-java</artifactId>
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/FormatMetadataUtils.java
b/paimon-format/src/main/java/org/apache/paimon/format/FormatMetadataUtils.java
new file mode 100644
index 0000000000..d5013bf8c8
--- /dev/null
+++
b/paimon-format/src/main/java/org/apache/paimon/format/FormatMetadataUtils.java
@@ -0,0 +1,137 @@
+/*
+ * 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;
+
+import org.apache.paimon.arrow.ArrowUtils;
+import org.apache.paimon.types.RowType;
+
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
+
+import javax.annotation.Nullable;
+
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.util.Base64;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.stream.Collectors;
+
+/** Utilities for format metadata encoded at file boundaries. */
+public class FormatMetadataUtils {
+
+ public static final String ARROW_SCHEMA_METADATA_KEY = "ARROW:schema";
+
+ private FormatMetadataUtils() {}
+
+ public static Map<String, String> encodeMetadata(Map<String, byte[]>
metadata) {
+ Map<String, String> encoded = new LinkedHashMap<>();
+ for (Map.Entry<String, byte[]> entry : metadata.entrySet()) {
+ encoded.put(entry.getKey(),
Base64.getEncoder().encodeToString(entry.getValue()));
+ }
+ return encoded;
+ }
+
+ /**
+ * Decodes base64-encoded metadata values. Values that are not valid
base64 are returned as
+ * UTF-8 bytes.
+ */
+ public static Map<String, byte[]> decodeMetadata(Map<String, String>
metadata) {
+ Map<String, byte[]> decoded = new LinkedHashMap<>();
+ for (Map.Entry<String, String> entry : metadata.entrySet()) {
+ try {
+ decoded.put(entry.getKey(),
Base64.getDecoder().decode(entry.getValue()));
+ } catch (IllegalArgumentException e) {
+ decoded.put(entry.getKey(),
entry.getValue().getBytes(StandardCharsets.UTF_8));
+ }
+ }
+ return decoded;
+ }
+
+ public static Optional<Schema> readArrowSchema(@Nullable String
encodedSchema) {
+ if (encodedSchema == null) {
+ return Optional.empty();
+ }
+ try {
+ byte[] schemaBytes = Base64.getDecoder().decode(encodedSchema);
+ return
Optional.of(Schema.deserializeMessage(ByteBuffer.wrap(schemaBytes)));
+ } catch (RuntimeException e) {
+ return Optional.empty();
+ }
+ }
+
+ public static byte[] serializeArrowSchema(Schema arrowSchema) {
+ return arrowSchema.serializeAsMessage();
+ }
+
+ /**
+ * Builds an Arrow schema from a Paimon row type and injects metadata into
top-level fields.
+ *
+ * <p>The keys of {@code fieldMetadata} are top-level field names. Nested
fields are converted
+ * from the {@link RowType} but do not receive metadata from this map. If
injected metadata
+ * conflicts with metadata produced during Arrow conversion, the Arrow
conversion metadata wins
+ * to preserve format-specific field information.
+ */
+ public static Schema buildArrowSchema(
+ RowType rowType, Map<String, Map<String, String>> fieldMetadata) {
+ List<Field> fields =
+ rowType.getFields().stream()
+ .map(
+ field ->
+ withMetadata(
+ ArrowUtils.toArrowField(
+ field.name(),
field.id(), field.type(), 0),
+
fieldMetadata.get(field.name())))
+ .collect(Collectors.toList());
+ return new Schema(fields);
+ }
+
+ /** Returns metadata for top-level Arrow fields only. */
+ public static Map<String, Map<String, String>> readFieldMetadata(Schema
arrowSchema) {
+ Map<String, Map<String, String>> result = new LinkedHashMap<>();
+ for (Field field : arrowSchema.getFields()) {
+ result.put(
+ field.getName(),
+ Collections.unmodifiableMap(new
LinkedHashMap<>(field.getMetadata())));
+ }
+ return Collections.unmodifiableMap(result);
+ }
+
+ private static Field withMetadata(Field field, @Nullable Map<String,
String> metadata) {
+ if (metadata == null || metadata.isEmpty()) {
+ return field;
+ }
+ FieldType fieldType = field.getFieldType();
+ Map<String, String> result = new LinkedHashMap<>();
+ result.putAll(metadata);
+ result.putAll(fieldType.getMetadata());
+ return new Field(
+ field.getName(),
+ new FieldType(
+ fieldType.isNullable(),
+ fieldType.getType(),
+ fieldType.getDictionary(),
+ result),
+ field.getChildren());
+ }
+}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
b/paimon-format/src/main/java/org/apache/paimon/format/SupportsReaderArrowSchema.java
similarity index 59%
copy from
paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
copy to
paimon-format/src/main/java/org/apache/paimon/format/SupportsReaderArrowSchema.java
index 808febd9a2..43e525ddd9 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/SupportsReaderArrowSchema.java
@@ -16,22 +16,15 @@
* limitations under the License.
*/
-package org.apache.paimon.format.parquet.writer;
+package org.apache.paimon.format;
-import org.apache.parquet.hadoop.ParquetWriter;
-import org.apache.parquet.io.OutputFile;
+import org.apache.arrow.vector.types.pojo.Schema;
import java.io.IOException;
-import java.io.Serializable;
+import java.util.Optional;
-/**
- * A builder to create a {@link ParquetWriter} from a Parquet {@link
OutputFile}.
- *
- * @param <T> The type of elements written by the writer.
- */
-@FunctionalInterface
-public interface ParquetBuilder<T> extends Serializable {
+/** Reader capability for formats that can recover Arrow schema from file
metadata. */
+public interface SupportsReaderArrowSchema {
- /** Creates and configures a parquet writer to the given output file. */
- ParquetWriter<T> createWriter(OutputFile out, String compression) throws
IOException;
+ Optional<Schema> readArrowSchema() throws IOException;
}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
b/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
index b1de74242a..fd291685fa 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
@@ -24,8 +24,10 @@ import org.apache.paimon.data.columnar.ColumnarRow;
import org.apache.paimon.data.columnar.ColumnarRowIterator;
import org.apache.paimon.data.columnar.VectorizedColumnBatch;
import org.apache.paimon.data.columnar.VectorizedRowIterator;
+import org.apache.paimon.format.FormatMetadataUtils;
import org.apache.paimon.format.FormatReaderFactory;
import org.apache.paimon.format.OrcFormatReaderContext;
+import org.apache.paimon.format.SupportsReaderArrowSchema;
import org.apache.paimon.format.fs.HadoopReadOnlyFileSystem;
import org.apache.paimon.format.orc.filter.OrcFilters;
import org.apache.paimon.fs.FileIO;
@@ -39,6 +41,7 @@ import org.apache.paimon.utils.Pair;
import org.apache.paimon.utils.Pool;
import org.apache.paimon.utils.RoaringBitmap32;
+import org.apache.arrow.vector.types.pojo.Schema;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch;
import org.apache.hadoop.hive.ql.io.sarg.SearchArgument;
@@ -54,7 +57,9 @@ import org.apache.orc.impl.RecordReaderImpl;
import javax.annotation.Nullable;
import java.io.IOException;
+import java.nio.charset.StandardCharsets;
import java.util.List;
+import java.util.Optional;
import static org.apache.paimon.format.orc.OrcTypeUtil.convertToOrcSchema;
import static
org.apache.paimon.format.orc.reader.AbstractOrcColumnVector.createPaimonVector;
@@ -104,7 +109,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
Pool<OrcReaderBatch> poolOfBatches =
createPoolOfBatches(context.filePath(), poolSize,
context.fileIO());
- RecordReader orcReader =
+ OrcRecordReader orcReader =
createRecordReader(
hadoopConfig,
schema,
@@ -224,12 +229,14 @@ public class OrcReaderFactory implements
FormatReaderFactory {
* batch is addressed by the starting row number of the batch, plus the
number of records to be
* skipped before.
*/
- private static final class OrcVectorizedReader implements
FileRecordReader<InternalRow> {
+ private static final class OrcVectorizedReader
+ implements FileRecordReader<InternalRow>,
SupportsReaderArrowSchema {
- private final RecordReader orcReader;
+ private final OrcRecordReader orcReader;
private final Pool<OrcReaderBatch> pool;
- private OrcVectorizedReader(final RecordReader orcReader, final
Pool<OrcReaderBatch> pool) {
+ private OrcVectorizedReader(
+ final OrcRecordReader orcReader, final Pool<OrcReaderBatch>
pool) {
this.orcReader = checkNotNull(orcReader, "orcReader");
this.pool = checkNotNull(pool, "pool");
}
@@ -240,8 +247,8 @@ public class OrcReaderFactory implements
FormatReaderFactory {
final OrcReaderBatch batch = getCachedEntry();
final VectorizedRowBatch orcVectorBatch =
batch.orcVectorizedRowBatch();
- long rowNumber = orcReader.getRowNumber();
- if (!nextBatch(orcReader, orcVectorBatch)) {
+ long rowNumber = orcReader.recordReader.getRowNumber();
+ if (!nextBatch(orcReader.recordReader, orcVectorBatch)) {
batch.recycle();
return null;
}
@@ -249,9 +256,31 @@ public class OrcReaderFactory implements
FormatReaderFactory {
return batch.convertAndGetIterator(orcVectorBatch, rowNumber);
}
+ @Override
+ public Optional<Schema> readArrowSchema() {
+ org.apache.orc.Reader fileReader = orcReader.fileReader;
+ if (!fileReader
+ .getMetadataKeys()
+ .contains(FormatMetadataUtils.ARROW_SCHEMA_METADATA_KEY)) {
+ return Optional.empty();
+ }
+ return FormatMetadataUtils.readArrowSchema(
+ StandardCharsets.UTF_8
+ .decode(
+ fileReader
+ .getMetadataValue(
+
FormatMetadataUtils.ARROW_SCHEMA_METADATA_KEY)
+ .duplicate())
+ .toString());
+ }
+
@Override
public void close() throws IOException {
- orcReader.close();
+ try {
+ orcReader.recordReader.close();
+ } finally {
+ orcReader.fileReader.close();
+ }
}
private OrcReaderBatch getCachedEntry() throws IOException {
@@ -264,7 +293,18 @@ public class OrcReaderFactory implements
FormatReaderFactory {
}
}
- private static RecordReader createRecordReader(
+ private static final class OrcRecordReader {
+
+ private final org.apache.orc.Reader fileReader;
+ private final RecordReader recordReader;
+
+ private OrcRecordReader(org.apache.orc.Reader fileReader, RecordReader
recordReader) {
+ this.fileReader = fileReader;
+ this.recordReader = recordReader;
+ }
+ }
+
+ private static OrcRecordReader createRecordReader(
org.apache.hadoop.conf.Configuration conf,
TypeDescription schema,
List<OrcFilters.Predicate> conjunctPredicates,
@@ -314,7 +354,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
// assign ids
schema.getId();
- return orcRowsReader;
+ return new OrcRecordReader(orcReader, orcRowsReader);
} catch (IOException e) {
// exception happened, we need to close the reader
IOUtils.closeQuietly(orcReader);
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/orc/writer/OrcBulkWriter.java
b/paimon-format/src/main/java/org/apache/paimon/format/orc/writer/OrcBulkWriter.java
index c44e3f26d6..f4cdbaa33f 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/orc/writer/OrcBulkWriter.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/orc/writer/OrcBulkWriter.java
@@ -20,7 +20,9 @@ package org.apache.paimon.format.orc.writer;
import org.apache.paimon.annotation.VisibleForTesting;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.format.FormatMetadataUtils;
import org.apache.paimon.format.FormatWriter;
+import org.apache.paimon.format.SupportsWriterMetadata;
import org.apache.paimon.fs.PositionOutputStream;
import org.apache.paimon.options.MemorySize;
@@ -28,19 +30,27 @@ import
org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch;
import org.apache.orc.Writer;
import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
import static org.apache.paimon.utils.Preconditions.checkNotNull;
+import static org.apache.paimon.utils.Preconditions.checkState;
/** A {@link FormatWriter} implementation that writes data in ORC format. */
-public class OrcBulkWriter implements FormatWriter {
+public class OrcBulkWriter implements FormatWriter, SupportsWriterMetadata {
private final Writer writer;
private final Vectorizer<InternalRow> vectorizer;
private final VectorizedRowBatch rowBatch;
private final PositionOutputStream underlyingStream;
+ private final Map<String, byte[]> metadata;
private long currentBatchMemoryUsage = 0;
private final long memoryLimit;
+ private boolean closed = false;
public OrcBulkWriter(
Vectorizer<InternalRow> vectorizer,
@@ -54,6 +64,16 @@ public class OrcBulkWriter implements FormatWriter {
this.rowBatch = vectorizer.getSchema().createRowBatch(batchSize);
this.underlyingStream = underlyingStream;
this.memoryLimit = memoryLimit.getBytes();
+ this.metadata = new HashMap<>();
+ }
+
+ @Override
+ public void addMetadata(Map<String, byte[]> metadata) {
+ checkState(!closed, "Cannot add metadata after writer is closed.");
+ for (Map.Entry<String, byte[]> entry : metadata.entrySet()) {
+ this.metadata.put(
+ entry.getKey(), Arrays.copyOf(entry.getValue(),
entry.getValue().length));
+ }
}
@Override
@@ -75,7 +95,14 @@ public class OrcBulkWriter implements FormatWriter {
@Override
public void close() throws IOException {
flush();
+ for (Map.Entry<String, String> entry :
+ FormatMetadataUtils.encodeMetadata(metadata).entrySet()) {
+ writer.addUserMetadata(
+ entry.getKey(),
+
ByteBuffer.wrap(entry.getValue().getBytes(StandardCharsets.UTF_8)));
+ }
writer.close();
+ this.closed = true;
}
@Override
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetWriterFactory.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetWriterFactory.java
index 282805897a..0d729cac7e 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetWriterFactory.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetWriterFactory.java
@@ -34,6 +34,8 @@ import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.io.OutputFile;
import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
/** A factory that creates a Parquet {@link FormatWriter}. */
public class ParquetWriterFactory implements FormatWriterFactory,
SupportsVariantInference {
@@ -57,8 +59,12 @@ public class ParquetWriterFactory implements
FormatWriterFactory, SupportsVarian
compression = null;
}
- final ParquetWriter<InternalRow> writer =
writerBuilder.createWriter(out, compression);
- return new ParquetBulkWriter(writer);
+ // Keep this exact map instance shared by ParquetBulkWriter and
WriteSupport. The writer
+ // collects metadata before close, and WriteSupport reads it when
finalizing the footer.
+ Map<String, byte[]> metadata = new HashMap<>();
+ final ParquetWriter<InternalRow> writer =
+ writerBuilder.createWriter(out, compression, () -> metadata);
+ return new ParquetBulkWriter(writer, metadata);
}
@Override
@@ -73,7 +79,11 @@ public class ParquetWriterFactory implements
FormatWriterFactory, SupportsVarian
ParquetBuilder<InternalRow> newBuilder =
((RowDataParquetBuilder) writerBuilder)
.withShreddingSchemas(inferredShreddingSchema);
- final ParquetWriter<InternalRow> writer = newBuilder.createWriter(out,
compression);
- return new ParquetBulkWriter(writer);
+ // Keep this exact map instance shared by ParquetBulkWriter and
WriteSupport. The writer
+ // collects metadata before close, and WriteSupport reads it when
finalizing the footer.
+ Map<String, byte[]> metadata = new HashMap<>();
+ final ParquetWriter<InternalRow> writer =
+ newBuilder.createWriter(out, compression, () -> metadata);
+ return new ParquetBulkWriter(writer, metadata);
}
}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
index af63a7d9f9..0470275386 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
@@ -20,6 +20,8 @@ package org.apache.paimon.format.parquet.reader;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.columnar.writable.WritableColumnVector;
+import org.apache.paimon.format.FormatMetadataUtils;
+import org.apache.paimon.format.SupportsReaderArrowSchema;
import org.apache.paimon.format.parquet.type.ParquetField;
import org.apache.paimon.format.parquet.type.ParquetPrimitiveField;
import org.apache.paimon.fs.FileIO;
@@ -41,6 +43,7 @@ import java.io.IOException;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
+import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
@@ -48,7 +51,8 @@ import static java.lang.String.format;
import static
org.apache.paimon.format.parquet.reader.ParquetReaderUtil.createReadableColumnVectors;
/** Record reader for parquet. */
-public class VectorizedParquetRecordReader implements
FileRecordReader<InternalRow> {
+public class VectorizedParquetRecordReader
+ implements FileRecordReader<InternalRow>, SupportsReaderArrowSchema {
private ParquetFileReader reader;
@@ -269,6 +273,16 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
}
}
+ @Override
+ public Optional<org.apache.arrow.vector.types.pojo.Schema>
readArrowSchema()
+ throws IOException {
+ return FormatMetadataUtils.readArrowSchema(
+ reader.getFooter()
+ .getFileMetaData()
+ .getKeyValueMetaData()
+ .get(FormatMetadataUtils.ARROW_SCHEMA_METADATA_KEY));
+ }
+
@Override
public void close() throws IOException {
if (reader != null) {
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
index 808febd9a2..e1dc534969 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBuilder.java
@@ -23,6 +23,8 @@ import org.apache.parquet.io.OutputFile;
import java.io.IOException;
import java.io.Serializable;
+import java.util.Map;
+import java.util.function.Supplier;
/**
* A builder to create a {@link ParquetWriter} from a Parquet {@link
OutputFile}.
@@ -34,4 +36,11 @@ public interface ParquetBuilder<T> extends Serializable {
/** Creates and configures a parquet writer to the given output file. */
ParquetWriter<T> createWriter(OutputFile out, String compression) throws
IOException;
+
+ default ParquetWriter<T> createWriter(
+ OutputFile out, String compression, Supplier<Map<String, byte[]>>
metadataSupplier)
+ throws IOException {
+ throw new UnsupportedOperationException(
+ "This ParquetBuilder does not support writer metadata.");
+ }
}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBulkWriter.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBulkWriter.java
index d7282f699f..42d17d5e82 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBulkWriter.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetBulkWriter.java
@@ -20,6 +20,7 @@ package org.apache.paimon.format.parquet.writer;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.format.FormatWriter;
+import org.apache.paimon.format.SupportsWriterMetadata;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.metadata.ParquetMetadata;
@@ -27,11 +28,15 @@ import org.apache.parquet.hadoop.metadata.ParquetMetadata;
import javax.annotation.Nullable;
import java.io.IOException;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
import static org.apache.paimon.utils.Preconditions.checkNotNull;
+import static org.apache.paimon.utils.Preconditions.checkState;
/** A simple {@link FormatWriter} implementation that wraps a {@link
ParquetWriter}. */
-public class ParquetBulkWriter implements FormatWriter {
+public class ParquetBulkWriter implements FormatWriter, SupportsWriterMetadata
{
/** The ParquetWriter to write to. */
private final ParquetWriter<InternalRow> parquetWriter;
@@ -39,13 +44,32 @@ public class ParquetBulkWriter implements FormatWriter {
/** Cached footer metadata after close, used to avoid re-reading the file
for stats. */
@Nullable private ParquetMetadata footerMetadata;
+ private final Map<String, byte[]> metadata;
+
+ private boolean closed = false;
+
/**
* Creates a new ParquetBulkWriter wrapping the given ParquetWriter.
*
* @param parquetWriter The ParquetWriter to write to.
*/
public ParquetBulkWriter(ParquetWriter<InternalRow> parquetWriter) {
+ this(parquetWriter, new HashMap<>());
+ }
+
+ public ParquetBulkWriter(
+ ParquetWriter<InternalRow> parquetWriter, Map<String, byte[]>
metadata) {
this.parquetWriter = checkNotNull(parquetWriter, "parquetWriter");
+ this.metadata = checkNotNull(metadata, "metadata");
+ }
+
+ @Override
+ public void addMetadata(Map<String, byte[]> metadata) {
+ checkState(!closed, "Cannot add metadata after writer is closed.");
+ for (Map.Entry<String, byte[]> entry : metadata.entrySet()) {
+ this.metadata.put(
+ entry.getKey(), Arrays.copyOf(entry.getValue(),
entry.getValue().length));
+ }
}
@Override
@@ -57,6 +81,7 @@ public class ParquetBulkWriter implements FormatWriter {
public void close() throws IOException {
parquetWriter.close();
this.footerMetadata = parquetWriter.getFooter();
+ this.closed = true;
}
@Override
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetRowDataBuilder.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetRowDataBuilder.java
index 14970e548e..60fd1097fe 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetRowDataBuilder.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/ParquetRowDataBuilder.java
@@ -19,12 +19,14 @@
package org.apache.paimon.format.parquet.writer;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.format.FormatMetadataUtils;
import org.apache.paimon.format.parquet.VariantUtils;
import org.apache.paimon.types.RowType;
import org.apache.hadoop.conf.Configuration;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.api.WriteSupport;
+import org.apache.parquet.hadoop.api.WriteSupport.FinalizedWriteContext;
import org.apache.parquet.io.OutputFile;
import org.apache.parquet.io.api.RecordConsumer;
import org.apache.parquet.schema.MessageType;
@@ -32,6 +34,8 @@ import org.apache.parquet.schema.MessageType;
import javax.annotation.Nullable;
import java.util.HashMap;
+import java.util.Map;
+import java.util.function.Supplier;
import static
org.apache.paimon.format.parquet.ParquetSchemaConverter.convertToParquetMessageType;
@@ -41,12 +45,20 @@ public class ParquetRowDataBuilder
private final RowType rowType;
@Nullable private final RowType shreddingSchemas;
+ private Supplier<Map<String, byte[]>> metadataSupplier;
public ParquetRowDataBuilder(
OutputFile path, RowType rowType, @Nullable RowType
shreddingSchemas) {
super(path);
this.rowType = rowType;
this.shreddingSchemas = shreddingSchemas;
+ this.metadataSupplier = HashMap::new;
+ }
+
+ public ParquetRowDataBuilder withMetadataSupplier(
+ Supplier<Map<String, byte[]>> metadataSupplier) {
+ this.metadataSupplier = metadataSupplier;
+ return this;
}
@Override
@@ -89,5 +101,11 @@ public class ParquetRowDataBuilder
public void write(InternalRow record) {
this.writer.write(record);
}
+
+ @Override
+ public FinalizedWriteContext finalizeWrite() {
+ return new FinalizedWriteContext(
+
FormatMetadataUtils.encodeMetadata(metadataSupplier.get()));
+ }
}
}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/RowDataParquetBuilder.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/RowDataParquetBuilder.java
index 2e84df0932..e87367f9ef 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/RowDataParquetBuilder.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/writer/RowDataParquetBuilder.java
@@ -35,6 +35,9 @@ import org.apache.parquet.io.OutputFile;
import javax.annotation.Nullable;
import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.function.Supplier;
/** A {@link ParquetBuilder} for {@link InternalRow}. */
public class RowDataParquetBuilder implements ParquetBuilder<InternalRow> {
@@ -58,8 +61,16 @@ public class RowDataParquetBuilder implements
ParquetBuilder<InternalRow> {
@Override
public ParquetWriter<InternalRow> createWriter(OutputFile out, String
compression)
throws IOException {
+ return createWriter(out, compression, HashMap::new);
+ }
+
+ @Override
+ public ParquetWriter<InternalRow> createWriter(
+ OutputFile out, String compression, Supplier<Map<String, byte[]>>
metadataSupplier)
+ throws IOException {
ParquetRowDataBuilder builder =
new ParquetRowDataBuilder(out, rowType, shreddingSchemas)
+ .withMetadataSupplier(metadataSupplier)
.withConf(conf)
.withCompressionCodec(getCompressionCodec(getCompression(compression)))
.withRowGroupSize(
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/FormatMetadataUtilsTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/FormatMetadataUtilsTest.java
new file mode 100644
index 0000000000..7e6f44a644
--- /dev/null
+++
b/paimon-format/src/test/java/org/apache/paimon/format/FormatMetadataUtilsTest.java
@@ -0,0 +1,152 @@
+/*
+ * 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;
+
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+
+import org.apache.arrow.vector.types.Types;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.junit.jupiter.api.Test;
+
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.Base64;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link FormatMetadataUtils}. */
+public class FormatMetadataUtilsTest {
+
+ @Test
+ public void testEncodeAndDecodeMetadata() {
+ Map<String, byte[]> metadata = new LinkedHashMap<>();
+ metadata.put("encoded",
"paimon-value".getBytes(StandardCharsets.UTF_8));
+
+ Map<String, String> encoded =
FormatMetadataUtils.encodeMetadata(metadata);
+ encoded.put("plain", "plain-value");
+ assertThat(encoded)
+ .containsEntry(
+ "encoded",
+ Base64.getEncoder()
+
.encodeToString("paimon-value".getBytes(StandardCharsets.UTF_8)));
+
+ Map<String, byte[]> decoded =
FormatMetadataUtils.decodeMetadata(encoded);
+
+ assertThat(new String(decoded.get("encoded"), StandardCharsets.UTF_8))
+ .isEqualTo("paimon-value");
+ assertThat(new String(decoded.get("plain"), StandardCharsets.UTF_8))
+ .isEqualTo("plain-value");
+ }
+
+ @Test
+ public void testReadArrowSchema() {
+ Map<String, String> fieldMetadata = new LinkedHashMap<>();
+ fieldMetadata.put("paimon.test.field-key", "field-value");
+ Schema schema =
+ new Schema(
+ Collections.singletonList(
+ new Field(
+ "field",
+ new FieldType(
+ true,
+
Types.MinorType.VARCHAR.getType(),
+ null,
+ fieldMetadata),
+ null)));
+ String encodedSchema =
+ Base64.getEncoder()
+
.encodeToString(FormatMetadataUtils.serializeArrowSchema(schema));
+
+
assertThat(FormatMetadataUtils.readArrowSchema(encodedSchema)).hasValue(schema);
+ assertThat(FormatMetadataUtils.readArrowSchema(null)).isEmpty();
+
assertThat(FormatMetadataUtils.readArrowSchema("not-base64")).isEmpty();
+ assertThat(
+ FormatMetadataUtils.readArrowSchema(
+ Base64.getEncoder()
+ .encodeToString(
+ "not-arrow-schema"
+
.getBytes(StandardCharsets.UTF_8))))
+ .isEmpty();
+ }
+
+ @Test
+ public void testReadFieldMetadata() {
+ assertThat(FormatMetadataUtils.readFieldMetadata(new
Schema(Collections.emptyList())))
+ .isEmpty();
+
+ Map<String, String> fieldMetadata = new LinkedHashMap<>();
+ fieldMetadata.put("paimon.test.field-key", "field-value");
+ Schema schema =
+ new Schema(
+ Arrays.asList(
+ new Field(
+ "with_metadata",
+ new FieldType(
+ true,
+ Types.MinorType.INT.getType(),
+ null,
+ fieldMetadata),
+ null),
+ new Field(
+ "without_metadata",
+ new FieldType(
+ true,
Types.MinorType.INT.getType(), null, null),
+ null)));
+
+ assertThat(FormatMetadataUtils.readFieldMetadata(schema))
+ .containsEntry("with_metadata", fieldMetadata)
+ .containsEntry("without_metadata", Collections.emptyMap());
+ }
+
+ @Test
+ public void testBuildArrowSchemaWithFieldMetadata() {
+ RowType rowType =
+ DataTypes.ROW(
+ DataTypes.FIELD(0, "id", DataTypes.INT()),
+ DataTypes.FIELD(
+ 1, "tags", DataTypes.MAP(DataTypes.STRING(),
DataTypes.INT())),
+ DataTypes.FIELD(
+ 2,
+ "nested",
+ DataTypes.ROW(
+ DataTypes.FIELD(3, "name",
DataTypes.STRING()),
+ DataTypes.FIELD(
+ 4, "scores",
DataTypes.ARRAY(DataTypes.INT())))));
+ Map<String, String> tagsMetadata = new LinkedHashMap<>();
+ tagsMetadata.put("paimon.test.tags", "enabled");
+
+ Map<String, Map<String, String>> fieldMetadata = new LinkedHashMap<>();
+ fieldMetadata.put("tags", tagsMetadata);
+
+ Schema schema = FormatMetadataUtils.buildArrowSchema(rowType,
fieldMetadata);
+
+ assertThat(schema.getFields())
+ .extracting(Field::getName)
+ .containsExactly("id", "tags", "nested");
+ assertThat(schema.findField("tags").getMetadata())
+ .containsEntry("paimon.test.tags", "enabled");
+
assertThat(schema.findField("nested").getMetadata()).doesNotContainKey("paimon.test.tags");
+ }
+}
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcFormatReadWriteTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcFormatReadWriteTest.java
index fb625c68da..74980503e6 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcFormatReadWriteTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcFormatReadWriteTest.java
@@ -24,27 +24,42 @@ import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.format.FileFormat;
import org.apache.paimon.format.FileFormatFactory;
+import org.apache.paimon.format.FormatMetadataUtils;
import org.apache.paimon.format.FormatReadWriteTest;
import org.apache.paimon.format.FormatReaderContext;
import org.apache.paimon.format.FormatWriter;
import org.apache.paimon.format.OrcOptions;
+import org.apache.paimon.format.SupportsReaderArrowSchema;
+import org.apache.paimon.format.SupportsWriterMetadata;
import org.apache.paimon.fs.PositionOutputStream;
import org.apache.paimon.options.Options;
+import org.apache.paimon.reader.FileRecordReader;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
+import org.apache.arrow.vector.types.Types;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
import org.junit.jupiter.api.Test;
import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
+import java.util.Map;
+import java.util.Optional;
import java.util.TimeZone;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** An orc {@link FormatReadWriteTest}. */
public class OrcFormatReadWriteTest extends FormatReadWriteTest {
@@ -81,6 +96,76 @@ public class OrcFormatReadWriteTest extends
FormatReadWriteTest {
super("orc");
}
+ @Test
+ public void testWriteMetadata() throws IOException {
+ RowType rowType =
+ DataTypes.ROW(
+ DataTypes.FIELD(0, "id", DataTypes.INT()),
+ DataTypes.FIELD(1, "name", DataTypes.STRING()));
+
+ PositionOutputStream out = fileIO.newOutputStream(file, false);
+ FormatWriter writer =
newFormat.createWriterFactory(rowType).create(out, "zstd");
+ Map<String, String> fieldMetadata = new HashMap<>();
+ fieldMetadata.put("paimon.test.field-key", "field-value");
+ fieldMetadata.put("paimon.test.field-version", "1");
+ Schema arrowSchema =
+ new Schema(
+ Arrays.asList(
+ new Field(
+ "id",
+ new FieldType(
+ true,
+ Types.MinorType.INT.getType(),
+ null,
+
Collections.singletonMap("PARQUET:field_id", "0")),
+ null),
+ new Field(
+ "name",
+ new FieldType(
+ true,
+
Types.MinorType.VARCHAR.getType(),
+ null,
+ fieldMetadata),
+ null)));
+ byte[] arrowSchemaBytes = arrowSchema.serializeAsMessage();
+ Map<String, byte[]> metadata = new HashMap<>();
+ metadata.put("paimon.test.key",
"paimon-test-value".getBytes(StandardCharsets.UTF_8));
+ metadata.put(FormatMetadataUtils.ARROW_SCHEMA_METADATA_KEY,
arrowSchemaBytes);
+ ((SupportsWriterMetadata) writer).addMetadata(metadata);
+ writer.addElement(GenericRow.of(1,
org.apache.paimon.data.BinaryString.fromString("one")));
+ writer.close();
+ assertThatThrownBy(() -> ((SupportsWriterMetadata)
writer).addMetadata(metadata))
+ .isInstanceOf(IllegalStateException.class);
+ out.close();
+
+ try (org.apache.orc.Reader reader =
+ OrcReaderFactory.createReader(
+ new org.apache.hadoop.conf.Configuration(false),
fileIO, file, null)) {
+ ByteBuffer value = reader.getMetadataValue("paimon.test.key");
+ Map<String, byte[]> decodedMetadata =
+ FormatMetadataUtils.decodeMetadata(
+ Collections.singletonMap(
+ "paimon.test.key",
+
StandardCharsets.UTF_8.decode(value.duplicate()).toString()));
+ assertThat(new String(decodedMetadata.get("paimon.test.key"),
StandardCharsets.UTF_8))
+ .isEqualTo("paimon-test-value");
+ }
+
+ FormatReaderContext context =
+ new FormatReaderContext(fileIO, file,
fileIO.getFileSize(file));
+ RowType emptyRowType = new RowType(Collections.emptyList());
+ try (FileRecordReader<InternalRow> reader =
+ newFormat
+ .createReaderFactory(emptyRowType, emptyRowType,
Collections.emptyList())
+ .createReader(context)) {
+ Optional<Schema> readArrowSchema =
+ ((SupportsReaderArrowSchema) reader).readArrowSchema();
+ assertThat(readArrowSchema).hasValue(arrowSchema);
+
assertThat(FormatMetadataUtils.readFieldMetadata(readArrowSchema.get()))
+ .containsEntry("name", fieldMetadata);
+ }
+ }
+
@Test
public void testTimestampLTZWithLegacyWriteAndRead() throws IOException {
RowType rowType =
DataTypes.ROW(DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE());
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
index 137da29aa6..97f64413ac 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetFormatReadWriteTest.java
@@ -20,29 +20,46 @@ package org.apache.paimon.format.parquet;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
import org.apache.paimon.format.FileFormat;
import org.apache.paimon.format.FileFormatFactory;
+import org.apache.paimon.format.FormatMetadataUtils;
import org.apache.paimon.format.FormatReadWriteTest;
+import org.apache.paimon.format.FormatReaderContext;
import org.apache.paimon.format.FormatWriter;
+import org.apache.paimon.format.SupportsReaderArrowSchema;
+import org.apache.paimon.format.SupportsWriterMetadata;
+import org.apache.paimon.format.parquet.writer.ParquetBuilder;
import org.apache.paimon.fs.PositionOutputStream;
import org.apache.paimon.options.Options;
+import org.apache.paimon.reader.FileRecordReader;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
+import org.apache.arrow.vector.types.Types;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
import org.apache.parquet.column.values.bloomfilter.BloomFilter;
import org.apache.parquet.hadoop.ParquetFileReader;
+import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.metadata.BlockMetaData;
import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.hadoop.metadata.ParquetMetadata;
+import org.apache.parquet.io.OutputFile;
import org.assertj.core.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
import java.util.concurrent.ThreadLocalRandom;
/** A parquet {@link FormatReadWriteTest}. */
@@ -58,6 +75,98 @@ public class ParquetFormatReadWriteTest extends
FormatReadWriteTest {
new FileFormatFactory.FormatContext(new Options(), 1024,
1024));
}
+ @Test
+ public void testWriteMetadata() throws Exception {
+ ParquetFileFormat format =
+ new ParquetFileFormat(
+ new FileFormatFactory.FormatContext(new Options(),
1024, 1024));
+ RowType rowType =
+ DataTypes.ROW(
+ DataTypes.FIELD(0, "id", DataTypes.INT()),
+ DataTypes.FIELD(1, "name", DataTypes.STRING()));
+
+ PositionOutputStream out = fileIO.newOutputStream(file, false);
+ FormatWriter writer = format.createWriterFactory(rowType).create(out,
"zstd");
+ Map<String, String> fieldMetadata = new HashMap<>();
+ fieldMetadata.put("paimon.test.field-key", "field-value");
+ fieldMetadata.put("paimon.test.field-version", "1");
+ Schema arrowSchema =
+ new Schema(
+ Arrays.asList(
+ new Field(
+ "id",
+ new FieldType(
+ true,
+ Types.MinorType.INT.getType(),
+ null,
+
Collections.singletonMap("PARQUET:field_id", "0")),
+ null),
+ new Field(
+ "name",
+ new FieldType(
+ true,
+
Types.MinorType.VARCHAR.getType(),
+ null,
+ fieldMetadata),
+ null)));
+ byte[] arrowSchemaBytes = arrowSchema.serializeAsMessage();
+ Map<String, byte[]> metadata = new HashMap<>();
+ metadata.put("paimon.test.key",
"paimon-test-value".getBytes(StandardCharsets.UTF_8));
+ metadata.put(FormatMetadataUtils.ARROW_SCHEMA_METADATA_KEY,
arrowSchemaBytes);
+ ((SupportsWriterMetadata) writer).addMetadata(metadata);
+ writer.addElement(GenericRow.of(1, BinaryString.fromString("one")));
+ writer.close();
+ Assertions.assertThatThrownBy(() -> ((SupportsWriterMetadata)
writer).addMetadata(metadata))
+ .isInstanceOf(IllegalStateException.class);
+ out.close();
+
+ try (ParquetFileReader reader =
+ ParquetUtil.getParquetReader(
+ fileIO, file, fileIO.getFileSize(file), new
Options())) {
+ Map<String, String> fileMetadata =
+ reader.getFooter().getFileMetaData().getKeyValueMetaData();
+ Map<String, byte[]> decodedMetadata =
FormatMetadataUtils.decodeMetadata(fileMetadata);
+ Assertions.assertThat(
+ new String(
+ decodedMetadata.get("paimon.test.key"),
StandardCharsets.UTF_8))
+ .isEqualTo("paimon-test-value");
+ }
+
+ FormatReaderContext context =
+ new FormatReaderContext(fileIO, file,
fileIO.getFileSize(file));
+ RowType emptyRowType = new RowType(Collections.emptyList());
+ try (FileRecordReader<InternalRow> reader =
+ format.createReaderFactory(emptyRowType, emptyRowType,
Collections.emptyList())
+ .createReader(context)) {
+ Optional<Schema> readArrowSchema =
+ ((SupportsReaderArrowSchema) reader).readArrowSchema();
+ Assertions.assertThat(readArrowSchema).hasValue(arrowSchema);
+
Assertions.assertThat(FormatMetadataUtils.readFieldMetadata(readArrowSchema.get()))
+ .containsEntry("name", fieldMetadata);
+ }
+ }
+
+ @Test
+ public void testUnsupportedMetadataBuilderFailsExplicitly() throws
Exception {
+ ParquetBuilder<InternalRow> unsupportedMetadataBuilder =
+ new ParquetBuilder<InternalRow>() {
+ @Override
+ public ParquetWriter<InternalRow> createWriter(
+ OutputFile out, String compression) {
+ throw new AssertionError("Two-argument createWriter
should not be called.");
+ }
+ };
+ ParquetWriterFactory factory = new
ParquetWriterFactory(unsupportedMetadataBuilder);
+ PositionOutputStream out = fileIO.newOutputStream(file, false);
+ try {
+ Assertions.assertThatThrownBy(() -> factory.create(out, "zstd"))
+ .isInstanceOf(UnsupportedOperationException.class)
+ .hasMessageContaining("does not support writer metadata");
+ } finally {
+ out.close();
+ }
+ }
+
@ParameterizedTest
@ValueSource(booleans = {true, false})
public void testEnableBloomFilter(boolean enabled) throws Exception {