This is an automated email from the ASF dual-hosted git repository.
tianchen pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/arrow.git
The following commit(s) were added to refs/heads/master by this push:
new b196d1a ARROW-6839: [Java] Add APIs to read and write
"custom_metadata" field of IPC file footer (#7231)
b196d1a is described below
commit b196d1a86312660d0900d5edbe8757e8d23c7e73
Author: Ji Liu <[email protected]>
AuthorDate: Tue Jun 16 14:07:47 2020 +0800
ARROW-6839: [Java] Add APIs to read and write "custom_metadata" field of
IPC file footer (#7231)
* ARROW-6839: [Java] Add APIs to read and write "custom_metadata" field of
IPC file footer
* remove dictionary in test
* extract kv write methods and simplify tests
Co-authored-by: tianchen <[email protected]>
---
.../apache/arrow/vector/ipc/ArrowFileReader.java | 12 ++++++
.../apache/arrow/vector/ipc/ArrowFileWriter.java | 11 +++++-
.../arrow/vector/ipc/message/ArrowFooter.java | 45 +++++++++++++++++++++-
.../arrow/vector/ipc/message/FBSerializables.java | 22 +++++++++++
.../org/apache/arrow/vector/types/pojo/Schema.java | 16 +-------
.../arrow/vector/ipc/TestArrowReaderWriter.java | 32 +++++++++++++++
6 files changed, 121 insertions(+), 17 deletions(-)
diff --git
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileReader.java
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileReader.java
index 4a9726a..82ddbbf 100644
--- a/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileReader.java
+++ b/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileReader.java
@@ -21,7 +21,9 @@ import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.SeekableByteChannel;
import java.util.Arrays;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import org.apache.arrow.flatbuf.Footer;
import org.apache.arrow.memory.BufferAllocator;
@@ -113,6 +115,16 @@ public class ArrowFileReader extends ArrowReader {
}
/**
+ * Get custom metadata.
+ */
+ public Map<String, String> getMetaData() {
+ if (footer != null) {
+ return footer.getMetaData();
+ }
+ return new HashMap<>();
+ }
+
+ /**
* Read a dictionary batch from the source, will be invoked after the schema
has been read and
* called N times, where N is the number of dictionaries indicated by the
schema Fields.
*
diff --git
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileWriter.java
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileWriter.java
index 673cc6c..fb1ca00 100644
--- a/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileWriter.java
+++ b/java/vector/src/main/java/org/apache/arrow/vector/ipc/ArrowFileWriter.java
@@ -21,6 +21,7 @@ import java.io.IOException;
import java.nio.channels.WritableByteChannel;
import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
import org.apache.arrow.util.VisibleForTesting;
import org.apache.arrow.vector.VectorSchemaRoot;
@@ -45,11 +46,19 @@ public class ArrowFileWriter extends ArrowWriter {
private final List<ArrowBlock> dictionaryBlocks = new ArrayList<>();
private final List<ArrowBlock> recordBlocks = new ArrayList<>();
+ private Map<String, String> metaData;
+
public ArrowFileWriter(VectorSchemaRoot root, DictionaryProvider provider,
WritableByteChannel out) {
super(root, provider, out);
}
public ArrowFileWriter(VectorSchemaRoot root, DictionaryProvider provider,
WritableByteChannel out,
+ Map<String, String> metaData) {
+ super(root, provider, out);
+ this.metaData = metaData;
+ }
+
+ public ArrowFileWriter(VectorSchemaRoot root, DictionaryProvider provider,
WritableByteChannel out,
IpcOption option) {
super(root, provider, out, option);
}
@@ -81,7 +90,7 @@ public class ArrowFileWriter extends ArrowWriter {
out.writeIntLittleEndian(0);
long footerStart = out.getCurrentPosition();
- out.write(new ArrowFooter(schema, dictionaryBlocks, recordBlocks), false);
+ out.write(new ArrowFooter(schema, dictionaryBlocks, recordBlocks,
metaData), false);
int footerLength = (int) (out.getCurrentPosition() - footerStart);
if (footerLength <= 0) {
throw new InvalidArrowFileException("invalid footer");
diff --git
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/ArrowFooter.java
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/ArrowFooter.java
index f704780..77d3b1e 100644
---
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/ArrowFooter.java
+++
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/ArrowFooter.java
@@ -18,12 +18,16 @@
package org.apache.arrow.vector.ipc.message;
import static
org.apache.arrow.vector.ipc.message.FBSerializables.writeAllStructsToVector;
+import static
org.apache.arrow.vector.ipc.message.FBSerializables.writeKeyValues;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import org.apache.arrow.flatbuf.Block;
import org.apache.arrow.flatbuf.Footer;
+import org.apache.arrow.flatbuf.KeyValue;
import org.apache.arrow.vector.types.pojo.Schema;
import com.google.flatbuffers.FlatBufferBuilder;
@@ -37,17 +41,30 @@ public class ArrowFooter implements FBSerializable {
private final List<ArrowBlock> recordBatches;
+ private final Map<String, String> metaData;
+
+ public ArrowFooter(Schema schema, List<ArrowBlock> dictionaries,
List<ArrowBlock> recordBatches) {
+ this(schema, dictionaries, recordBatches, null);
+ }
+
/**
* Constructs a new instance.
*
* @param schema The schema for record batches in the file.
* @param dictionaries The dictionaries relevant to the file.
* @param recordBatches The recordBatches written to the file.
+ * @param metaData user-defined k-v meta data.
*/
- public ArrowFooter(Schema schema, List<ArrowBlock> dictionaries,
List<ArrowBlock> recordBatches) {
+ public ArrowFooter(
+ Schema schema,
+ List<ArrowBlock> dictionaries,
+ List<ArrowBlock> recordBatches,
+ Map<String, String> metaData) {
+
this.schema = schema;
this.dictionaries = dictionaries;
this.recordBatches = recordBatches;
+ this.metaData = metaData;
}
/**
@@ -57,7 +74,8 @@ public class ArrowFooter implements FBSerializable {
this(
Schema.convertSchema(footer.schema()),
dictionaries(footer),
- recordBatches(footer)
+ recordBatches(footer),
+ metaData(footer)
);
}
@@ -84,6 +102,18 @@ public class ArrowFooter implements FBSerializable {
return dictionaries;
}
+ private static Map<String, String> metaData(Footer footer) {
+ Map<String, String> metaData = new HashMap<>();
+
+ int metaDataLength = footer.customMetadataLength();
+ for (int i = 0; i < metaDataLength; i++) {
+ KeyValue kv = footer.customMetadata(i);
+ metaData.put(kv.key(), kv.value());
+ }
+
+ return metaData;
+ }
+
public Schema getSchema() {
return schema;
}
@@ -96,6 +126,10 @@ public class ArrowFooter implements FBSerializable {
return recordBatches;
}
+ public Map<String, String> getMetaData() {
+ return metaData;
+ }
+
@Override
public int writeTo(FlatBufferBuilder builder) {
int schemaIndex = schema.getSchema(builder);
@@ -103,10 +137,17 @@ public class ArrowFooter implements FBSerializable {
int dicsOffset = writeAllStructsToVector(builder, dictionaries);
Footer.startRecordBatchesVector(builder, recordBatches.size());
int rbsOffset = writeAllStructsToVector(builder, recordBatches);
+
+ int metaDataOffset = 0;
+ if (metaData != null) {
+ metaDataOffset = writeKeyValues(builder, metaData);
+ }
+
Footer.startFooter(builder);
Footer.addSchema(builder, schemaIndex);
Footer.addDictionaries(builder, dicsOffset);
Footer.addRecordBatches(builder, rbsOffset);
+ Footer.addCustomMetadata(builder, metaDataOffset);
return Footer.endFooter(builder);
}
diff --git
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/FBSerializables.java
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/FBSerializables.java
index f139d62..26736ed 100644
---
a/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/FBSerializables.java
+++
b/java/vector/src/main/java/org/apache/arrow/vector/ipc/message/FBSerializables.java
@@ -19,7 +19,11 @@ package org.apache.arrow.vector.ipc.message;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.Iterator;
import java.util.List;
+import java.util.Map;
+
+import org.apache.arrow.flatbuf.KeyValue;
import com.google.flatbuffers.FlatBufferBuilder;
@@ -42,4 +46,22 @@ public class FBSerializables {
}
return builder.endVector();
}
+
+ /**
+ * Writes map data with string type.
+ */
+ public static int writeKeyValues(FlatBufferBuilder builder, Map<String,
String> metaData) {
+ int[] metadataOffsets = new int[metaData.size()];
+ Iterator<Map.Entry<String, String>> metadataIterator =
metaData.entrySet().iterator();
+ for (int i = 0; i < metadataOffsets.length; i++) {
+ Map.Entry<String, String> kv = metadataIterator.next();
+ int keyOffset = builder.createString(kv.getKey());
+ int valueOffset = builder.createString(kv.getValue());
+ KeyValue.startKeyValue(builder);
+ KeyValue.addKey(builder, keyOffset);
+ KeyValue.addValue(builder, valueOffset);
+ metadataOffsets[i] = KeyValue.endKeyValue(builder);
+ }
+ return org.apache.arrow.flatbuf.Field.createCustomMetadataVector(builder,
metadataOffsets);
+ }
}
diff --git
a/java/vector/src/main/java/org/apache/arrow/vector/types/pojo/Schema.java
b/java/vector/src/main/java/org/apache/arrow/vector/types/pojo/Schema.java
index 487f8b4..7ada43e 100644
--- a/java/vector/src/main/java/org/apache/arrow/vector/types/pojo/Schema.java
+++ b/java/vector/src/main/java/org/apache/arrow/vector/types/pojo/Schema.java
@@ -26,16 +26,15 @@ import java.util.AbstractMap;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
-import java.util.Iterator;
import java.util.List;
import java.util.Map;
-import java.util.Map.Entry;
import java.util.Objects;
import java.util.stream.Collectors;
import org.apache.arrow.flatbuf.KeyValue;
import org.apache.arrow.util.Collections2;
import org.apache.arrow.util.Preconditions;
+import org.apache.arrow.vector.ipc.message.FBSerializables;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonInclude;
@@ -179,18 +178,7 @@ public class Schema {
fieldOffsets[i] = fields.get(i).getField(builder);
}
int fieldsOffset =
org.apache.arrow.flatbuf.Schema.createFieldsVector(builder, fieldOffsets);
- int[] metadataOffsets = new int[metadata.size()];
- Iterator<Entry<String, String>> metadataIterator =
metadata.entrySet().iterator();
- for (int i = 0; i < metadataOffsets.length; i++) {
- Entry<String, String> kv = metadataIterator.next();
- int keyOffset = builder.createString(kv.getKey());
- int valueOffset = builder.createString(kv.getValue());
- KeyValue.startKeyValue(builder);
- KeyValue.addKey(builder, keyOffset);
- KeyValue.addValue(builder, valueOffset);
- metadataOffsets[i] = KeyValue.endKeyValue(builder);
- }
- int metadataOffset =
org.apache.arrow.flatbuf.Field.createCustomMetadataVector(builder,
metadataOffsets);
+ int metadataOffset = FBSerializables.writeKeyValues(builder, metadata);
org.apache.arrow.flatbuf.Schema.startSchema(builder);
org.apache.arrow.flatbuf.Schema.addFields(builder, fieldsOffset);
org.apache.arrow.flatbuf.Schema.addCustomMetadata(builder, metadataOffset);
diff --git
a/java/vector/src/test/java/org/apache/arrow/vector/ipc/TestArrowReaderWriter.java
b/java/vector/src/test/java/org/apache/arrow/vector/ipc/TestArrowReaderWriter.java
index 0804856..15a19ed 100644
---
a/java/vector/src/test/java/org/apache/arrow/vector/ipc/TestArrowReaderWriter.java
+++
b/java/vector/src/test/java/org/apache/arrow/vector/ipc/TestArrowReaderWriter.java
@@ -37,8 +37,10 @@ import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
+import java.util.Map;
import java.util.stream.Collectors;
import org.apache.arrow.flatbuf.FieldNode;
@@ -751,4 +753,34 @@ public class TestArrowReaderWriter {
assertEquals(10, arrBuf.getInt(0));
}
}
+
+ @Test
+ public void testCustomMetaData() throws IOException {
+
+ VarCharVector vector = newVarCharVector("varchar1", allocator);
+
+ List<Field> fields = Arrays.asList(vector.getField());
+ List<FieldVector> vectors = Collections2.asImmutableList(vector);
+ Map<String, String> metadata = new HashMap<>();
+ metadata.put("key1", "value1");
+ metadata.put("key2", "value2");
+ try (VectorSchemaRoot root = new VectorSchemaRoot(fields, vectors,
vector.getValueCount());
+ ByteArrayOutputStream out = new ByteArrayOutputStream();
+ ArrowFileWriter writer = new ArrowFileWriter(root, null,
newChannel(out), metadata);) {
+
+ writer.start();
+ writer.end();
+
+ try (SeekableReadChannel channel = new SeekableReadChannel(
+ new ByteArrayReadableSeekableByteChannel(out.toByteArray()));
+ ArrowFileReader reader = new ArrowFileReader(channel, allocator)) {
+ reader.getVectorSchemaRoot();
+
+ Map<String, String> readMeta = reader.getMetaData();
+ assertEquals(2, readMeta.size());
+ assertEquals("value1", readMeta.get("key1"));
+ assertEquals("value2", readMeta.get("key2"));
+ }
+ }
+ }
}