Repository: nifi
Updated Branches:
  refs/heads/master 4847ed28a -> 6874a5d82


NIFI-2593: This closes #891. Fixed handling of nested records/structs in 
ConvertAvroToORC


Project: http://git-wip-us.apache.org/repos/asf/nifi/repo
Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/6874a5d8
Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/6874a5d8
Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/6874a5d8

Branch: refs/heads/master
Commit: 6874a5d82da27c1dfa41dfdd3fb8db25584c3bc0
Parents: 4847ed2
Author: Matt Burgess <[email protected]>
Authored: Wed Aug 17 20:38:15 2016 -0400
Committer: joewitt <[email protected]>
Committed: Thu Aug 18 15:15:51 2016 -0400

----------------------------------------------------------------------
 .../hadoop/hive/ql/io/orc/NiFiOrcUtils.java     | 20 +++-
 .../processors/hive/TestConvertAvroToORC.java   | 97 ++++++++++++++++++++
 .../src/test/resources/array_of_records.avsc    | 38 ++++++++
 3 files changed, 153 insertions(+), 2 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/nifi/blob/6874a5d8/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/main/java/org/apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java
----------------------------------------------------------------------
diff --git 
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/main/java/org/apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java
 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/main/java/org/apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java
index 78430bf..b8c9dab 100644
--- 
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/main/java/org/apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java
+++ 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/main/java/org/apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java
@@ -165,9 +165,25 @@ public class NiFiOrcUtils {
                 });
                 return mapWritable;
             }
-
+            if (o instanceof GenericData.Record) {
+                GenericData.Record record = (GenericData.Record) o;
+                TypeInfo recordSchema = 
NiFiOrcUtils.getOrcField(record.getSchema());
+                List<Schema.Field> recordFields = 
record.getSchema().getFields();
+                if (recordFields != null) {
+                    Object[] fieldObjects = new Object[recordFields.size()];
+                    for (int i = 0; i < recordFields.size(); i++) {
+                        Schema.Field field = recordFields.get(i);
+                        Schema fieldSchema = field.schema();
+                        Object fieldObject = record.get(field.name());
+                        fieldObjects[i] = 
NiFiOrcUtils.convertToORCObject(NiFiOrcUtils.getOrcField(fieldSchema), 
fieldObject);
+                    }
+                    return NiFiOrcUtils.createOrcStruct(recordSchema, 
fieldObjects);
+                }
+            }
+            throw new IllegalArgumentException("Error converting object of 
type " + o.getClass().getName() + " to ORC type " + typeInfo.getTypeName());
+        } else {
+            return null;
         }
-        return null;
     }
 
 

http://git-wip-us.apache.org/repos/asf/nifi/blob/6874a5d8/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/processors/hive/TestConvertAvroToORC.java
----------------------------------------------------------------------
diff --git 
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/processors/hive/TestConvertAvroToORC.java
 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/processors/hive/TestConvertAvroToORC.java
index ad40198..1c269fa 100644
--- 
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/processors/hive/TestConvertAvroToORC.java
+++ 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/processors/hive/TestConvertAvroToORC.java
@@ -16,13 +16,16 @@
  */
 package org.apache.nifi.processors.hive;
 
+import org.apache.avro.Schema;
 import org.apache.avro.file.DataFileWriter;
 import org.apache.avro.generic.GenericData;
 import org.apache.avro.generic.GenericDatumWriter;
+import org.apache.avro.generic.GenericRecord;
 import org.apache.avro.io.DatumWriter;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.fs.FileSystem;
 import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hive.ql.io.orc.NiFiOrcUtils;
 import org.apache.hadoop.hive.ql.io.orc.OrcFile;
 import org.apache.hadoop.hive.ql.io.orc.OrcStruct;
 import org.apache.hadoop.hive.ql.io.orc.Reader;
@@ -41,10 +44,14 @@ import org.junit.Before;
 import org.junit.Test;
 
 import java.io.ByteArrayOutputStream;
+import java.io.File;
 import java.io.FileOutputStream;
 import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.HashMap;
+import java.util.LinkedList;
+import java.util.List;
 import java.util.Map;
 import java.util.TreeMap;
 
