diff --git a/iceberg/iceberg-handler/src/test/queries/positive/variant_type_projection.q b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_projection.q new file mode 100644 index 000000000000..2076693c5fa5 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/queries/positive/variant_type_projection.q @@ -0,0 +1,66 @@ +-- SORT_QUERY_RESULTS +set hive.explain.user=false; + +drop table if exists variant_proj; + +CREATE EXTERNAL TABLE variant_proj ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +); + +INSERT INTO variant_proj VALUES +(1, parse_json('{"name": "John", "age": 30}')), +(2, parse_json('[1, 2, 3]')), +(3, parse_json('"plain string"')), +(4, parse_json('42')); + +INSERT INTO variant_proj (id) SELECT 5; + +-- raw projection: variants are returned as their logical JSON +EXPLAIN SELECT id, data FROM variant_proj; + +SELECT id, data FROM variant_proj; + +SELECT * FROM variant_proj; + +-- raw projection across a shuffle +SELECT id, data FROM variant_proj ORDER BY id DESC; + +-- variant nested inside a struct is decoded within the folded output +-- (array/map<..,variant> reads fail in Iceberg column projection; tracked separately) +drop table if exists variant_nested; + +CREATE EXTERNAL TABLE variant_nested ( + id INT, + s STRUCT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +); + +INSERT INTO variant_nested SELECT + 1, + named_struct('label', 'x', 'v', parse_json('{"k": 1}')); + +SELECT * FROM variant_nested; + +SELECT id, to_json(s.v) FROM variant_nested; + +-- a plain struct with the variant shape keeps the generic struct rendering +drop table if exists bin_struct; + +CREATE EXTERNAL TABLE bin_struct ( + id INT, + s STRUCT +) STORED BY ICEBERG; + +INSERT INTO bin_struct SELECT 1, named_struct('metadata', cast('x' as binary), 'value', cast('y' as binary)); + +SELECT * FROM bin_struct; + +drop table variant_proj; +drop table variant_nested; +drop table bin_struct; diff --git a/iceberg/iceberg-handler/src/test/results/positive/variant_type_projection.q.out b/iceberg/iceberg-handler/src/test/results/positive/variant_type_projection.q.out new file mode 100644 index 000000000000..e4165c150de3 --- /dev/null +++ b/iceberg/iceberg-handler/src/test/results/positive/variant_type_projection.q.out @@ -0,0 +1,235 @@ +PREHOOK: query: drop table if exists variant_proj +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_proj +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_proj ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_proj +POSTHOOK: query: CREATE EXTERNAL TABLE variant_proj ( + id INT, + data VARIANT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_proj +PREHOOK: query: INSERT INTO variant_proj VALUES +(1, parse_json('{"name": "John", "age": 30}')), +(2, parse_json('[1, 2, 3]')), +(3, parse_json('"plain string"')), +(4, parse_json('42')) +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_proj +POSTHOOK: query: INSERT INTO variant_proj VALUES +(1, parse_json('{"name": "John", "age": 30}')), +(2, parse_json('[1, 2, 3]')), +(3, parse_json('"plain string"')), +(4, parse_json('42')) +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_proj +PREHOOK: query: INSERT INTO variant_proj (id) SELECT 5 +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_proj +POSTHOOK: query: INSERT INTO variant_proj (id) SELECT 5 +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_proj +PREHOOK: query: EXPLAIN SELECT id, data FROM variant_proj +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_proj +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: EXPLAIN SELECT id, data FROM variant_proj +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_proj +POSTHOOK: Output: hdfs://### HDFS PATH ### +STAGE DEPENDENCIES: + Stage-0 is a root stage + +STAGE PLANS: + Stage: Stage-0 + Fetch Operator + limit: -1 + Processor Tree: + TableScan + alias: variant_proj + Select Operator + expressions: id (type: int), data (type: struct) + outputColumnNames: _col0, _col1 + ListSink + +PREHOOK: query: SELECT id, data FROM variant_proj +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_proj +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT id, data FROM variant_proj +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_proj +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"age":30,"name":"John"} +2 [1,2,3] +3 "plain string" +4 42 +5 NULL +PREHOOK: query: SELECT * FROM variant_proj +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_proj +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT * FROM variant_proj +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_proj +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"age":30,"name":"John"} +2 [1,2,3] +3 "plain string" +4 42 +5 NULL +PREHOOK: query: SELECT id, data FROM variant_proj ORDER BY id DESC +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_proj +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT id, data FROM variant_proj ORDER BY id DESC +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_proj +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"age":30,"name":"John"} +2 [1,2,3] +3 "plain string" +4 42 +5 NULL +PREHOOK: query: drop table if exists variant_nested +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists variant_nested +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE variant_nested ( + id INT, + s STRUCT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +) +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_nested +POSTHOOK: query: CREATE EXTERNAL TABLE variant_nested ( + id INT, + s STRUCT +) STORED BY ICEBERG +TBLPROPERTIES ( + 'format-version'='3' +) +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_nested +PREHOOK: query: INSERT INTO variant_nested SELECT + 1, + named_struct('label', 'x', 'v', parse_json('{"k": 1}')) +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@variant_nested +POSTHOOK: query: INSERT INTO variant_nested SELECT + 1, + named_struct('label', 'x', 'v', parse_json('{"k": 1}')) +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@variant_nested +PREHOOK: query: SELECT * FROM variant_nested +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_nested +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT * FROM variant_nested +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_nested +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"label":"x","v":{"k":1}} +PREHOOK: query: SELECT id, to_json(s.v) FROM variant_nested +PREHOOK: type: QUERY +PREHOOK: Input: default@variant_nested +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT id, to_json(s.v) FROM variant_nested +POSTHOOK: type: QUERY +POSTHOOK: Input: default@variant_nested +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"k":1} +PREHOOK: query: drop table if exists bin_struct +PREHOOK: type: DROPTABLE +PREHOOK: Output: database:default +POSTHOOK: query: drop table if exists bin_struct +POSTHOOK: type: DROPTABLE +POSTHOOK: Output: database:default +PREHOOK: query: CREATE EXTERNAL TABLE bin_struct ( + id INT, + s STRUCT +) STORED BY ICEBERG +PREHOOK: type: CREATETABLE +PREHOOK: Output: database:default +PREHOOK: Output: default@bin_struct +POSTHOOK: query: CREATE EXTERNAL TABLE bin_struct ( + id INT, + s STRUCT +) STORED BY ICEBERG +POSTHOOK: type: CREATETABLE +POSTHOOK: Output: database:default +POSTHOOK: Output: default@bin_struct +PREHOOK: query: INSERT INTO bin_struct SELECT 1, named_struct('metadata', cast('x' as binary), 'value', cast('y' as binary)) +PREHOOK: type: QUERY +PREHOOK: Input: _dummy_database@_dummy_table +PREHOOK: Output: default@bin_struct +POSTHOOK: query: INSERT INTO bin_struct SELECT 1, named_struct('metadata', cast('x' as binary), 'value', cast('y' as binary)) +POSTHOOK: type: QUERY +POSTHOOK: Input: _dummy_database@_dummy_table +POSTHOOK: Output: default@bin_struct +PREHOOK: query: SELECT * FROM bin_struct +PREHOOK: type: QUERY +PREHOOK: Input: default@bin_struct +PREHOOK: Output: hdfs://### HDFS PATH ### +POSTHOOK: query: SELECT * FROM bin_struct +POSTHOOK: type: QUERY +POSTHOOK: Input: default@bin_struct +POSTHOOK: Output: hdfs://### HDFS PATH ### +1 {"metadata":x,"value":y} +PREHOOK: query: drop table variant_proj +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_proj +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_proj +POSTHOOK: query: drop table variant_proj +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_proj +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_proj +PREHOOK: query: drop table variant_nested +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@variant_nested +PREHOOK: Output: database:default +PREHOOK: Output: default@variant_nested +POSTHOOK: query: drop table variant_nested +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@variant_nested +POSTHOOK: Output: database:default +POSTHOOK: Output: default@variant_nested +PREHOOK: query: drop table bin_struct +PREHOOK: type: DROPTABLE +PREHOOK: Input: default@bin_struct +PREHOOK: Output: database:default +PREHOOK: Output: default@bin_struct +POSTHOOK: query: drop table bin_struct +POSTHOOK: type: DROPTABLE +POSTHOOK: Input: default@bin_struct +POSTHOOK: Output: database:default +POSTHOOK: Output: default@bin_struct diff --git a/jdbc/src/java/org/apache/hive/jdbc/JdbcColumn.java b/jdbc/src/java/org/apache/hive/jdbc/JdbcColumn.java index 2747ae751d46..abcaca99a58f 100644 --- a/jdbc/src/java/org/apache/hive/jdbc/JdbcColumn.java +++ b/jdbc/src/java/org/apache/hive/jdbc/JdbcColumn.java @@ -170,6 +170,8 @@ static Type typeStringToHiveType(String type) throws SQLException { return Type.NULL_TYPE; case serdeConstants.UNKNOWN_TYPE_NAME: return Type.UNKNOWN_TYPE; + case serdeConstants.VARIANT_TYPE_NAME: + return Type.VARIANT_TYPE; default: throw new SQLException("Unrecognized column type: " + type); } @@ -228,6 +230,8 @@ static String getColumnTypeName(String type) throws SQLException { return serdeConstants.VOID_TYPE_NAME; case serdeConstants.UNKNOWN_TYPE_NAME: return serdeConstants.UNKNOWN_TYPE_NAME; + case serdeConstants.VARIANT_TYPE_NAME: + return serdeConstants.VARIANT_TYPE_NAME; case "map": return serdeConstants.MAP_TYPE_NAME; case "array": diff --git a/serde/src/java/org/apache/hadoop/hive/serde2/DelimitedJSONSerDe.java b/serde/src/java/org/apache/hadoop/hive/serde2/DelimitedJSONSerDe.java index 76c20b43b350..f7536fa15fdd 100644 --- a/serde/src/java/org/apache/hadoop/hive/serde2/DelimitedJSONSerDe.java +++ b/serde/src/java/org/apache/hadoop/hive/serde2/DelimitedJSONSerDe.java @@ -58,7 +58,7 @@ protected void serializeField(ByteStream.Output out, Object obj, ObjectInspector if (!objInspector.getCategory().equals(Category.PRIMITIVE) || (objInspector.getTypeName().equalsIgnoreCase(serdeConstants.BINARY_TYPE_NAME))) { //do this for all complex types and binary try { - serialize(out, SerDeUtils.getJSONString(obj, objInspector, serdeParams.getNullSequence().toString()), + serialize(out, SerDeUtils.getJSONString(obj, objInspector, serdeParams.getNullSequence().toString(), true), PrimitiveObjectInspectorFactory.javaStringObjectInspector, serdeParams.getSeparators(), 1, serdeParams.getNullSequence(), serdeParams.isEscaped(), serdeParams.getEscapeChar(), serdeParams.getNeedsEscape()); diff --git a/serde/src/java/org/apache/hadoop/hive/serde2/SerDeUtils.java b/serde/src/java/org/apache/hadoop/hive/serde2/SerDeUtils.java index 849425ea395a..322b3e92a2f4 100644 --- a/serde/src/java/org/apache/hadoop/hive/serde2/SerDeUtils.java +++ b/serde/src/java/org/apache/hadoop/hive/serde2/SerDeUtils.java @@ -19,6 +19,7 @@ package org.apache.hadoop.hive.serde2; import java.nio.charset.Charset; +import java.time.ZoneOffset; import java.util.List; import java.util.Map; import java.util.Properties; @@ -31,6 +32,7 @@ import org.apache.hadoop.hive.serde2.objectinspector.StructField; import org.apache.hadoop.hive.serde2.objectinspector.StructObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.UnionObjectInspector; +import org.apache.hadoop.hive.serde2.objectinspector.VariantObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.primitive.BinaryObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.primitive.BooleanObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.primitive.ByteObjectInspector; @@ -48,6 +50,7 @@ import org.apache.hadoop.hive.serde2.objectinspector.primitive.StringObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.primitive.TimestampObjectInspector; import org.apache.hadoop.hive.serde2.objectinspector.primitive.TimestampLocalTZObjectInspector; +import org.apache.hadoop.hive.serde2.variant.Variant; import org.apache.hadoop.io.BytesWritable; import org.apache.hadoop.io.Text; import org.slf4j.Logger; @@ -72,6 +75,8 @@ public final class SerDeUtils { // lower case null is used within json objects private static final String JSON_NULL = "null"; + // the physical struct shape identifying variant values at runtime + private static final String VARIANT_TYPE_NAME = VariantObjectInspector.get().getTypeName(); public static final String LIST_SINK_OUTPUT_FORMATTER = "list.sink.output.formatter"; public static final String LIST_SINK_OUTPUT_PROTOCOL = "list.sink.output.protocol"; public static final Logger LOG = LoggerFactory.getLogger(SerDeUtils.class.getName()); @@ -180,7 +185,7 @@ public static Object toThriftPayload(Object val, ObjectInspector valOI, int vers } // for now, expose non-primitive as a string // TODO: expose non-primitive as a structured object while maintaining JDBC compliance - return SerDeUtils.getJSONString(val, valOI); + return SerDeUtils.getJSONString(val, valOI, JSON_NULL, true); } public static String getJSONString(Object o, ObjectInspector oi) { @@ -197,13 +202,23 @@ public static String getJSONString(Object o, ObjectInspector oi) { * @return */ public static String getJSONString(Object o, ObjectInspector oi, String nullStr) { + return getJSONString(o, oi, nullStr, false); + } + + /** + * Use this for client-facing output: variant values, recognized by their physical + * {@code struct} shape, are rendered as their logical JSON + * value instead of the raw metadata/value bytes. + */ + public static String getJSONString(Object o, ObjectInspector oi, String nullStr, boolean decodeVariant) { StringBuilder sb = new StringBuilder(); - buildJSONString(sb, o, oi, nullStr); + buildJSONString(sb, o, oi, nullStr, decodeVariant); return sb.toString(); } - static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, String nullStr) { + static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, String nullStr, + boolean decodeVariant) { switch (oi.getCategory()) { case PRIMITIVE: { PrimitiveObjectInspector poi = (PrimitiveObjectInspector) oi; @@ -321,7 +336,7 @@ static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, Stri if (i > 0) { sb.append(COMMA); } - buildJSONString(sb, olist.get(i), listElementObjectInspector, JSON_NULL); + buildJSONString(sb, olist.get(i), listElementObjectInspector, JSON_NULL, decodeVariant); } sb.append(RBRACKET); } @@ -345,9 +360,9 @@ static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, Stri sb.append(COMMA); } Map.Entry e = (Map.Entry) entry; - buildJSONString(sb, e.getKey(), mapKeyObjectInspector, JSON_NULL); + buildJSONString(sb, e.getKey(), mapKeyObjectInspector, JSON_NULL, decodeVariant); sb.append(COLON); - buildJSONString(sb, e.getValue(), mapValueObjectInspector, JSON_NULL); + buildJSONString(sb, e.getValue(), mapValueObjectInspector, JSON_NULL, decodeVariant); } sb.append(RBRACE); } @@ -355,24 +370,27 @@ static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, Stri } case STRUCT: { StructObjectInspector soi = (StructObjectInspector) oi; - List structFields = soi.getAllStructFieldRefs(); if (o == null) { sb.append(nullStr); - } else { - sb.append(LBRACE); - for (int i = 0; i < structFields.size(); i++) { - if (i > 0) { - sb.append(COMMA); - } - sb.append(QUOTE); - sb.append(structFields.get(i).getFieldName()); - sb.append(QUOTE); - sb.append(COLON); - buildJSONString(sb, soi.getStructFieldData(o, structFields.get(i)), - structFields.get(i).getFieldObjectInspector(), JSON_NULL); + break; + } + if (decodeVariant && appendVariantJson(sb, o, soi)) { + break; + } + List structFields = soi.getAllStructFieldRefs(); + sb.append(LBRACE); + for (int i = 0; i < structFields.size(); i++) { + if (i > 0) { + sb.append(COMMA); } - sb.append(RBRACE); + sb.append(QUOTE); + sb.append(structFields.get(i).getFieldName()); + sb.append(QUOTE); + sb.append(COLON); + buildJSONString(sb, soi.getStructFieldData(o, structFields.get(i)), + structFields.get(i).getFieldObjectInspector(), JSON_NULL, decodeVariant); } + sb.append(RBRACE); break; } case UNION: { @@ -384,7 +402,7 @@ static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, Stri sb.append(uoi.getTag(o)); sb.append(COLON); buildJSONString(sb, uoi.getField(o), - uoi.getObjectInspectors().get(uoi.getTag(o)), JSON_NULL); + uoi.getObjectInspectors().get(uoi.getTag(o)), JSON_NULL, decodeVariant); sb.append(RBRACE); } break; @@ -394,6 +412,33 @@ static void buildJSONString(StringBuilder sb, Object o, ObjectInspector oi, Stri } } + /** + * Renders a variant value as its logical JSON and returns true. Variant values are recognized + * by their physical {@code struct} shape, independent of the + * ObjectInspector flavor; returns false without output when the struct has a different shape + * or the bytes are not valid variant encoding, so the caller falls back to the generic struct + * rendering. + */ + private static boolean appendVariantJson(StringBuilder sb, Object o, StructObjectInspector soi) { + if (!VARIANT_TYPE_NAME.equals(soi.getTypeName())) { + return false; + } + List fields = soi.getAllStructFieldRefs(); + byte[] metadata = ((BinaryObjectInspector) fields.get(0).getFieldObjectInspector()) + .getPrimitiveJavaObject(soi.getStructFieldData(o, fields.get(0))); + byte[] value = ((BinaryObjectInspector) fields.get(1).getFieldObjectInspector()) + .getPrimitiveJavaObject(soi.getStructFieldData(o, fields.get(1))); + if (metadata == null || value == null) { + return false; + } + try { + sb.append(new Variant(value, metadata).toJson(ZoneOffset.UTC)); + return true; + } catch (RuntimeException e) { + return false; + } + } + /** * return false though element is null if nullsafe flag is true for that */ diff --git a/serde/src/java/org/apache/hadoop/hive/serde2/thrift/Type.java b/serde/src/java/org/apache/hadoop/hive/serde2/thrift/Type.java index 6773b8687005..9dd1568da98f 100644 --- a/serde/src/java/org/apache/hadoop/hive/serde2/thrift/Type.java +++ b/serde/src/java/org/apache/hadoop/hive/serde2/thrift/Type.java @@ -112,7 +112,11 @@ public enum Type { true, false), UNKNOWN_TYPE(serdeConstants.UNKNOWN_TYPE_NAME.toUpperCase(), java.sql.Types.NULL, - TTypeId.UNKNOWN_TYPE); + TTypeId.UNKNOWN_TYPE), + VARIANT_TYPE(serdeConstants.VARIANT_TYPE_NAME.toUpperCase(), + java.sql.Types.OTHER, + TTypeId.VARIANT_TYPE, + true, false); private final String name; private final TTypeId tType; @@ -268,6 +272,9 @@ public static Type getType(TypeInfo typeInfo) { case UNKNOWN: { return Type.UNKNOWN_TYPE; } + case VARIANT: { + return Type.VARIANT_TYPE; + } default: { throw new RuntimeException("Unrecognized type: " + typeInfo.getCategory()); } diff --git a/serde/src/test/org/apache/hadoop/hive/serde2/TestSerDeUtils.java b/serde/src/test/org/apache/hadoop/hive/serde2/TestSerDeUtils.java new file mode 100644 index 000000000000..11b883c7f18f --- /dev/null +++ b/serde/src/test/org/apache/hadoop/hive/serde2/TestSerDeUtils.java @@ -0,0 +1,100 @@ +/* + * 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.hadoop.hive.serde2; + +import java.util.Arrays; +import java.util.List; +import java.util.Map; + +import org.apache.hadoop.hive.serde2.objectinspector.ObjectInspector; +import org.apache.hadoop.hive.serde2.thrift.Type; +import org.apache.hadoop.hive.serde2.typeinfo.TypeInfoUtils; +import org.apache.hadoop.hive.serde2.variant.Variant; +import org.apache.hadoop.hive.serde2.variant.VariantBuilder; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class TestSerDeUtils { + + private static ObjectInspector oi(String typeString) { + return TypeInfoUtils.getStandardJavaObjectInspectorFromTypeInfo( + TypeInfoUtils.getTypeInfoFromTypeString(typeString)); + } + + private static List variant(String json) throws Exception { + Variant v = VariantBuilder.parseJson(json, false); + return List.of(v.getMetadata(), v.getValue()); + } + + @Test + public void testGetJSONStringDecodesVariant() throws Exception { + String json = SerDeUtils.getJSONString(variant("{\"name\":\"John\",\"age\":30}"), + oi("struct"), "null", true); + assertEquals("{\"age\":30,\"name\":\"John\"}", json); + } + + @Test + public void testGetJSONStringDecodesNestedVariant() throws Exception { + Object row = Arrays.asList("x", variant("{\"k\":1}")); + String json = SerDeUtils.getJSONString(row, oi("struct"), "null", true); + assertEquals("{\"label\":\"x\",\"v\":{\"k\":1}}", json); + } + + @Test + public void testGetJSONStringDecodesVariantInListAndMap() throws Exception { + Object list = Arrays.asList(variant("1"), null); + assertEquals("[1,null]", SerDeUtils.getJSONString(list, oi("array"), "null", true)); + + Object map = Map.of("k", variant("{\"x\":true}")); + assertEquals("{\"k\":{\"x\":true}}", SerDeUtils.getJSONString(map, oi("map"), "null", true)); + } + + @Test + public void testNonVariantBytesKeepGenericRendering() throws Exception { + Object row = Arrays.asList(new byte[] { 2 }, new byte[] { 65 }); + String json = SerDeUtils.getJSONString(row, oi("struct"), "null", true); + assertEquals("{\"metadata\":\u0002,\"value\":A}", json); + } + + @Test + public void testNullVariantFieldsKeepGenericRendering() throws Exception { + Object row = Arrays.asList(null, null); + String json = SerDeUtils.getJSONString(row, oi("struct"), "null", true); + assertEquals("{\"metadata\":null,\"value\":null}", json); + } + + @Test + public void testVariantIsNotDecodedByDefault() throws Exception { + String json = SerDeUtils.getJSONString(variant("{\"k\":1}"), oi("struct")); + assertTrue(json, json.startsWith("{\"metadata\":")); + } + + @Test + public void testToThriftPayloadDecodesVariant() throws Exception { + Object payload = SerDeUtils.toThriftPayload(variant("{\"k\":1}"), oi("struct"), 10); + assertEquals("{\"k\":1}", payload); + } + + @Test + public void testThriftTypeForVariant() { + assertEquals(Type.VARIANT_TYPE, Type.getType("variant")); + assertEquals(Type.VARIANT_TYPE, Type.getType(TypeInfoUtils.getTypeInfoFromTypeString("variant"))); + } +} \ No newline at end of file diff --git a/service-rpc/if/TCLIService.thrift b/service-rpc/if/TCLIService.thrift index ee58cd9979d2..3ab61adb718e 100644 --- a/service-rpc/if/TCLIService.thrift +++ b/service-rpc/if/TCLIService.thrift @@ -95,7 +95,8 @@ enum TTypeId { INTERVAL_YEAR_MONTH_TYPE, INTERVAL_DAY_TIME_TYPE, TIMESTAMPLOCALTZ_TYPE, - UNKNOWN_TYPE + UNKNOWN_TYPE, + VARIANT_TYPE } const set PRIMITIVE_TYPES = [ @@ -126,6 +127,7 @@ const set COMPLEX_TYPES = [ TTypeId.STRUCT_TYPE TTypeId.UNION_TYPE TTypeId.USER_DEFINED_TYPE + TTypeId.VARIANT_TYPE ] const set COLLECTION_TYPES = [ @@ -157,6 +159,7 @@ const map TYPE_NAMES = { TTypeId.INTERVAL_DAY_TIME_TYPE: "INTERVAL_DAY_TIME" TTypeId.TIMESTAMPLOCALTZ_TYPE: "TIMESTAMP WITH LOCAL TIME ZONE" TTypeId.UNKNOWN_TYPE: "UNKNOWN" + TTypeId.VARIANT_TYPE: "VARIANT" } // Thrift does not support recursively defined types or forward declarations, diff --git a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_constants.cpp b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_constants.cpp index ae0e3039f30b..583be0f6dfd7 100644 --- a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_constants.cpp +++ b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_constants.cpp @@ -36,6 +36,7 @@ TCLIServiceConstants::TCLIServiceConstants() { COMPLEX_TYPES.insert((TTypeId::type)12); COMPLEX_TYPES.insert((TTypeId::type)13); COMPLEX_TYPES.insert((TTypeId::type)14); + COMPLEX_TYPES.insert((TTypeId::type)24); COLLECTION_TYPES.insert((TTypeId::type)10); COLLECTION_TYPES.insert((TTypeId::type)11); @@ -63,6 +64,7 @@ TCLIServiceConstants::TCLIServiceConstants() { TYPE_NAMES.insert(std::make_pair((TTypeId::type)13, "UNIONTYPE")); TYPE_NAMES.insert(std::make_pair((TTypeId::type)23, "UNKNOWN")); TYPE_NAMES.insert(std::make_pair((TTypeId::type)18, "VARCHAR")); + TYPE_NAMES.insert(std::make_pair((TTypeId::type)24, "VARIANT")); CHARACTER_MAXIMUM_LENGTH = "characterMaximumLength"; diff --git a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.cpp b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.cpp index c8b96da3ed5b..95ac818fa9e8 100644 --- a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.cpp +++ b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.cpp @@ -84,7 +84,8 @@ int _kTTypeIdValues[] = { TTypeId::INTERVAL_YEAR_MONTH_TYPE, TTypeId::INTERVAL_DAY_TIME_TYPE, TTypeId::TIMESTAMPLOCALTZ_TYPE, - TTypeId::UNKNOWN_TYPE + TTypeId::UNKNOWN_TYPE, + TTypeId::VARIANT_TYPE }; const char* _kTTypeIdNames[] = { "BOOLEAN_TYPE", @@ -110,9 +111,10 @@ const char* _kTTypeIdNames[] = { "INTERVAL_YEAR_MONTH_TYPE", "INTERVAL_DAY_TIME_TYPE", "TIMESTAMPLOCALTZ_TYPE", - "UNKNOWN_TYPE" + "UNKNOWN_TYPE", + "VARIANT_TYPE" }; -const std::map _TTypeId_VALUES_TO_NAMES(::apache::thrift::TEnumIterator(24, _kTTypeIdValues, _kTTypeIdNames), ::apache::thrift::TEnumIterator(-1, nullptr, nullptr)); +const std::map _TTypeId_VALUES_TO_NAMES(::apache::thrift::TEnumIterator(25, _kTTypeIdValues, _kTTypeIdNames), ::apache::thrift::TEnumIterator(-1, nullptr, nullptr)); std::ostream& operator<<(std::ostream& out, const TTypeId::type& val) { std::map::const_iterator it = _TTypeId_VALUES_TO_NAMES.find(val); diff --git a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.h b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.h index b80d355e03a9..c8172baee826 100644 --- a/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.h +++ b/service-rpc/src/gen/thrift/gen-cpp/TCLIService_types.h @@ -68,7 +68,8 @@ struct TTypeId { INTERVAL_YEAR_MONTH_TYPE = 20, INTERVAL_DAY_TIME_TYPE = 21, TIMESTAMPLOCALTZ_TYPE = 22, - UNKNOWN_TYPE = 23 + UNKNOWN_TYPE = 23, + VARIANT_TYPE = 24 }; }; diff --git a/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TCLIServiceConstants.java b/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TCLIServiceConstants.java index e525f2cdae32..8fa4b327fb0a 100644 --- a/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TCLIServiceConstants.java +++ b/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TCLIServiceConstants.java @@ -39,6 +39,7 @@ COMPLEX_TYPES.add(org.apache.hive.service.rpc.thrift.TTypeId.STRUCT_TYPE); COMPLEX_TYPES.add(org.apache.hive.service.rpc.thrift.TTypeId.UNION_TYPE); COMPLEX_TYPES.add(org.apache.hive.service.rpc.thrift.TTypeId.USER_DEFINED_TYPE); + COMPLEX_TYPES.add(org.apache.hive.service.rpc.thrift.TTypeId.VARIANT_TYPE); } public static final java.util.Set COLLECTION_TYPES = java.util.EnumSet.noneOf(TTypeId.class); @@ -72,6 +73,7 @@ TYPE_NAMES.put(org.apache.hive.service.rpc.thrift.TTypeId.UNION_TYPE, "UNIONTYPE"); TYPE_NAMES.put(org.apache.hive.service.rpc.thrift.TTypeId.UNKNOWN_TYPE, "UNKNOWN"); TYPE_NAMES.put(org.apache.hive.service.rpc.thrift.TTypeId.VARCHAR_TYPE, "VARCHAR"); + TYPE_NAMES.put(org.apache.hive.service.rpc.thrift.TTypeId.VARIANT_TYPE, "VARIANT"); } public static final java.lang.String CHARACTER_MAXIMUM_LENGTH = "characterMaximumLength"; diff --git a/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TTypeId.java b/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TTypeId.java index 57108255c1c6..0b566dc79457 100644 --- a/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TTypeId.java +++ b/service-rpc/src/gen/thrift/gen-javabean/org/apache/hive/service/rpc/thrift/TTypeId.java @@ -32,7 +32,8 @@ public enum TTypeId implements org.apache.thrift.TEnum { INTERVAL_YEAR_MONTH_TYPE(20), INTERVAL_DAY_TIME_TYPE(21), TIMESTAMPLOCALTZ_TYPE(22), - UNKNOWN_TYPE(23); + UNKNOWN_TYPE(23), + VARIANT_TYPE(24); private final int value; @@ -102,6 +103,8 @@ public static TTypeId findByValue(int value) { return TIMESTAMPLOCALTZ_TYPE; case 23: return UNKNOWN_TYPE; + case 24: + return VARIANT_TYPE; default: return null; } diff --git a/service-rpc/src/gen/thrift/gen-php/Constant.php b/service-rpc/src/gen/thrift/gen-php/Constant.php index 59b6c40636c3..5f45ad5ba616 100644 --- a/service-rpc/src/gen/thrift/gen-php/Constant.php +++ b/service-rpc/src/gen/thrift/gen-php/Constant.php @@ -57,6 +57,7 @@ protected static function init_COMPLEX_TYPES() 12 => true, 13 => true, 14 => true, + 24 => true, ); } @@ -94,6 +95,7 @@ protected static function init_TYPE_NAMES() 13 => "UNIONTYPE", 23 => "UNKNOWN", 18 => "VARCHAR", + 24 => "VARIANT", ); } diff --git a/service-rpc/src/gen/thrift/gen-php/TTypeId.php b/service-rpc/src/gen/thrift/gen-php/TTypeId.php index ac799c4a173d..aa4a6b1a7c2a 100644 --- a/service-rpc/src/gen/thrift/gen-php/TTypeId.php +++ b/service-rpc/src/gen/thrift/gen-php/TTypeId.php @@ -64,6 +64,8 @@ final class TTypeId const UNKNOWN_TYPE = 23; + const VARIANT_TYPE = 24; + static public $__names = array( 0 => 'BOOLEAN_TYPE', 1 => 'TINYINT_TYPE', @@ -89,6 +91,7 @@ final class TTypeId 21 => 'INTERVAL_DAY_TIME_TYPE', 22 => 'TIMESTAMPLOCALTZ_TYPE', 23 => 'UNKNOWN_TYPE', + 24 => 'VARIANT_TYPE', ); } diff --git a/service-rpc/src/gen/thrift/gen-py/TCLIService/constants.py b/service-rpc/src/gen/thrift/gen-py/TCLIService/constants.py index 769dc57c9b8b..61c8107da3ec 100644 --- a/service-rpc/src/gen/thrift/gen-py/TCLIService/constants.py +++ b/service-rpc/src/gen/thrift/gen-py/TCLIService/constants.py @@ -39,6 +39,7 @@ 12, 13, 14, + 24, )) COLLECTION_TYPES = set(( 10, @@ -68,6 +69,7 @@ 13: "UNIONTYPE", 23: "UNKNOWN", 18: "VARCHAR", + 24: "VARIANT", } CHARACTER_MAXIMUM_LENGTH = "characterMaximumLength" PRECISION = "precision" diff --git a/service-rpc/src/gen/thrift/gen-py/TCLIService/ttypes.py b/service-rpc/src/gen/thrift/gen-py/TCLIService/ttypes.py index 8d72578a4970..ddeffbd83358 100644 --- a/service-rpc/src/gen/thrift/gen-py/TCLIService/ttypes.py +++ b/service-rpc/src/gen/thrift/gen-py/TCLIService/ttypes.py @@ -83,6 +83,7 @@ class TTypeId(object): INTERVAL_DAY_TIME_TYPE = 21 TIMESTAMPLOCALTZ_TYPE = 22 UNKNOWN_TYPE = 23 + VARIANT_TYPE = 24 _VALUES_TO_NAMES = { 0: "BOOLEAN_TYPE", @@ -109,6 +110,7 @@ class TTypeId(object): 21: "INTERVAL_DAY_TIME_TYPE", 22: "TIMESTAMPLOCALTZ_TYPE", 23: "UNKNOWN_TYPE", + 24: "VARIANT_TYPE", } _NAMES_TO_VALUES = { @@ -136,6 +138,7 @@ class TTypeId(object): "INTERVAL_DAY_TIME_TYPE": 21, "TIMESTAMPLOCALTZ_TYPE": 22, "UNKNOWN_TYPE": 23, + "VARIANT_TYPE": 24, } diff --git a/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_constants.rb b/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_constants.rb index 8aca726814bf..c1771bd2af12 100644 --- a/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_constants.rb +++ b/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_constants.rb @@ -35,6 +35,7 @@ 12, 13, 14, + 24, ]) COLLECTION_TYPES = Set.new([ @@ -66,6 +67,7 @@ 13 => %q"UNIONTYPE", 23 => %q"UNKNOWN", 18 => %q"VARCHAR", + 24 => %q"VARIANT", } CHARACTER_MAXIMUM_LENGTH = %q"characterMaximumLength" diff --git a/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_types.rb b/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_types.rb index 787ddcd7728e..f17f5d4bb018 100644 --- a/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_types.rb +++ b/service-rpc/src/gen/thrift/gen-rb/t_c_l_i_service_types.rb @@ -47,8 +47,9 @@ module TTypeId INTERVAL_DAY_TIME_TYPE = 21 TIMESTAMPLOCALTZ_TYPE = 22 UNKNOWN_TYPE = 23 - VALUE_MAP = {0 => "BOOLEAN_TYPE", 1 => "TINYINT_TYPE", 2 => "SMALLINT_TYPE", 3 => "INT_TYPE", 4 => "BIGINT_TYPE", 5 => "FLOAT_TYPE", 6 => "DOUBLE_TYPE", 7 => "STRING_TYPE", 8 => "TIMESTAMP_TYPE", 9 => "BINARY_TYPE", 10 => "ARRAY_TYPE", 11 => "MAP_TYPE", 12 => "STRUCT_TYPE", 13 => "UNION_TYPE", 14 => "USER_DEFINED_TYPE", 15 => "DECIMAL_TYPE", 16 => "NULL_TYPE", 17 => "DATE_TYPE", 18 => "VARCHAR_TYPE", 19 => "CHAR_TYPE", 20 => "INTERVAL_YEAR_MONTH_TYPE", 21 => "INTERVAL_DAY_TIME_TYPE", 22 => "TIMESTAMPLOCALTZ_TYPE", 23 => "UNKNOWN_TYPE"} - VALID_VALUES = Set.new([BOOLEAN_TYPE, TINYINT_TYPE, SMALLINT_TYPE, INT_TYPE, BIGINT_TYPE, FLOAT_TYPE, DOUBLE_TYPE, STRING_TYPE, TIMESTAMP_TYPE, BINARY_TYPE, ARRAY_TYPE, MAP_TYPE, STRUCT_TYPE, UNION_TYPE, USER_DEFINED_TYPE, DECIMAL_TYPE, NULL_TYPE, DATE_TYPE, VARCHAR_TYPE, CHAR_TYPE, INTERVAL_YEAR_MONTH_TYPE, INTERVAL_DAY_TIME_TYPE, TIMESTAMPLOCALTZ_TYPE, UNKNOWN_TYPE]).freeze + VARIANT_TYPE = 24 + VALUE_MAP = {0 => "BOOLEAN_TYPE", 1 => "TINYINT_TYPE", 2 => "SMALLINT_TYPE", 3 => "INT_TYPE", 4 => "BIGINT_TYPE", 5 => "FLOAT_TYPE", 6 => "DOUBLE_TYPE", 7 => "STRING_TYPE", 8 => "TIMESTAMP_TYPE", 9 => "BINARY_TYPE", 10 => "ARRAY_TYPE", 11 => "MAP_TYPE", 12 => "STRUCT_TYPE", 13 => "UNION_TYPE", 14 => "USER_DEFINED_TYPE", 15 => "DECIMAL_TYPE", 16 => "NULL_TYPE", 17 => "DATE_TYPE", 18 => "VARCHAR_TYPE", 19 => "CHAR_TYPE", 20 => "INTERVAL_YEAR_MONTH_TYPE", 21 => "INTERVAL_DAY_TIME_TYPE", 22 => "TIMESTAMPLOCALTZ_TYPE", 23 => "UNKNOWN_TYPE", 24 => "VARIANT_TYPE"} + VALID_VALUES = Set.new([BOOLEAN_TYPE, TINYINT_TYPE, SMALLINT_TYPE, INT_TYPE, BIGINT_TYPE, FLOAT_TYPE, DOUBLE_TYPE, STRING_TYPE, TIMESTAMP_TYPE, BINARY_TYPE, ARRAY_TYPE, MAP_TYPE, STRUCT_TYPE, UNION_TYPE, USER_DEFINED_TYPE, DECIMAL_TYPE, NULL_TYPE, DATE_TYPE, VARCHAR_TYPE, CHAR_TYPE, INTERVAL_YEAR_MONTH_TYPE, INTERVAL_DAY_TIME_TYPE, TIMESTAMPLOCALTZ_TYPE, UNKNOWN_TYPE, VARIANT_TYPE]).freeze end module TStatusCode diff --git a/service/src/java/org/apache/hive/service/cli/ColumnValue.java b/service/src/java/org/apache/hive/service/cli/ColumnValue.java index b994d6fc846b..31692853425f 100644 --- a/service/src/java/org/apache/hive/service/cli/ColumnValue.java +++ b/service/src/java/org/apache/hive/service/cli/ColumnValue.java @@ -216,6 +216,7 @@ public static TColumnValue toTColumnValue(TypeDescriptor typeDescriptor, Object case STRUCT_TYPE: case UNION_TYPE: case USER_DEFINED_TYPE: + case VARIANT_TYPE: return stringValue((String)value); case NULL_TYPE: return stringValue((String)value);