This is an automated email from the ASF dual-hosted git repository.
mattyb149 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 87fd9f9 NIFI-7958 ConvertAvroToORC - decimal support
87fd9f9 is described below
commit 87fd9f968fceaccac8c115cc4cd83c608d4e04b7
Author: Arek Burdach <[email protected]>
AuthorDate: Tue Oct 27 15:51:12 2020 +0100
NIFI-7958 ConvertAvroToORC - decimal support
Signed-off-by: Matthew Burgess <[email protected]>
This closes #4628
---
.../apache/hadoop/hive/ql/io/orc/NiFiOrcUtils.java | 36 +++++++++++++++++++---
.../nifi/processors/hive/TestConvertAvroToORC.java | 17 +++++++---
.../org/apache/nifi/util/orc/TestNiFiOrcUtils.java | 19 ++++++++----
3 files changed, 58 insertions(+), 14 deletions(-)
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 f75a9e9..70e320a 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
@@ -30,6 +30,7 @@ import
org.apache.hadoop.hive.serde2.objectinspector.ObjectInspector;
import org.apache.hadoop.hive.serde2.objectinspector.ObjectInspectorFactory;
import
org.apache.hadoop.hive.serde2.objectinspector.SettableStructObjectInspector;
import org.apache.hadoop.hive.serde2.objectinspector.StructField;
+import org.apache.hadoop.hive.serde2.typeinfo.DecimalTypeInfo;
import org.apache.hadoop.hive.serde2.typeinfo.ListTypeInfo;
import org.apache.hadoop.hive.serde2.typeinfo.MapTypeInfo;
import org.apache.hadoop.hive.serde2.typeinfo.TypeInfo;
@@ -120,6 +121,10 @@ public class NiFiOrcUtils {
if (o instanceof String || o instanceof Utf8 || o instanceof
GenericData.EnumSymbol) {
return new Text(o.toString());
}
+ if (o instanceof ByteBuffer && typeInfo instanceof
DecimalTypeInfo) {
+ ByteBuffer buffer = (ByteBuffer) o;
+ return new HiveDecimalWritable(buffer.array(),
((DecimalTypeInfo) typeInfo).scale());
+ }
if (o instanceof ByteBuffer) {
return new BytesWritable(((ByteBuffer) o).array());
}
@@ -253,13 +258,17 @@ public class NiFiOrcUtils {
case INT:
case LONG:
case BOOLEAN:
- case BYTES:
case DOUBLE:
case FLOAT:
case STRING:
case NULL:
return getPrimitiveOrcTypeFromPrimitiveAvroType(fieldType);
-
+ case BYTES:
+ if (isLogicalType(fieldSchema)){
+ return getLogicalTypeInfo(fieldSchema);
+ } else {
+ return getPrimitiveOrcTypeFromPrimitiveAvroType(fieldType);
+ }
case UNION:
List<Schema> unionFieldSchemas = fieldSchema.getTypes();
@@ -311,6 +320,21 @@ public class NiFiOrcUtils {
}
+ private static boolean isLogicalType(Schema schema){
+ return schema.getProp("logicalType") != null;
+ }
+
+ private static TypeInfo getLogicalTypeInfo(Schema schema){
+ String type = schema.getProp("logicalType");
+ switch (type){
+ case "decimal":
+ int precision = schema.getJsonProp("precision").asInt(10);
+ int scale = schema.getJsonProp("scale").asInt(2);
+ return new DecimalTypeInfo(precision, scale);
+ }
+ throw new IllegalArgumentException("Logical type " + type + " is not
supported!");
+ }
+
public static Schema.Type getAvroSchemaTypeOfObject(Object o) {
if (o == null) {
return Schema.Type.NULL;
@@ -380,7 +404,11 @@ public class NiFiOrcUtils {
case NULL: // Hive has no null type, we picked boolean as the ORC
type so use it for Hive DDL too. All values are necessarily null.
return "BOOLEAN";
case BYTES:
- return "BINARY";
+ if (isLogicalType(avroSchema)){
+ return
getLogicalTypeInfo(avroSchema).toString().toUpperCase();
+ } else {
+ return "BINARY";
+ }
case DOUBLE:
return "DOUBLE";
case FLOAT:
@@ -577,4 +605,4 @@ public class NiFiOrcUtils {
}
return memoryManager;
}
-}
\ No newline at end of file
+}
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 f34a647..9248b27 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
@@ -33,6 +33,7 @@ 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;
import org.apache.hadoop.hive.ql.io.orc.RecordReader;
+import org.apache.hadoop.hive.serde2.io.HiveDecimalWritable;
import org.apache.hadoop.hive.serde2.objectinspector.StructObjectInspector;
import org.apache.hadoop.hive.serde2.typeinfo.TypeInfo;
import org.apache.hadoop.io.DoubleWritable;
@@ -51,6 +52,8 @@ import java.io.ByteArrayOutputStream;
import java.io.InputStream;
import java.io.File;
import java.io.FileOutputStream;
+import java.math.BigDecimal;
+import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
@@ -240,7 +243,9 @@ public class TestConvertAvroToORC {
put("key2", 2.0);
}};
- GenericData.Record record =
TestNiFiOrcUtils.buildComplexAvroRecord(10, mapData1, "DEF", 3.0f,
Arrays.asList(10, 20));
+ BigDecimal sampleBigDecimal = new BigDecimal("12.34");
+ ByteBuffer bigDecimalAsBytes =
ByteBuffer.wrap(sampleBigDecimal.unscaledValue().toByteArray());
+ GenericData.Record record =
TestNiFiOrcUtils.buildComplexAvroRecord(10, mapData1, "DEF", 3.0f,
Arrays.asList(10, 20), bigDecimalAsBytes);
DatumWriter<GenericData.Record> writer = new
GenericDatumWriter<>(record.getSchema());
DataFileWriter<GenericData.Record> fileWriter = new
DataFileWriter<>(writer);
@@ -254,7 +259,7 @@ public class TestConvertAvroToORC {
put("key2", 4.0);
}};
- record = TestNiFiOrcUtils.buildComplexAvroRecord(null, mapData2,
"XYZ", 4L, Arrays.asList(100, 200));
+ record = TestNiFiOrcUtils.buildComplexAvroRecord(null, mapData2,
"XYZ", 4L, Arrays.asList(100, 200), bigDecimalAsBytes);
fileWriter.append(record);
fileWriter.flush();
@@ -272,7 +277,7 @@ public class TestConvertAvroToORC {
// 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 complex_record " +
- "(myInt INT, myMap MAP<STRING, DOUBLE>, myEnum STRING,
myLongOrFloat UNIONTYPE<BIGINT, FLOAT>, myIntList ARRAY<INT>)"
+ "(myInt INT, myMap MAP<STRING, DOUBLE>, myEnum STRING,
myLongOrFloat UNIONTYPE<BIGINT, FLOAT>, myIntList ARRAY<INT>, myDecimal
DECIMAL(10,2))"
+ " STORED AS ORC",
resultFlowFile.getAttribute(ConvertAvroToORC.HIVE_DDL_ATTRIBUTE));
assertEquals("2",
resultFlowFile.getAttribute(ConvertAvroToORC.RECORD_COUNT_ATTRIBUTE));
assertEquals("test.orc",
resultFlowFile.getAttribute(CoreAttributes.FILENAME.key()));
@@ -309,6 +314,10 @@ public class TestConvertAvroToORC {
assertNotNull(mapValue);
assertTrue(mapValue instanceof DoubleWritable);
assertEquals(2.0, ((DoubleWritable) mapValue).get(), Double.MIN_VALUE);
+
+ Object decimalFieldObject = inspector.getStructFieldData(o,
inspector.getStructFieldRef("myDecimal"));
+ assertTrue(decimalFieldObject instanceof HiveDecimalWritable);
+ assertEquals(sampleBigDecimal, ((HiveDecimalWritable)
decimalFieldObject).getHiveDecimal().bigDecimalValue());
}
@Test
@@ -498,4 +507,4 @@ public class TestConvertAvroToORC {
assertTrue(el0 instanceof Map);
assertEquals(new Text("v1"), ((Map) el0).get(new Text("key1")));
}
-}
\ No newline at end of file
+}
diff --git
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/util/orc/TestNiFiOrcUtils.java
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/util/orc/TestNiFiOrcUtils.java
index 74cdd2a..c743fb0 100644
---
a/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/util/orc/TestNiFiOrcUtils.java
+++
b/nifi-nar-bundles/nifi-hive-bundle/nifi-hive-processors/src/test/java/org/apache/nifi/util/orc/TestNiFiOrcUtils.java
@@ -244,7 +244,8 @@ public class TestNiFiOrcUtils {
"MAP<STRING, DOUBLE>",
"STRING",
"UNIONTYPE<BIGINT, FLOAT>",
- "ARRAY<INT>"
+ "ARRAY<INT>",
+ "DECIMAL(10,2)"
};
Schema testSchema = buildComplexAvroSchema();
@@ -253,7 +254,7 @@ public class TestNiFiOrcUtils {
assertEquals(expectedTypes[i],
NiFiOrcUtils.getHiveTypeFromAvroType(fields.get(i).schema()));
}
- assertEquals("STRUCT<myInt:INT, myMap:MAP<STRING, DOUBLE>,
myEnum:STRING, myLongOrFloat:UNIONTYPE<BIGINT, FLOAT>, myIntList:ARRAY<INT>>",
+ assertEquals("STRUCT<myInt:INT, myMap:MAP<STRING, DOUBLE>,
myEnum:STRING, myLongOrFloat:UNIONTYPE<BIGINT, FLOAT>, myIntList:ARRAY<INT>,
myDecimal:DECIMAL(10,2)>",
NiFiOrcUtils.getHiveTypeFromAvroType(testSchema));
}
@@ -270,7 +271,7 @@ public class TestNiFiOrcUtils {
Schema avroSchema = buildComplexAvroSchema();
String ddl = NiFiOrcUtils.generateHiveDDL(avroSchema, "myHiveTable");
assertEquals("CREATE EXTERNAL TABLE IF NOT EXISTS myHiveTable "
- + "(myInt INT, myMap MAP<STRING, DOUBLE>, myEnum STRING,
myLongOrFloat UNIONTYPE<BIGINT, FLOAT>, myIntList ARRAY<INT>)"
+ + "(myInt INT, myMap MAP<STRING, DOUBLE>, myEnum STRING,
myLongOrFloat UNIONTYPE<BIGINT, FLOAT>, myIntList ARRAY<INT>, myDecimal
DECIMAL(10,2))"
+ " STORED AS ORC", ddl);
}
@@ -365,10 +366,15 @@ public class TestNiFiOrcUtils {
builder.name("myEnum").type().enumeration("myEnum").symbols("ABC",
"DEF", "XYZ").enumDefault("ABC");
builder.name("myLongOrFloat").type().unionOf().longType().and().floatType().endUnion().noDefault();
builder.name("myIntList").type().array().items().intType().noDefault();
+ builder.name("myDecimal").type().bytesBuilder()
+ .prop("logicalType", "decimal")
+ .prop("precision", "10")
+ .prop("scale", "2")
+ .endBytes().noDefault();
return builder.endRecord();
}
- public static GenericData.Record buildComplexAvroRecord(Integer i,
Map<String, Double> m, String e, Object unionVal, List<Integer> intArray) {
+ public static GenericData.Record buildComplexAvroRecord(Integer i,
Map<String, Double> m, String e, Object unionVal, List<Integer> intArray,
ByteBuffer decimal) {
Schema schema = buildComplexAvroSchema();
GenericData.Record row = new GenericData.Record(schema);
row.put("myInt", i);
@@ -376,6 +382,7 @@ public class TestNiFiOrcUtils {
row.put("myEnum", e);
row.put("myLongOrFloat", unionVal);
row.put("myIntList", intArray);
+ row.put("myDecimal", decimal);
return row;
}
@@ -403,7 +410,7 @@ public class TestNiFiOrcUtils {
}
public static TypeInfo buildComplexOrcSchema() {
- return
TypeInfoUtils.getTypeInfoFromTypeString("struct<myInt:int,myMap:map<string,double>,myEnum:string,myLongOrFloat:uniontype<int>,myIntList:array<int>>");
+ return
TypeInfoUtils.getTypeInfoFromTypeString("struct<myInt:int,myMap:map<string,double>,myEnum:string,myLongOrFloat:uniontype<int>,myIntList:array<int>,myDecimal:decimal(10,2)>");
}
public static Schema buildNestedComplexAvroSchema() {
@@ -455,4 +462,4 @@ public class TestNiFiOrcUtils {
return TypeInfoFactory.getPrimitiveTypeInfo("string");
}
}
-}
\ No newline at end of file
+}