This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new cf09697e60 [core] Support DATE, TIMESTAMP and TIMESTAMP_NTZ in variant
shredding schema (#9230)
cf09697e60 is described below
commit cf09697e609f00d1659842d3505760ee5284bc9d
Author: Juntao Zhang <[email protected]>
AuthorDate: Sat Aug 15 20:12:19 2026 +0800
[core] Support DATE, TIMESTAMP and TIMESTAMP_NTZ in variant shredding
schema (#9230)
---
.../paimon/data/variant/BaseVariantReader.java | 12 ++
.../paimon/data/variant/PaimonShreddingUtils.java | 37 +++--
.../apache/paimon/data/variant/ShreddingUtils.java | 8 +-
.../org/apache/paimon/data/variant/VariantGet.java | 13 ++
.../data/variant/GenericVariantBuilderHelper.java | 136 +++++++++++++++++
.../variant/InferVariantShreddingSchemaTest.java | 54 ++++---
.../data/variant/PaimonShreddingUtilsTest.java | 166 ++++++++++++---------
7 files changed, 319 insertions(+), 107 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/variant/BaseVariantReader.java
b/paimon-common/src/main/java/org/apache/paimon/data/variant/BaseVariantReader.java
index 7b8872973f..ca650aa05a 100644
---
a/paimon-common/src/main/java/org/apache/paimon/data/variant/BaseVariantReader.java
+++
b/paimon-common/src/main/java/org/apache/paimon/data/variant/BaseVariantReader.java
@@ -32,13 +32,16 @@ import org.apache.paimon.types.BooleanType;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.DateType;
import org.apache.paimon.types.DecimalType;
import org.apache.paimon.types.DoubleType;
import org.apache.paimon.types.FloatType;
import org.apache.paimon.types.IntType;
+import org.apache.paimon.types.LocalZonedTimestampType;
import org.apache.paimon.types.MapType;
import org.apache.paimon.types.RowType;
import org.apache.paimon.types.SmallIntType;
+import org.apache.paimon.types.TimestampType;
import org.apache.paimon.types.TinyIntType;
import org.apache.paimon.types.VariantType;
@@ -503,6 +506,15 @@ public class BaseVariantReader {
typedValueIdx,
((DecimalType) scalaType).getPrecision(),
((DecimalType) scalaType).getScale());
+ } else if (scalaType instanceof DateType) {
+ i = row.getInt(typedValueIdx);
+ } else if (scalaType instanceof TimestampType) {
+ i = row.getTimestamp(typedValueIdx, ((TimestampType)
scalaType).getPrecision());
+ } else if (scalaType instanceof LocalZonedTimestampType) {
+ i =
+ row.getTimestamp(
+ typedValueIdx,
+ ((LocalZonedTimestampType)
scalaType).getPrecision());
} else {
throw new UnsupportedOperationException("Unsupported scalar
type: " + scalaType);
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/variant/PaimonShreddingUtils.java
b/paimon-common/src/main/java/org/apache/paimon/data/variant/PaimonShreddingUtils.java
index cd7884bc64..6fe539412f 100644
---
a/paimon-common/src/main/java/org/apache/paimon/data/variant/PaimonShreddingUtils.java
+++
b/paimon-common/src/main/java/org/apache/paimon/data/variant/PaimonShreddingUtils.java
@@ -25,6 +25,7 @@ import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.columnar.RowToColumnConverter;
import org.apache.paimon.data.columnar.heap.CastedRowColumnVector;
import org.apache.paimon.data.columnar.writable.WritableBytesVector;
@@ -126,6 +127,11 @@ public class PaimonShreddingUtils {
throw new UnsupportedOperationException();
}
+ @Override
+ public Timestamp getTimestamp(int ordinal, int precision) {
+ return row.getTimestamp(ordinal, precision);
+ }
+
@Override
public ShreddingUtils.ShreddedRow getStruct(int ordinal, int
numFields) {
return new PaimonShreddedRow(row.getRow(ordinal, numFields));
@@ -279,6 +285,9 @@ public class PaimonShreddingUtils {
case BIGINT:
case FLOAT:
case DOUBLE:
+ case DATE:
+ case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+ case TIMESTAMP_WITHOUT_TIME_ZONE:
builder.field(VARIANT_VALUE_FIELD_NAME, DataTypes.BYTES());
builder.field(TYPED_VALUE_FIELD_NAME, dataType);
break;
@@ -318,6 +327,9 @@ public class PaimonShreddingUtils {
typedIdx = i;
switch (field.type().getTypeRoot()) {
case ROW:
+ if (!(dataType instanceof RowType)) {
+ throw invalidVariantShreddingSchema(rowType);
+ }
RowType r = (RowType) dataType;
List<DataField> rFields = r.getFields();
// The struct must not be empty or contain
duplicate field names.
@@ -327,15 +339,11 @@ public class PaimonShreddingUtils {
}
objectSchema = new
VariantSchema.ObjectField[rFields.size()];
for (int index = 0; index < rFields.size();
index++) {
- if (field.type() instanceof RowType) {
- DataField f = rFields.get(index);
- objectSchema[index] =
- new VariantSchema.ObjectField(
- f.name(),
-
buildVariantSchema((RowType) f.type(), false));
- } else {
- throw
invalidVariantShreddingSchema(rowType);
- }
+ DataField f = rFields.get(index);
+ objectSchema[index] =
+ new VariantSchema.ObjectField(
+ f.name(),
+ buildVariantSchema((RowType)
f.type(), false));
}
break;
case ARRAY:
@@ -374,15 +382,23 @@ public class PaimonShreddingUtils {
case DOUBLE:
scalarSchema = new VariantSchema.DoubleType();
break;
+ case CHAR:
case VARCHAR:
scalarSchema = new VariantSchema.StringType();
break;
case BINARY:
+ case VARBINARY:
scalarSchema = new VariantSchema.BinaryType();
break;
case DATE:
scalarSchema = new VariantSchema.DateType();
break;
+ case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+ scalarSchema = new VariantSchema.TimestampType();
+ break;
+ case TIMESTAMP_WITHOUT_TIME_ZONE:
+ scalarSchema = new
VariantSchema.TimestampNTZType();
+ break;
case DECIMAL:
DecimalType d = (DecimalType) dataType;
scalarSchema =
@@ -513,6 +529,9 @@ public class PaimonShreddingUtils {
} else if (schema.scalarSchema instanceof
VariantSchema.DecimalType) {
VariantSchema.DecimalType dt = (VariantSchema.DecimalType)
schema.scalarSchema;
paimonValue = Decimal.fromBigDecimal((BigDecimal) result,
dt.precision, dt.scale);
+ } else if (schema.scalarSchema instanceof
VariantSchema.TimestampType
+ || schema.scalarSchema instanceof
VariantSchema.TimestampNTZType) {
+ paimonValue = Timestamp.fromMicros((Long) result);
} else {
paimonValue = result;
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/variant/ShreddingUtils.java
b/paimon-common/src/main/java/org/apache/paimon/data/variant/ShreddingUtils.java
index f87d5f6178..511d014b18 100644
---
a/paimon-common/src/main/java/org/apache/paimon/data/variant/ShreddingUtils.java
+++
b/paimon-common/src/main/java/org/apache/paimon/data/variant/ShreddingUtils.java
@@ -18,6 +18,8 @@
package org.apache.paimon.data.variant;
+import org.apache.paimon.data.Timestamp;
+
import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.UUID;
@@ -61,6 +63,8 @@ public class ShreddingUtils {
UUID getUuid(int ordinal);
+ Timestamp getTimestamp(int ordinal, int precision);
+
ShreddedRow getStruct(int ordinal, int numFields);
ShreddedRow getArray(int ordinal);
@@ -133,10 +137,10 @@ public class ShreddingUtils {
} else if (scalar instanceof VariantSchema.DateType) {
builder.appendDate(row.getInt(typedIdx));
} else if (scalar instanceof VariantSchema.TimestampType) {
- builder.appendTimestamp(row.getLong(typedIdx));
+ builder.appendTimestamp(row.getTimestamp(typedIdx,
6).toMicros());
} else {
assert scalar instanceof VariantSchema.TimestampNTZType;
- builder.appendTimestampNtz(row.getLong(typedIdx));
+ builder.appendTimestampNtz(row.getTimestamp(typedIdx,
6).toMicros());
}
} else if (schema.arraySchema != null) {
VariantSchema elementSchema = schema.arraySchema;
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantGet.java
b/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantGet.java
index 51233624f9..e25c6df5bd 100644
--- a/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantGet.java
+++ b/paimon-common/src/main/java/org/apache/paimon/data/variant/VariantGet.java
@@ -25,6 +25,7 @@ import org.apache.paimon.data.Decimal;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.variant.GenericVariantUtil.Type;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.DataType;
@@ -146,10 +147,22 @@ public class VariantGet {
input = (int) v.getLong();
inputType = DataTypes.DATE();
break;
+ case TIMESTAMP:
+ input = Timestamp.fromMicros(v.getLong());
+ inputType = DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE();
+ break;
+ case TIMESTAMP_NTZ:
+ input = Timestamp.fromMicros(v.getLong());
+ inputType = DataTypes.TIMESTAMP();
+ break;
case FLOAT:
input = v.getFloat();
inputType = DataTypes.FLOAT();
break;
+ case BINARY:
+ input = v.getBinary();
+ inputType = DataTypes.BYTES();
+ break;
default:
// todo: support other types
throw new IllegalArgumentException("Unsupported type: " +
v.getType());
diff --git
a/paimon-common/src/test/java/org/apache/paimon/data/variant/GenericVariantBuilderHelper.java
b/paimon-common/src/test/java/org/apache/paimon/data/variant/GenericVariantBuilderHelper.java
new file mode 100644
index 0000000000..921747850a
--- /dev/null
+++
b/paimon-common/src/test/java/org/apache/paimon/data/variant/GenericVariantBuilderHelper.java
@@ -0,0 +1,136 @@
+/*
+ * 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.paimon.data.variant;
+
+import org.apache.paimon.types.ArrayType;
+import org.apache.paimon.types.DataField;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.RowType;
+
+import java.math.BigDecimal;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+/** Test helper for building {@link GenericVariant} values from typed Java
objects. */
+public final class GenericVariantBuilderHelper {
+
+ private GenericVariantBuilderHelper() {}
+
+ public static GenericVariant build(RowType rowType, Map<String, Object>
values) {
+ GenericVariantBuilder builder = new GenericVariantBuilder(false);
+ appendValue(builder, rowType, values);
+ return builder.result();
+ }
+
+ /** Build a variant object from a RowType schema and a map of field
values. */
+ private static void appendObject(
+ GenericVariantBuilder builder, RowType rowType, Map<String,
Object> values) {
+ int start = builder.getWritePos();
+ ArrayList<GenericVariantBuilder.FieldEntry> fields = new ArrayList<>();
+ for (DataField field : rowType.getFields()) {
+ String key = field.name();
+ if (!values.containsKey(key)) {
+ continue;
+ }
+ fields.add(
+ new GenericVariantBuilder.FieldEntry(
+ key, builder.addKey(key), builder.getWritePos() -
start));
+ Object value = values.get(key);
+ if (value == null) {
+ builder.appendNull();
+ } else {
+ appendValue(builder, field.type(), value);
+ }
+ }
+ builder.finishWritingObject(start, fields);
+ }
+
+ /** Build a variant array from an element type and a list of element
values. */
+ private static void appendArray(
+ GenericVariantBuilder builder, DataType elementType, List<Object>
values) {
+ int start = builder.getWritePos();
+ ArrayList<Integer> offsets = new ArrayList<>();
+ for (Object value : values) {
+ offsets.add(builder.getWritePos() - start);
+ if (value == null) {
+ builder.appendNull();
+ } else {
+ appendValue(builder, elementType, value);
+ }
+ }
+ builder.finishWritingArray(start, offsets);
+ }
+
+ @SuppressWarnings("unchecked")
+ private static void appendValue(
+ GenericVariantBuilder builder, DataType dataType, Object value) {
+ switch (dataType.getTypeRoot()) {
+ case VARCHAR:
+ case CHAR:
+ builder.appendString((String) value);
+ break;
+ case TINYINT:
+ builder.appendLong((Byte) value);
+ break;
+ case SMALLINT:
+ builder.appendLong((Short) value);
+ break;
+ case INTEGER:
+ builder.appendLong((Integer) value);
+ break;
+ case BIGINT:
+ builder.appendLong((Long) value);
+ break;
+ case FLOAT:
+ builder.appendFloat((Float) value);
+ break;
+ case DOUBLE:
+ builder.appendDouble((Double) value);
+ break;
+ case BOOLEAN:
+ builder.appendBoolean((Boolean) value);
+ break;
+ case BINARY:
+ case VARBINARY:
+ builder.appendBinary((byte[]) value);
+ break;
+ case DECIMAL:
+ builder.appendDecimal((BigDecimal) value);
+ break;
+ case DATE:
+ builder.appendDate((Integer) value);
+ break;
+ case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+ builder.appendTimestamp((Long) value);
+ break;
+ case TIMESTAMP_WITHOUT_TIME_ZONE:
+ builder.appendTimestampNtz((Long) value);
+ break;
+ case ROW:
+ appendObject(builder, (RowType) dataType, (Map<String,
Object>) value);
+ break;
+ case ARRAY:
+ appendArray(builder, ((ArrayType) dataType).getElementType(),
(List<Object>) value);
+ break;
+ default:
+ throw new IllegalArgumentException("Unsupported type: " +
dataType);
+ }
+ }
+}
diff --git
a/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
b/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
index ad541b0157..f5d02342a2 100644
---
a/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/data/variant/InferVariantShreddingSchemaTest.java
@@ -27,9 +27,13 @@ import org.apache.paimon.types.RowType;
import org.junit.jupiter.api.Test;
+import java.math.BigDecimal;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import static
org.apache.paimon.data.variant.PaimonShreddingUtils.variantShreddingSchema;
import static org.assertj.core.api.Assertions.assertThat;
@@ -441,35 +445,39 @@ public class InferVariantShreddingSchemaTest {
// Schema: row<v: variant>
RowType schema = RowType.of(new DataType[] {DataTypes.VARIANT()}, new
String[] {"v"});
- String json =
- "{"
- + "\"string\": \"test\", "
- + "\"long\": 123456789, "
- + "\"double\": 3.14159, "
- + "\"boolean\": true, "
- + "\"null\": null"
- + "}";
-
- GenericVariant variant = GenericVariant.fromJson(json);
+ DataField f1 = new DataField(0, "binary", DataTypes.BYTES());
+ DataField f2 = new DataField(1, "boolean", DataTypes.BOOLEAN());
+ DataField f3 = new DataField(2, "date", DataTypes.DATE());
+ DataField f4 = new DataField(3, "decimal", DataTypes.DECIMAL(18, 5));
+ DataField f5 = new DataField(4, "double", DataTypes.DOUBLE());
+ DataField f6 = new DataField(5, "float", DataTypes.FLOAT());
+ DataField f7 = new DataField(6, "long", DataTypes.BIGINT());
+ DataField f8 = new DataField(7, "null", DataTypes.VARIANT());
+ DataField f9 = new DataField(8, "string", DataTypes.STRING());
+ DataField f10 = new DataField(9, "timestamp",
DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE());
+ DataField f11 = new DataField(10, "timestampntz",
DataTypes.TIMESTAMP());
+ RowType objectType = RowType.of(f1, f2, f3, f4, f5, f6, f7, f8, f9,
f10, f11);
+
+ Map<String, Object> values = new HashMap<>();
+ values.put("string", "test");
+ values.put("long", 123456789L);
+ values.put("decimal", new BigDecimal("3.14159"));
+ values.put("double", 1.0123456789012345678901234567890123456789D);
+ values.put("boolean", true);
+ values.put("null", null);
+ values.put("date", 20000);
+ values.put("timestamp", 1_234_567_890_123_456L);
+ values.put("timestampntz", 9_876_543_210_123_456L);
+ values.put("float", 3.14f);
+ values.put("binary", "bytes".getBytes(StandardCharsets.UTF_8));
+ GenericVariant variant = GenericVariantBuilderHelper.build(objectType,
values);
List<InternalRow> rows = Arrays.asList(GenericRow.of(variant));
InferVariantShreddingSchema inferrer =
defaultInferVariantShreddingSchema(schema);
RowType inferredSchema = inferrer.inferSchema(rows);
- // All primitive types: boolean, decimal (3.14159 becomes DECIMAL),
bigint, variant (null),
- // string
- RowType expectedType =
- RowType.of(
- new DataType[] {
- DataTypes.BOOLEAN(),
- DataTypes.DECIMAL(18, 5),
- DataTypes.BIGINT(),
- DataTypes.VARIANT(),
- DataTypes.STRING()
- },
- new String[] {"boolean", "double", "long", "null",
"string"});
assertThat(inferredSchema.getField("v").type())
- .isEqualTo(variantShreddingSchema(expectedType));
+ .isEqualTo(variantShreddingSchema(objectType));
}
@Test
diff --git
a/paimon-common/src/test/java/org/apache/paimon/data/variant/PaimonShreddingUtilsTest.java
b/paimon-common/src/test/java/org/apache/paimon/data/variant/PaimonShreddingUtilsTest.java
index d918045702..d3835b6247 100644
---
a/paimon-common/src/test/java/org/apache/paimon/data/variant/PaimonShreddingUtilsTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/data/variant/PaimonShreddingUtilsTest.java
@@ -23,18 +23,29 @@ import org.apache.paimon.data.Decimal;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.Timestamp;
import org.apache.paimon.data.variant.PaimonShreddingUtils.FieldToExtract;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.MapType;
import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.DateTimeUtils;
+
+import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableMap;
import org.junit.jupiter.api.Test;
+import java.math.BigDecimal;
+import java.nio.charset.StandardCharsets;
import java.time.ZoneOffset;
+import java.util.Arrays;
import java.util.HashMap;
+import java.util.Map;
+import java.util.TimeZone;
+import static
org.apache.paimon.data.variant.PaimonShreddingUtils.assembleVariant;
import static
org.apache.paimon.data.variant.PaimonShreddingUtils.assembleVariantStruct;
import static
org.apache.paimon.data.variant.PaimonShreddingUtils.buildFieldsToExtract;
import static
org.apache.paimon.data.variant.PaimonShreddingUtils.buildVariantSchema;
@@ -67,28 +78,33 @@ public class PaimonShreddingUtilsTest {
DataField f9 = new DataField(9, "decimal", DataTypes.DECIMAL(5, 2));
DataField f10 = new DataField(10, "boolean", DataTypes.BOOLEAN());
DataField f11 = new DataField(11, "nullField", DataTypes.INT());
- RowType allTypes = RowType.of(f1, f2, f3, f4, f5, f6, f7, f8, f9, f10,
f11);
-
- String json =
- "{\n"
- + " \"object\": {\n"
- + " \"name\": \"Apache Paimon\",\n"
- + " \"age\": 3\n"
- + " },\n"
- + " \"array\": [1, 2, 3, 4, 5],\n"
- + " \"string\": \"Hello, World!\",\n"
- + " \"tinyint\": 1,\n"
- + " \"smallint\": 3000,\n"
- + " \"int\": 400000,\n"
- + " \"long\": 12345678901234,\n"
- + " \"double\":
1.0123456789012345678901234567890123456789,\n"
- + " \"decimal\": 100.99,\n"
- + " \"boolean\": true,\n"
- + " \"nullField\": null\n"
- + "}\n";
-
- GenericVariant v = GenericVariant.fromJson(json);
- GenericRow expert =
+ DataField f12 = new DataField(12, "date", DataTypes.DATE());
+ DataField f13 = new DataField(13, "timestamp",
DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE());
+ DataField f14 = new DataField(14, "timestampntz",
DataTypes.TIMESTAMP());
+ DataField f15 = new DataField(15, "float", DataTypes.FLOAT());
+ DataField f16 = new DataField(16, "binary", DataTypes.BYTES());
+ RowType allTypes =
+ RowType.of(f1, f2, f3, f4, f5, f6, f7, f8, f9, f10, f11, f12,
f13, f14, f15, f16);
+
+ Map<String, Object> values = new HashMap<>();
+ values.put("object", ImmutableMap.of("name", "Apache Paimon", "age",
3));
+ values.put("array", Arrays.asList(1, 2, 3, 4, 5));
+ values.put("string", "Hello, World!");
+ values.put("tinyint", (byte) 1);
+ values.put("smallint", (short) 3000);
+ values.put("int", 400000);
+ values.put("long", 12345678901234L);
+ values.put("double", 1.0123456789012345678901234567890123456789D);
+ values.put("decimal", new BigDecimal("100.99"));
+ values.put("boolean", true);
+ values.put("nullField", null);
+ values.put("date", 20000);
+ values.put("timestamp", 1_234_567_890_123_456L);
+ values.put("timestampntz", 9_876_543_210_123_456L);
+ values.put("float", 3.14f);
+ values.put("binary", "bytes".getBytes(StandardCharsets.UTF_8));
+ GenericVariant v = GenericVariantBuilderHelper.build(allTypes, values);
+ GenericRow expected =
GenericRow.of(
GenericRow.of(BinaryString.fromString("Apache
Paimon"), 3),
new GenericArray(new Integer[] {1, 2, 3, 4, 5}),
@@ -100,63 +116,29 @@ public class PaimonShreddingUtilsTest {
1.0123456789012345678901234567890123456789D,
Decimal.fromBigDecimal(new
java.math.BigDecimal("100.99"), 5, 2),
true,
- null);
+ null,
+ 20000,
+ Timestamp.fromMicros(1_234_567_890_123_456L),
+ Timestamp.fromMicros(9_876_543_210_123_456L),
+ 3.14f,
+ "bytes".getBytes(StandardCharsets.UTF_8));
// shredding to real type
- RowType shreddedType = new RowType(allTypes.getFields());
- RowType shreddingSchema = variantShreddingSchema(shreddedType);
- VariantSchema variantSchema = buildVariantSchema(shreddingSchema);
- FieldToExtract[] fields = new FieldToExtract[allTypes.getFieldCount()];
- for (int i = 0; i < allTypes.getFields().size(); i++) {
- fields[i] =
- buildFieldsToExtract(
- allTypes.getFields().get(i).type(),
- "$." + allTypes.getFields().get(i).name(),
- castArgs,
- variantSchema);
- }
- assertThat(assembleVariantStruct(castShredded(v, variantSchema),
variantSchema, fields))
- .isEqualTo(expert);
+ assertVariantStructEquals(new RowType(allTypes.getFields()), allTypes,
v, expected);
// no shredding
- shreddedType = RowType.of();
- shreddingSchema = variantShreddingSchema(shreddedType);
- variantSchema = buildVariantSchema(shreddingSchema);
- fields = new FieldToExtract[allTypes.getFieldCount()];
- for (int i = 0; i < allTypes.getFields().size(); i++) {
- fields[i] =
- buildFieldsToExtract(
- allTypes.getFields().get(i).type(),
- "$." + allTypes.getFields().get(i).name(),
- castArgs,
- variantSchema);
- }
-
- assertThat(assembleVariantStruct(castShredded(v, variantSchema),
variantSchema, fields))
- .isEqualTo(expert);
+ assertVariantStructEquals(RowType.of(), allTypes, v, expected);
// shredding to string, then cast to the real type
- shreddedType =
+ RowType shreddedType =
RowType.of(
allTypes.getFields().stream()
.map(a -> a.newType(DataTypes.STRING()))
.toArray(DataField[]::new));
- shreddingSchema = variantShreddingSchema(shreddedType);
- variantSchema = buildVariantSchema(shreddingSchema);
- fields = new FieldToExtract[allTypes.getFieldCount()];
- for (int i = 0; i < allTypes.getFields().size(); i++) {
- fields[i] =
- buildFieldsToExtract(
- allTypes.getFields().get(i).type(),
- "$." + allTypes.getFields().get(i).name(),
- castArgs,
- variantSchema);
- }
- assertThat(assembleVariantStruct(castShredded(v, variantSchema),
variantSchema, fields))
- .isEqualTo(expert);
+ assertVariantStructEquals(shreddedType, allTypes, v, expected);
// shredding to real type, then cast to the string
- expert =
+ expected =
GenericRow.of(
BinaryString.fromString("{\"age\":3,\"name\":\"Apache
Paimon\"}"),
BinaryString.fromString("[1,2,3,4,5]"),
@@ -168,12 +150,25 @@ public class PaimonShreddingUtilsTest {
BinaryString.fromString("1.0123456789012346"),
BinaryString.fromString("100.99"),
BinaryString.fromString("true"),
- null);
+ null,
+
BinaryString.fromString(DateTimeUtils.formatDate(20000)),
+ BinaryString.fromString(
+ DateTimeUtils.formatTimestamp(
+
Timestamp.fromMicros(1_234_567_890_123_456L),
+ TimeZone.getDefault(),
+ 6)),
+ BinaryString.fromString(
+ DateTimeUtils.formatTimestamp(
+
Timestamp.fromMicros(9_876_543_210_123_456L),
+ DateTimeUtils.UTC_ZONE,
+ 6)),
+ BinaryString.fromString("3.14"),
+ BinaryString.fromString("bytes"));
shreddedType = new RowType(allTypes.getFields());
- shreddingSchema = variantShreddingSchema(shreddedType);
- variantSchema = buildVariantSchema(shreddingSchema);
- fields = new FieldToExtract[allTypes.getFieldCount()];
+ RowType shreddingSchema = variantShreddingSchema(shreddedType);
+ VariantSchema variantSchema = buildVariantSchema(shreddingSchema);
+ FieldToExtract[] fields = new FieldToExtract[allTypes.getFieldCount()];
for (int i = 0; i < allTypes.getFields().size(); i++) {
fields[i] =
buildFieldsToExtract(
@@ -183,7 +178,7 @@ public class PaimonShreddingUtilsTest {
variantSchema);
}
assertThat(assembleVariantStruct(castShredded(v, variantSchema),
variantSchema, fields))
- .isEqualTo(expert);
+ .isEqualTo(expected);
// no shredding, then cast to the string
shreddedType = RowType.of();
@@ -200,7 +195,7 @@ public class PaimonShreddingUtilsTest {
}
assertThat(assembleVariantStruct(castShredded(v, variantSchema),
variantSchema, fields))
- .isEqualTo(expert);
+ .isEqualTo(expected);
// cast struct to map
shreddedType = RowType.of(f1);
@@ -229,6 +224,31 @@ public class PaimonShreddingUtilsTest {
})));
}
+ private static void assertVariantStructEquals(
+ RowType shreddedType, RowType allTypes, GenericVariant v,
GenericRow expected) {
+ VariantCastArgs castArgs = new VariantCastArgs(true, ZoneOffset.UTC);
+
+ RowType shreddingSchema = variantShreddingSchema(shreddedType);
+ VariantSchema variantSchema = buildVariantSchema(shreddingSchema);
+
+ FieldToExtract[] fieldsToExtract = new
FieldToExtract[allTypes.getFieldCount()];
+ for (int i = 0; i < allTypes.getFields().size(); i++) {
+ fieldsToExtract[i] =
+ buildFieldsToExtract(
+ allTypes.getFields().get(i).type(),
+ "$." + allTypes.getFields().get(i).name(),
+ castArgs,
+ variantSchema);
+ }
+
+ InternalRow inputRow = castShredded(v, variantSchema);
+ InternalRow outputRow = assembleVariantStruct(inputRow, variantSchema,
fieldsToExtract);
+ assertThat(outputRow).isEqualTo(expected);
+
+ Variant rebuilt = assembleVariant(inputRow, variantSchema);
+
assertThat(rebuilt.toJson(castArgs.zoneId())).isEqualTo(v.toJson(castArgs.zoneId()));
+ }
+
@Test
void testAssembleVariantStructWithShredding() {
VariantCastArgs castArgs = new VariantCastArgs(true, ZoneOffset.UTC);