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"));
+      }
+    }
+  }
 }

Reply via email to