@@ -209,6 +216,96 @@ public class TestConvertAvroToORC {
         assertNotNull(mapValue);
         assertTrue(mapValue instanceof DoubleWritable);
         assertEquals(2.0, ((DoubleWritable) mapValue).get(), Double.MIN_VALUE);
+    }
+
+    @Test
+    public void test_onTrigger_array_of_records() throws Exception {
+        final Schema schema = new Schema.Parser().parse(new 
File("src/test/resources/array_of_records.avsc"));
+        List<GenericRecord> innerRecords = new LinkedList<>();
+
+        final GenericRecord outerRecord = new GenericData.Record(schema);
+
+        Schema arraySchema = schema.getField("records").schema();
+        Schema innerRecordSchema = arraySchema.getElementType();
+        final GenericRecord innerRecord1 = new 
GenericData.Record(innerRecordSchema);
+        innerRecord1.put("name", "Joe");
+        innerRecord1.put("age", 42);
+
+        innerRecords.add(innerRecord1);
+
+        final GenericRecord innerRecord2 = new 
GenericData.Record(innerRecordSchema);
+        innerRecord2.put("name", "Mary");
+        innerRecord2.put("age", 28);
+
+        innerRecords.add(innerRecord2);
+
+        GenericData.Array<GenericRecord> array = new 
GenericData.Array<>(arraySchema, innerRecords);
+        outerRecord.put("records", array);
+
+        final DatumWriter<GenericRecord> datumWriter = new 
GenericDatumWriter<>(schema);
+        ByteArrayOutputStream out = new ByteArrayOutputStream();
+        try (DataFileWriter<GenericRecord> dataFileWriter = new 
DataFileWriter<>(datumWriter)) {
+            dataFileWriter.create(schema, out);
+            dataFileWriter.append(outerRecord);
+        }
+        out.close();
+
+        // Build a flow file from the Avro record
+        Map<String, String> attributes = new HashMap<String, String>() {{
+            put(CoreAttributes.FILENAME.key(), "test");
+        }};
+        runner.enqueue(out.toByteArray(), attributes);
+        runner.run();
+
+        runner.assertAllFlowFilesTransferred(ConvertAvroToORC.REL_SUCCESS, 1);
+
+        // Write the flow file out to disk, since the ORC Reader needs a path
+        MockFlowFile resultFlowFile = 
runner.getFlowFilesForRelationship(ConvertAvroToORC.REL_SUCCESS).get(0);
+        assertEquals("CREATE EXTERNAL TABLE IF NOT EXISTS 
org_apache_nifi_outer_record " +
+                "(records ARRAY<STRUCT<name:STRING, age:INT>>)"
+                + " STORED AS ORC", 
resultFlowFile.getAttribute(ConvertAvroToORC.HIVE_DDL_ATTRIBUTE));
+        assertEquals("1", 
resultFlowFile.getAttribute(ConvertAvroToORC.RECORD_COUNT_ATTRIBUTE));
+        assertEquals("test.orc", 
resultFlowFile.getAttribute(CoreAttributes.FILENAME.key()));
+        byte[] resultContents = runner.getContentAsByteArray(resultFlowFile);
+        FileOutputStream fos = new FileOutputStream("target/test1.orc");
+        fos.write(resultContents);
+        fos.flush();
+        fos.close();
+
+        Configuration conf = new Configuration();
+        FileSystem fs = FileSystem.getLocal(conf);
+        Reader reader = OrcFile.createReader(new Path("target/test1.orc"), 
OrcFile.readerOptions(conf).filesystem(fs));
+        RecordReader rows = reader.rows();
+        Object o = rows.next(null);
+        assertNotNull(o);
+        assertTrue(o instanceof OrcStruct);
+        StructObjectInspector inspector = (StructObjectInspector) 
OrcStruct.createObjectInspector(NiFiOrcUtils.getOrcField(schema));
+
+        // Verify the record contains an array
+        Object arrayFieldObject = inspector.getStructFieldData(o, 
inspector.getStructFieldRef("records"));
+        assertTrue(arrayFieldObject instanceof ArrayList);
+        ArrayList<?> arrayField = (ArrayList<?>) arrayFieldObject;
+        assertEquals(2, arrayField.size());
+
+        // Verify the first element. Should be a record with two fields "name" 
and "age"
+        Object element = arrayField.get(0);
+        assertTrue(element instanceof OrcStruct);
+        StructObjectInspector elementInspector = (StructObjectInspector) 
OrcStruct.createObjectInspector(NiFiOrcUtils.getOrcField(innerRecordSchema));
+        Object nameObject = elementInspector.getStructFieldData(element, 
elementInspector.getStructFieldRef("name"));
+        assertTrue(nameObject instanceof Text);
+        assertEquals("Joe", nameObject.toString());
+        Object ageObject = elementInspector.getStructFieldData(element, 
elementInspector.getStructFieldRef("age"));
+        assertTrue(ageObject instanceof IntWritable);
+        assertEquals(42, ((IntWritable) ageObject).get());
 
+        // Verify the first element. Should be a record with two fields "name" 
and "age"
+        element = arrayField.get(1);
+        assertTrue(element instanceof OrcStruct);
+        nameObject = elementInspector.getStructFieldData(element, 
elementInspector.getStructFieldRef("name"));
+        assertTrue(nameObject instanceof Text);
+        assertEquals("Mary", nameObject.toString());
+        ageObject = elementInspector.getStructFieldData(element, 
elementInspector.getStructFieldRef("age"));
+        assertTrue(ageObject instanceof IntWritable);
+        assertEquals(28, ((IntWritable) ageObject).get());
     }
 }
\ No newline at end of file

http://git-wip-us.apache.org/repos/asf/nifi/blob/6874a5d8/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/resources/array_of_records.avsc
----------------------------------------------------------------------
diff --git 
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/resources/array_of_records.avsc
 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/resources/array_of_records.avsc
new file mode 100644
index 0000000..19cac6e
--- /dev/null
+++ 
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/resources/array_of_records.avsc
@@ -0,0 +1,38 @@
+/*
+ * 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.
+ */
+ {
+  "namespace" : "org.apache.nifi",
+  "name" : "outer_record",
+  "type" : "record",
+  "fields" : [ {
+    "name" : "records",
+    "type" : {
+      "type" : "array",
+      "items" : {
+        "type" : "record",
+        "name" : "inner_record",
+        "fields" : [ {
+          "name" : "name",
+          "type" : "string"
+        }, {
+          "name" : "age",
+          "type" : "int"
+        } ]
+      }
+    }
+  } ]
+}
\ No newline at end of file

Reply via email to