This is an automated email from the ASF dual-hosted git repository.
leaves12138 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 c35ef892ec Revert "[format] Support writing and reading Arrow schema
metadata for file formats" (#8346)
c35ef892ec is described below
commit c35ef892ecbe40d81e714367baae8a5984f4d4a4
Author: lxy <[email protected]>
AuthorDate: Wed Jun 24 19:05:01 2026 +0800
Revert "[format] Support writing and reading Arrow schema metadata for file
formats" (#8346)
Reverts apache/paimon#8321 to avoid introducing an Arrow dependency into
format.
---
.../paimon/format/SupportsWriterMetadata.java | 32 -----
paimon-format/pom.xml | 12 --
.../apache/paimon/format/FormatMetadataUtils.java | 137 -------------------
.../paimon/format/SupportsReaderArrowSchema.java | 30 ----
.../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, 16 insertions(+), 727 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
b/paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
deleted file mode 100644
index 89677a86e7..0000000000
---
a/paimon-common/src/main/java/org/apache/paimon/format/SupportsWriterMetadata.java
+++ /dev/null
@@ -1,32 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.paimon.format;
-
-import java.util.Map;
-
-/** Writer capability for adding format-specific file metadata before closing
the file. */
-public interface SupportsWriterMetadata {
-
- /**
- * 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 037e488168..b4136797b8 100644
--- a/paimon-format/pom.xml
+++ b/paimon-format/pom.xml
@@ -47,18 +47,6 @@ 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
deleted file mode 100644
index d5013bf8c8..0000000000
---
a/paimon-format/src/main/java/org/apache/paimon/format/FormatMetadataUtils.java
+++ /dev/null
@@ -1,137 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.paimon.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/SupportsReaderArrowSchema.java
b/paimon-format/src/main/java/org/apache/paimon/format/SupportsReaderArrowSchema.java
deleted file mode 100644
index 43e525ddd9..0000000000
---
a/paimon-format/src/main/java/org/apache/paimon/format/SupportsReaderArrowSchema.java
+++ /dev/null
@@ -1,30 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.paimon.format;
-
-import org.apache.arrow.vector.types.pojo.Schema;
-
-import java.io.IOException;
-import java.util.Optional;
-
-/** Reader capability for formats that can recover Arrow schema from file
metadata. */
-public interface SupportsReaderArrowSchema {
-
- 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 fd291685fa..b1de74242a 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,10 +24,8 @@ 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;
@@ -41,7 +39,6 @@ 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;
@@ -57,9 +54,7 @@ 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;
@@ -109,7 +104,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
Pool<OrcReaderBatch> poolOfBatches =
createPoolOfBatches(context.filePath(), poolSize,
context.fileIO());
- OrcRecordReader orcReader =
+ RecordReader orcReader =
createRecordReader(
hadoopConfig,
schema,
@@ -229,14 +224,12 @@ 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>,
SupportsReaderArrowSchema {
+ private static final class OrcVectorizedReader implements
FileRecordReader<InternalRow> {
- private final OrcRecordReader orcReader;
+ private final RecordReader orcReader;
private final Pool<OrcReaderBatch> pool;
- private OrcVectorizedReader(
- final OrcRecordReader orcReader, final Pool<OrcReaderBatch>
pool) {
+ private OrcVectorizedReader(final RecordReader orcReader, final
Pool<OrcReaderBatch> pool) {
this.orcReader = checkNotNull(orcReader, "orcReader");
this.pool = checkNotNull(pool, "pool");
}
@@ -247,8 +240,8 @@ public class OrcReaderFactory implements
FormatReaderFactory {
final OrcReaderBatch batch = getCachedEntry();
final VectorizedRowBatch orcVectorBatch =
batch.orcVectorizedRowBatch();
- long rowNumber = orcReader.recordReader.getRowNumber();
- if (!nextBatch(orcReader.recordReader, orcVectorBatch)) {
+ long rowNumber = orcReader.getRowNumber();
+ if (!nextBatch(orcReader, orcVectorBatch)) {
batch.recycle();
return null;
}
@@ -256,31 +249,9 @@ 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 {
- try {
- orcReader.recordReader.close();
- } finally {
- orcReader.fileReader.close();
- }
+ orcReader.close();
}
private OrcReaderBatch getCachedEntry() throws IOException {
@@ -293,18 +264,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
}
}
- 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(
+ private static RecordReader createRecordReader(
org.apache.hadoop.conf.Configuration conf,
TypeDescription schema,
List<OrcFilters.Predicate> conjunctPredicates,
@@ -354,7 +314,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
// assign ids
schema.getId();
- return new OrcRecordReader(orcReader, orcRowsReader);
+ return 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 f4cdbaa33f..c44e3f26d6 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,9 +20,7 @@ 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;
@@ -30,27 +28,19 @@ 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, SupportsWriterMetadata {
+public class OrcBulkWriter implements FormatWriter {
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,
@@ -64,16 +54,6 @@ public class OrcBulkWriter implements FormatWriter,
SupportsWriterMetadata {
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
@@ -95,14 +75,7 @@ public class OrcBulkWriter implements FormatWriter,
SupportsWriterMetadata {
@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 0d729cac7e..282805897a 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,8 +34,6 @@ 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 {
@@ -59,12 +57,8 @@ public class ParquetWriterFactory implements
FormatWriterFactory, SupportsVarian
compression = null;
}
- // 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);
+ final ParquetWriter<InternalRow> writer =
writerBuilder.createWriter(out, compression);
+ return new ParquetBulkWriter(writer);
}
@Override
@@ -79,11 +73,7 @@ public class ParquetWriterFactory implements
FormatWriterFactory, SupportsVarian
ParquetBuilder<InternalRow> newBuilder =
((RowDataParquetBuilder) writerBuilder)
.withShreddingSchemas(inferredShreddingSchema);
- // 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);
+ final ParquetWriter<InternalRow> writer = newBuilder.createWriter(out,
compression);
+ return new ParquetBulkWriter(writer);
}
}
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 0470275386..af63a7d9f9 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,8 +20,6 @@ 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;
@@ -43,7 +41,6 @@ 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;
@@ -51,8 +48,7 @@ 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>, SupportsReaderArrowSchema {
+public class VectorizedParquetRecordReader implements
FileRecordReader<InternalRow> {
private ParquetFileReader reader;
@@ -273,16 +269,6 @@ public class VectorizedParquetRecordReader
}
}
- @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 e1dc534969..808febd9a2 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,8 +23,6 @@ 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}.
@@ -36,11 +34,4 @@ 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 42d17d5e82..d7282f699f 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,7 +20,6 @@ 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;
@@ -28,15 +27,11 @@ 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, SupportsWriterMetadata
{
+public class ParquetBulkWriter implements FormatWriter {
/** The ParquetWriter to write to. */
private final ParquetWriter<InternalRow> parquetWriter;
@@ -44,32 +39,13 @@ public class ParquetBulkWriter implements FormatWriter,
SupportsWriterMetadata {
/** 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
@@ -81,7 +57,6 @@ public class ParquetBulkWriter implements FormatWriter,
SupportsWriterMetadata {
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 60fd1097fe..14970e548e 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,14 +19,12 @@
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;
@@ -34,8 +32,6 @@ 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;
@@ -45,20 +41,12 @@ 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
@@ -101,11 +89,5 @@ 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 e87367f9ef..2e84df0932 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,9 +35,6 @@ 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> {
@@ -61,16 +58,8 @@ 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
deleted file mode 100644
index 7e6f44a644..0000000000
---
a/paimon-format/src/test/java/org/apache/paimon/format/FormatMetadataUtilsTest.java
+++ /dev/null
@@ -1,152 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.paimon.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 74980503e6..fb625c68da 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,42 +24,27 @@ 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 {
@@ -96,76 +81,6 @@ 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 97f64413ac..137da29aa6 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,46 +20,29 @@ 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}. */
@@ -75,98 +58,6 @@ 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 {