This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 0451738 [INLONG-2614][Sort] Support array and map data structures in
Hive sink and ClickHouse sink (#2615)
0451738 is described below
commit 0451738536d653d73944f699276ad12f7b88ba60
Author: Kevin Wen <[email protected]>
AuthorDate: Tue Feb 22 19:40:06 2022 +0800
[INLONG-2614][Sort] Support array and map data structures in Hive sink and
ClickHouse sink (#2615)
---
.../flink/clickhouse/ClickHouseRowConverter.java | 70 +++++
.../inlong/sort/flink/hive/HiveSinkHelper.java | 2 +-
.../sort/flink/hive/formats/TextRowWriter.java | 231 +++++++++++++++-
.../hive/formats/parquet/ParquetRowWriter.java | 294 ++++++++++++++++++++-
.../formats/parquet/ParquetRowWriterBuilder.java | 2 +
.../formats/parquet/ParquetSchemaConverter.java | 12 +
.../clickhouse/ClickHouseRowConverterTest.java | 66 +++++
.../sort/flink/hive/formats/TextRowWriterTest.java | 47 +++-
.../formats/parquet/ParquetBulkWriterTest.java | 49 +++-
9 files changed, 755 insertions(+), 18 deletions(-)
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverter.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverter.java
index 08f9756..5e5c4f9 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverter.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverter.java
@@ -18,6 +18,9 @@
package org.apache.inlong.sort.flink.clickhouse;
+import org.apache.commons.lang3.ArrayUtils;
+import
org.apache.flink.shaded.guava18.com.google.common.annotations.VisibleForTesting;
+import org.apache.inlong.sort.formats.common.ArrayTypeInfo;
import org.apache.inlong.sort.formats.common.BooleanTypeInfo;
import org.apache.inlong.sort.formats.common.ByteTypeInfo;
import org.apache.inlong.sort.formats.common.DateTypeInfo;
@@ -27,12 +30,16 @@ import org.apache.inlong.sort.formats.common.FloatTypeInfo;
import org.apache.inlong.sort.formats.common.FormatInfo;
import org.apache.inlong.sort.formats.common.IntTypeInfo;
import org.apache.inlong.sort.formats.common.LongTypeInfo;
+import org.apache.inlong.sort.formats.common.MapTypeInfo;
import org.apache.inlong.sort.formats.common.ShortTypeInfo;
import org.apache.inlong.sort.formats.common.StringTypeInfo;
import org.apache.inlong.sort.formats.common.TimeTypeInfo;
import org.apache.inlong.sort.formats.common.TimestampTypeInfo;
import org.apache.inlong.sort.formats.common.TypeInfo;
import org.apache.flink.types.Row;
+import ru.yandex.clickhouse.ClickHouseArray;
+import ru.yandex.clickhouse.domain.ClickHouseDataType;
+import ru.yandex.clickhouse.util.Utils;
import java.math.BigDecimal;
import java.sql.Date;
@@ -40,6 +47,7 @@ import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.Time;
import java.sql.Timestamp;
+import java.util.Map;
public class ClickHouseRowConverter {
@@ -83,8 +91,70 @@ public class ClickHouseRowConverter {
statement.setDate(index + 1, (Date) value);
} else if (typeInfo instanceof TimestampTypeInfo) {
statement.setTimestamp(index + 1, (Timestamp) value);
+ } else if (typeInfo instanceof ArrayTypeInfo) {
+ TypeInfo elementTypeInfo = ((ArrayTypeInfo)
typeInfo).getElementTypeInfo();
+ statement.setArray(index + 1, new ClickHouseArray(
+ getClickHouseDataTypeFromTypeInfo(elementTypeInfo),
toObjectArray(elementTypeInfo, value)));
+ } else if (typeInfo instanceof MapTypeInfo) {
+ statement.setObject(index + 1,
Utils.mapOf(toKeyValuePairObjectArray(value)));
} else {
throw new IllegalArgumentException("Unsupported TypeInfo " +
typeInfo.getClass().getName());
}
}
+
+ private static ClickHouseDataType
getClickHouseDataTypeFromTypeInfo(TypeInfo typeInfo) {
+ if (typeInfo instanceof StringTypeInfo) {
+ return ClickHouseDataType.String;
+ } else if (typeInfo instanceof BooleanTypeInfo || typeInfo instanceof
ByteTypeInfo) {
+ return ClickHouseDataType.Int8;
+ } else if (typeInfo instanceof ShortTypeInfo) {
+ return ClickHouseDataType.Int16;
+ } else if (typeInfo instanceof IntTypeInfo) {
+ return ClickHouseDataType.Int32;
+ } else if (typeInfo instanceof LongTypeInfo) {
+ return ClickHouseDataType.Int64;
+ } else if (typeInfo instanceof FloatTypeInfo) {
+ return ClickHouseDataType.Float32;
+ } else if (typeInfo instanceof DoubleTypeInfo) {
+ return ClickHouseDataType.Float64;
+ } else {
+ throw new IllegalArgumentException("Unsupported TypeInfo " +
typeInfo.getClass().getName());
+ }
+ }
+
+ @VisibleForTesting
+ static Object toObjectArray(TypeInfo typeInfo, Object object) {
+ if (typeInfo instanceof BooleanTypeInfo && object instanceof
boolean[]) {
+ return ArrayUtils.toObject((boolean[]) object);
+ } else if (typeInfo instanceof ByteTypeInfo && object instanceof
byte[]) {
+ return ArrayUtils.toObject((byte[]) object);
+ } else if (typeInfo instanceof ShortTypeInfo && object instanceof
short[]) {
+ return ArrayUtils.toObject((short[]) object);
+ } else if (typeInfo instanceof IntTypeInfo && object instanceof int[])
{
+ return ArrayUtils.toObject((int[]) object);
+ } else if (typeInfo instanceof LongTypeInfo && object instanceof
long[]) {
+ return ArrayUtils.toObject((long[]) object);
+ } else if (typeInfo instanceof FloatTypeInfo && object instanceof
float[]) {
+ return ArrayUtils.toObject((float[]) object);
+ } else if (typeInfo instanceof DoubleTypeInfo && object instanceof
double[]) {
+ return ArrayUtils.toObject((double[]) object);
+ } else {
+ return object;
+ }
+ }
+
+ @VisibleForTesting
+ static Object[] toKeyValuePairObjectArray(Object input) {
+ Map<?, ?> mapValue = (Map<?, ?>) input;
+ int size = mapValue.size();
+ Object[] kvps = new Object[size * 2];
+ int i = 0;
+ for (Map.Entry<?, ?> entry : mapValue.entrySet()) {
+ kvps[i] = entry.getKey();
+ kvps[i + 1] = entry.getValue();
+ i += 2;
+ }
+
+ return kvps;
+ }
}
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/HiveSinkHelper.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/HiveSinkHelper.java
index f61eb06..7c7e6f9 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/HiveSinkHelper.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/HiveSinkHelper.java
@@ -51,7 +51,7 @@ public class HiveSinkHelper {
return ParquetRowWriterBuilder.createWriterFactory(
rowType, (ParquetFileFormat) hiveFileFormat);
} else if (hiveFileFormat instanceof TextFileFormat) {
- return new TextRowWriter.Factory((TextFileFormat) hiveFileFormat,
config);
+ return new TextRowWriter.Factory((TextFileFormat) hiveFileFormat,
fieldTypes, config);
} else if (hiveFileFormat instanceof OrcFileFormat) {
return OrcBulkWriterFactory.createWriterFactory(rowType,
fieldTypes, config);
} else {
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriter.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriter.java
index 34a551c..70e975c 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriter.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriter.java
@@ -23,13 +23,18 @@ import static
com.google.common.base.Preconditions.checkNotNull;
import java.io.IOException;
import java.io.OutputStream;
import java.nio.charset.StandardCharsets;
+import java.util.Map;
import java.util.zip.GZIPOutputStream;
+
import org.anarres.lzo.LzoAlgorithm;
import org.anarres.lzo.LzoCompressor;
import org.anarres.lzo.LzoLibrary;
import org.anarres.lzo.LzopOutputStream;
import org.apache.flink.api.common.serialization.BulkWriter;
import org.apache.flink.core.fs.FSDataOutputStream;
+import
org.apache.flink.shaded.guava18.com.google.common.annotations.VisibleForTesting;
+import org.apache.flink.table.types.logical.ArrayType;
+import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.types.Row;
import org.apache.inlong.sort.configuration.Configuration;
import org.apache.inlong.sort.configuration.Constants;
@@ -37,25 +42,41 @@ import
org.apache.inlong.sort.protocol.sink.HiveSinkInfo.TextFileFormat;
public class TextRowWriter implements BulkWriter<Row> {
+ private static final String NULL_STRING = "null";
+
+ private static final String ARRAY_SPLITTER = ",";
+ private static final String ARRAY_START_SYMBOL = "[";
+ private static final String ARRAY_END_SYMBOL = "]";
+
+ private static final String MAP_START_SYMBOL = "{";
+ private static final String MAP_END_SYMBOL = "}";
+ private static final String MAP_ENTRY_SPLITTER = ",";
+ private static final String MAP_KEY_VALUE_SPLITTER = "=";
+
private final OutputStream outputStream;
private final TextFileFormat textFileFormat;
private final int bufferSize;
+ private final LogicalType[] fieldTypes;
+
public TextRowWriter(
FSDataOutputStream fsDataOutputStream,
TextFileFormat textFileFormat,
- Configuration config) throws IOException {
+ Configuration config,
+ LogicalType[] fieldTypes) throws IOException {
this.bufferSize =
checkNotNull(config).getInteger(Constants.SINK_HIVE_TEXT_BUFFER_SIZE);
this.outputStream =
getCompressionOutputStream(checkNotNull(fsDataOutputStream), textFileFormat);
this.textFileFormat = checkNotNull(textFileFormat);
+ this.fieldTypes = checkNotNull(fieldTypes);
}
@Override
public void addElement(Row row) throws IOException {
for (int i = 0; i < row.getArity(); i++) {
-
outputStream.write(String.valueOf(row.getField(i)).getBytes(StandardCharsets.UTF_8));
+ String fieldStr = convertField(row.getField(i), fieldTypes[i]);
+ outputStream.write(fieldStr.getBytes(StandardCharsets.UTF_8));
if (i != row.getArity() - 1) {
outputStream.write(textFileFormat.getSplitter());
}
@@ -63,6 +84,200 @@ public class TextRowWriter implements BulkWriter<Row> {
outputStream.write(10); // start a new line
}
+ @VisibleForTesting
+ static String convertField(Object field, LogicalType fieldType) {
+ if (field == null) {
+ return NULL_STRING;
+ }
+
+ switch (fieldType.getTypeRoot()) {
+ case ARRAY:
+ return convertArray(field, ((ArrayType)
fieldType).getElementType());
+ case MAP:
+ return convertMap((Map<?, ?>) field);
+ default:
+ return String.valueOf(field);
+ }
+ }
+
+ private static String convertArray(Object input, LogicalType elementType) {
+ switch (elementType.getTypeRoot()) {
+ case BOOLEAN:
+ return convertBooleanArray(input);
+ case TINYINT:
+ return convertByteArray(input);
+ case SMALLINT:
+ return convertShortArray(input);
+ case INTEGER:
+ return convertIntArray(input);
+ case BIGINT:
+ return convertLongArray(input);
+ case FLOAT:
+ return convertFloatArray(input);
+ case DOUBLE:
+ return convertDoubleArray(input);
+ default:
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertObjectArray(Object[] objArray) {
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < objArray.length; i++) {
+ stringBuilder.append(objArray[i]);
+ if (i != objArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ }
+
+ private static String convertBooleanArray(Object input) {
+ if (input instanceof boolean[]) {
+ boolean[] inputArray = (boolean[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertByteArray(Object input) {
+ if (input instanceof byte[]) {
+ byte[] inputArray = (byte[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertShortArray(Object input) {
+ if (input instanceof short[]) {
+ short[] inputArray = (short[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertIntArray(Object input) {
+ if (input instanceof int[]) {
+ int[] inputArray = (int[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertLongArray(Object input) {
+ if (input instanceof long[]) {
+ long[] inputArray = (long[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertFloatArray(Object input) {
+ if (input instanceof float[]) {
+ float[] inputArray = (float[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertDoubleArray(Object input) {
+ if (input instanceof double[]) {
+ double[] inputArray = (double[]) input;
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(ARRAY_START_SYMBOL);
+ for (int i = 0; i < inputArray.length; i++) {
+ stringBuilder.append(inputArray[i]);
+ if (i != inputArray.length - 1) {
+ stringBuilder.append(ARRAY_SPLITTER);
+ }
+ }
+ stringBuilder.append(ARRAY_END_SYMBOL);
+ return stringBuilder.toString();
+ } else {
+ return convertObjectArray((Object[]) input);
+ }
+ }
+
+ private static String convertMap(Map<?, ?> inputMap) {
+ StringBuilder stringBuilder = new StringBuilder();
+ stringBuilder.append(MAP_START_SYMBOL);
+ int i = 0;
+ int mapSize = inputMap.size();
+ for (Map.Entry<?, ?> entry : inputMap.entrySet()) {
+ ++i;
+ stringBuilder.append(entry.getKey());
+ stringBuilder.append(MAP_KEY_VALUE_SPLITTER);
+ stringBuilder.append(entry.getValue());
+ if (i != mapSize) {
+ stringBuilder.append(MAP_ENTRY_SPLITTER);
+ }
+ }
+ stringBuilder.append(MAP_END_SYMBOL);
+ return stringBuilder.toString();
+ }
+
@Override
public void flush() throws IOException {
outputStream.flush();
@@ -97,19 +312,17 @@ public class TextRowWriter implements BulkWriter<Row> {
private final Configuration config;
- public Factory(
- TextFileFormat textFileFormat,
- Configuration config) {
+ private final LogicalType[] fieldTypes;
+
+ public Factory(TextFileFormat textFileFormat, LogicalType[]
fieldTypes, Configuration config) {
this.textFileFormat = checkNotNull(textFileFormat);
+ this.fieldTypes = checkNotNull(fieldTypes);
this.config = checkNotNull(config);
}
@Override
public BulkWriter<Row> create(FSDataOutputStream fsDataOutputStream)
throws IOException {
- return new TextRowWriter(
- fsDataOutputStream,
- textFileFormat,
- config);
+ return new TextRowWriter(fsDataOutputStream, textFileFormat,
config, fieldTypes);
}
}
}
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriter.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriter.java
index 6252e11..b36ada8 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriter.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriter.java
@@ -20,7 +20,11 @@ package org.apache.inlong.sort.flink.hive.formats.parquet;
import java.sql.Timestamp;
import java.util.Date;
+import java.util.Map;
+
+import org.apache.flink.table.types.logical.ArrayType;
import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.MapType;
import org.apache.flink.table.types.logical.RowType;
import org.apache.flink.types.Row;
import org.apache.parquet.io.api.Binary;
@@ -31,6 +35,12 @@ import org.apache.parquet.schema.Type;
/** Writes a record to the Parquet API with the expected schema in order to be
written to a file. */
public class ParquetRowWriter {
+ static final String ARRAY_FIELD_NAME = "list";
+
+ static final String MAP_ENTITY_FIELD_NAME = "key_value";
+ static final String MAP_KEY_FIELD_NAME = "key";
+ static final String MAP_VALUE_FIELD_NAME = "value";
+
private final RecordConsumer recordConsumer;
private final FieldWriter[] filedWriters;
@@ -104,7 +114,15 @@ public class ParquetRowWriter {
throw new UnsupportedOperationException("Unsupported type:
" + type);
}
} else {
- throw new IllegalArgumentException("Unsupported data type: " + t);
+ switch (t.getTypeRoot()) {
+ case ARRAY:
+ return new ArrayWriter(((ArrayType) t).getElementType(),
+
type.asGroupType().getType(0).asGroupType().getType(0));
+ case MAP:
+ return new MapWriter((MapType) t, type);
+ default:
+ throw new UnsupportedOperationException("Unsupported type:
" + type);
+ }
}
}
@@ -184,4 +202,278 @@ public class ParquetRowWriter {
recordConsumer.addBinary(Binary.fromReusedByteArray((byte[])
row.getField(ordinal)));
}
}
+
+ private class ArrayWriter implements FieldWriter {
+
+ private final LogicalType elementTypeFlink;
+
+ private final Type elementTypeParquet;
+
+ public ArrayWriter(LogicalType elementTypeFlink, Type
elementTypeParquet) {
+ this.elementTypeFlink = elementTypeFlink;
+ this.elementTypeParquet = elementTypeParquet;
+ }
+
+ @Override
+ public void write(Row row, int ordinal) {
+ if (elementTypeParquet.isPrimitive()) {
+ switch (elementTypeFlink.getTypeRoot()) {
+ case CHAR:
+ case VARCHAR:
+ case DECIMAL:
+ case DATE:
+ case TIME_WITHOUT_TIME_ZONE:
+ case TIMESTAMP_WITHOUT_TIME_ZONE:
+ writeObjectArray(row.getField(ordinal));
+ break;
+ case BOOLEAN:
+ writeBooleanArray(row.getField(ordinal));
+ break;
+ case TINYINT:
+ writeTinyIntArray(row.getField(ordinal));
+ break;
+ case SMALLINT:
+ writeShortArray(row.getField(ordinal));
+ break;
+ case INTEGER:
+ writeIntArray(row.getField(ordinal));
+ break;
+ case BIGINT:
+ writeLongArray(row.getField(ordinal));
+ break;
+ case FLOAT:
+ writeFloatArray(row.getField(ordinal));
+ break;
+ case DOUBLE:
+ writeDoubleArray(row.getField(ordinal));
+ break;
+ default:
+ throw new UnsupportedOperationException(
+ "Unsupported element type in array: " +
elementTypeParquet);
+ }
+ } else {
+ throw new UnsupportedOperationException("Unsupported element
type in array: " + elementTypeParquet);
+ }
+ }
+
+ private void writeObjectArray(Object input) {
+ recordConsumer.startGroup();
+ if (input != null) {
+ Object[] inputArray = (Object[]) input;
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (Object ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ }
+ recordConsumer.endGroup();
+ }
+
+ private void writeBooleanArray(Object input) {
+ if (input instanceof boolean[]) {
+ boolean[] inputArray = (boolean[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (boolean ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeTinyIntArray(Object input) {
+ if (input instanceof byte[]) {
+ byte[] inputArray = (byte[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (byte ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeShortArray(Object input) {
+ if (input instanceof short[]) {
+ short[] inputArray = (short[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (short ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeIntArray(Object input) {
+ if (input instanceof int[]) {
+ int[] inputArray = (int[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (int ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeLongArray(Object input) {
+ if (input instanceof long[]) {
+ long[] inputArray = (long[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (long ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeFloatArray(Object input) {
+ if (input instanceof float[]) {
+ float[] inputArray = (float[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (float ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ private void writeDoubleArray(Object input) {
+ if (input instanceof double[]) {
+ double[] inputArray = (double[]) input;
+ recordConsumer.startGroup();
+ if (inputArray.length > 0) {
+ recordConsumer.startField(ARRAY_FIELD_NAME, 0);
+ for (double ele : inputArray) {
+ startGroupAndField(elementTypeParquet.getName());
+ FieldWriter writer = createWriter(elementTypeFlink,
elementTypeParquet);
+ writer.write(Row.of(ele), 0);
+ endGroupAndField(elementTypeParquet.getName());
+ }
+ recordConsumer.endField(ARRAY_FIELD_NAME, 0);
+ }
+ recordConsumer.endGroup();
+ } else {
+ writeObjectArray(input);
+ }
+ }
+
+ }
+
+ private class MapWriter implements FieldWriter {
+
+ private final LogicalType keyTypeFlink;
+
+ private final LogicalType valueTypeFlink;
+
+ private final Type keyTypeParquet;
+
+ private final Type valueTypeParquet;
+
+ public MapWriter(MapType mapTypeFlink, Type mapTypeParquet) {
+ this.keyTypeFlink = mapTypeFlink.getKeyType();
+ this.valueTypeFlink = mapTypeFlink.getValueType();
+ GroupType groupType =
mapTypeParquet.asGroupType().getType(0).asGroupType();
+ this.keyTypeParquet = groupType.getType(0);
+ this.valueTypeParquet = groupType.getType(1);
+ }
+
+ @Override
+ public void write(Row row, int ordinal) {
+ recordConsumer.startGroup();
+ Object inputField = row.getField(ordinal);
+ if (inputField != null) {
+ Map<?, ?> inputMap = (Map<?, ?>) inputField;
+ if (inputMap.size() > 0) {
+ FieldWriter keyWriter = createWriter(keyTypeFlink,
keyTypeParquet);
+ FieldWriter valueWriter = createWriter(valueTypeFlink,
valueTypeParquet);
+ recordConsumer.startField(MAP_ENTITY_FIELD_NAME, 0);
+ for (Map.Entry<?, ?> entry : inputMap.entrySet()) {
+ recordConsumer.startGroup();
+
+ recordConsumer.startField(MAP_KEY_FIELD_NAME, 0);
+ keyWriter.write(Row.of(entry.getKey()), 0);
+ recordConsumer.endField(MAP_KEY_FIELD_NAME,0);
+
+ Object value = entry.getValue();
+ if (value != null) {
+ recordConsumer.startField(MAP_VALUE_FIELD_NAME, 1);
+ valueWriter.write(Row.of(value), 0);
+ recordConsumer.endField(MAP_VALUE_FIELD_NAME, 1);
+ }
+
+ recordConsumer.endGroup();
+ }
+ recordConsumer.endField(MAP_ENTITY_FIELD_NAME, 0);
+ }
+ }
+ recordConsumer.endGroup();
+ }
+ }
+
+ private void startGroupAndField(String fieldName) {
+ recordConsumer.startGroup();
+ recordConsumer.startField(fieldName, 0);
+ }
+
+ private void endGroupAndField(String fieldName) {
+ recordConsumer.endField(fieldName, 0);
+ recordConsumer.endGroup();
+ }
}
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriterBuilder.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriterBuilder.java
index 010b0e9..28b5f32 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriterBuilder.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriterBuilder.java
@@ -95,6 +95,8 @@ public class ParquetRowWriterBuilder extends
ParquetWriter.Builder<Row, ParquetR
/** Flink Row {@link ParquetBuilder}. */
public static class FlinkParquetBuilder implements ParquetBuilder<Row> {
+ private static final long serialVersionUID = 5891262136717753537L;
+
private final RowType rowType;
private final ParquetFileFormat parquetFileFormat;
diff --git
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetSchemaConverter.java
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetSchemaConverter.java
index bd3c4de..4b7eaf5 100644
---
a/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetSchemaConverter.java
+++
b/inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetSchemaConverter.java
@@ -25,7 +25,9 @@ import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.MapTypeInfo;
import org.apache.flink.api.java.typeutils.ObjectArrayTypeInfo;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
+import org.apache.flink.table.types.logical.ArrayType;
import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.MapType;
import org.apache.flink.table.types.logical.RowType;
import org.apache.parquet.schema.GroupType;
@@ -44,6 +46,7 @@ import java.util.List;
/** Schema converter converts Parquet schema to and from Flink internal types.
*/
public class ParquetSchemaConverter {
private static final Logger LOGGER =
LoggerFactory.getLogger(ParquetSchemaConverter.class);
+ public static final String MAP_KEY = "key";
public static final String MAP_VALUE = "value";
public static final String LIST_ARRAY_TYPE = "array";
public static final String LIST_ELEMENT = "element";
@@ -601,6 +604,15 @@ public class ParquetSchemaConverter {
return Types.primitive(PrimitiveType.PrimitiveTypeName.INT64,
repetition)
.as(OriginalType.TIMESTAMP_MILLIS)
.named(name);
+ case ARRAY:
+ return Types.list(repetition)
+ .setElementType(convertToParquetType(LIST_ELEMENT,
((ArrayType) type).getElementType()))
+ .named(name);
+ case MAP:
+ return Types.map(repetition)
+ .key(convertToParquetType(MAP_KEY, ((MapType)
type).getKeyType()))
+ .value(convertToParquetType(MAP_VALUE, ((MapType)
type).getValueType()))
+ .named(name);
default:
throw new UnsupportedOperationException("Unsupported type: " +
type);
}
diff --git
a/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverterTest.java
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverterTest.java
new file mode 100644
index 0000000..2a776a3
--- /dev/null
+++
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/clickhouse/ClickHouseRowConverterTest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.inlong.sort.flink.clickhouse;
+
+import org.apache.inlong.sort.formats.common.IntTypeInfo;
+import org.apache.inlong.sort.formats.common.StringTypeInfo;
+import org.junit.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotEquals;
+import static org.junit.Assert.assertTrue;
+
+public class ClickHouseRowConverterTest {
+
+ @Test
+ public void testToObjectArray() {
+ int[] testArray = new int[] {1, 2, 3};
+ Object obj = ClickHouseRowConverter.toObjectArray(new IntTypeInfo(),
testArray);
+ assertTrue(obj instanceof Integer[]);
+ Integer[] integers = (Integer[]) obj;
+ assertEquals(Integer.valueOf(1), integers[0]);
+ assertEquals(Integer.valueOf(2), integers[1]);
+ assertEquals(Integer.valueOf(3), integers[2]);
+
+ String[] strs = new String[] {"f1", "f2", "f3"};
+ Object strsArray = ClickHouseRowConverter.toObjectArray(new
StringTypeInfo(), strs);
+ assertEquals(strs, strsArray);
+ }
+
+ @Test
+ public void testToKeyValuePairObjectArray() {
+ Map<String, Double> testMap1 = new HashMap<>();
+ testMap1.put("f1", 1.0);
+ testMap1.put("f2", 2.0);
+ Object[] objects1 =
ClickHouseRowConverter.toKeyValuePairObjectArray(testMap1);
+ assertEquals(4, objects1.length);
+ assertTrue(objects1[0].equals("f1") || objects1[0].equals("f2"));
+ assertTrue(objects1[2].equals("f1") || objects1[2].equals("f2"));
+ assertNotEquals(objects1[0], objects1[2]);
+ assertTrue(objects1[1].equals(1.0) || objects1[1].equals(2.0));
+ assertTrue(objects1[3].equals(1.0) || objects1[3].equals(2.0));
+ assertNotEquals(objects1[1], objects1[3]);
+
+ Map<String, Integer> testMap2 = new HashMap<>();
+ Object[] objects2 =
ClickHouseRowConverter.toKeyValuePairObjectArray(testMap2);
+ assertEquals(0, objects2.length);
+ }
+}
diff --git
a/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriterTest.java
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriterTest.java
index 03265c2..71d09d7 100644
---
a/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriterTest.java
+++
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/TextRowWriterTest.java
@@ -29,8 +29,18 @@ import java.io.InputStreamReader;
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 org.apache.flink.core.fs.local.LocalDataOutputStream;
+import org.apache.flink.table.types.logical.ArrayType;
+import org.apache.flink.table.types.logical.BigIntType;
+import org.apache.flink.table.types.logical.CharType;
+import org.apache.flink.table.types.logical.DoubleType;
+import org.apache.flink.table.types.logical.IntType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.MapType;
import org.apache.flink.types.Row;
import org.apache.hadoop.io.compress.CompressionInputStream;
import org.apache.inlong.sort.configuration.Configuration;
@@ -50,7 +60,8 @@ public class TextRowWriterTest {
TextRowWriter textRowWriter = new TextRowWriter(
new LocalDataOutputStream(file),
new TextFileFormat(','),
- new Configuration()
+ new Configuration(),
+ new LogicalType[] {new CharType(), new IntType()}
);
textRowWriter.addElement(Row.of("zhangsan", 1));
@@ -78,7 +89,8 @@ public class TextRowWriterTest {
TextRowWriter textRowWriter = new TextRowWriter(
new LocalDataOutputStream(gzipFile),
new TextFileFormat(',', CompressionType.GZIP),
- new Configuration()
+ new Configuration(),
+ new LogicalType[] {new CharType(), new IntType()}
);
textRowWriter.addElement(Row.of("zhangsan", 1));
@@ -96,7 +108,8 @@ public class TextRowWriterTest {
TextRowWriter textRowWriter = new TextRowWriter(
new LocalDataOutputStream(lzoFile),
new TextFileFormat(',', CompressionType.LZO),
- new Configuration()
+ new Configuration(),
+ new LogicalType[] {new CharType(), new IntType()}
);
textRowWriter.addElement(Row.of("zhangsan", 1));
@@ -142,4 +155,32 @@ public class TextRowWriterTest {
return false;
}
}
+
+ @Test
+ public void testConvertArrayField() {
+ int[] ints = new int[] {1, 2, 3};
+ assertEquals("[1,2,3]", TextRowWriter.convertField(ints, new
ArrayType(new IntType())));
+
+ String[] strings = new String[] {"f1", "f2", null};
+ assertEquals("[f1,f2,null]", TextRowWriter.convertField(strings, new
ArrayType(new CharType())));
+
+ Double[] doubles = new Double[] {1.0, null, 2.0};
+ assertEquals("[1.0,null,2.0]", TextRowWriter.convertField(doubles, new
ArrayType(new DoubleType())));
+
+ long[] longs = new long[] {};
+ assertEquals("[]", TextRowWriter.convertField(longs, new ArrayType(new
BigIntType())));
+ }
+
+ @Test
+ public void testConvertMapField() {
+ Map<String, Double> map = new HashMap<>();
+ map.put("f1", 1.0);
+ map.put("f2", null);
+ map.put("f3", 3.0);
+ assertEquals("{f1=1.0,f2=null,f3=3.0}",
+ TextRowWriter.convertField(map, new MapType(new CharType(),
new DoubleType())));
+
+ Map<Double, Integer> emptyMap = new HashMap<>();
+ assertEquals("{}", TextRowWriter.convertField(emptyMap, new
MapType(new DoubleType(), new IntType())));
+ }
}
diff --git
a/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetBulkWriterTest.java
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetBulkWriterTest.java
index 7b0a681..4834ff8 100644
---
a/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetBulkWriterTest.java
+++
b/inlong-sort/sort-connectors/src/test/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetBulkWriterTest.java
@@ -17,6 +17,10 @@
package org.apache.inlong.sort.flink.hive.formats.parquet;
+import static
org.apache.inlong.sort.flink.hive.formats.parquet.ParquetRowWriter.ARRAY_FIELD_NAME;
+import static
org.apache.inlong.sort.flink.hive.formats.parquet.ParquetRowWriter.MAP_ENTITY_FIELD_NAME;
+import static
org.apache.inlong.sort.flink.hive.formats.parquet.ParquetRowWriter.MAP_KEY_FIELD_NAME;
+import static
org.apache.inlong.sort.flink.hive.formats.parquet.ParquetRowWriter.MAP_VALUE_FIELD_NAME;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
@@ -27,12 +31,16 @@ import java.math.BigDecimal;
import java.sql.Date;
import java.sql.Time;
import java.sql.Timestamp;
+import java.util.HashMap;
+import java.util.Map;
+
import org.apache.flink.core.fs.local.LocalDataOutputStream;
import org.apache.flink.formats.parquet.ParquetBulkWriter;
import org.apache.flink.types.Row;
import org.apache.hadoop.fs.Path;
import org.apache.inlong.sort.configuration.Configuration;
import org.apache.inlong.sort.flink.hive.HiveSinkHelper;
+import org.apache.inlong.sort.formats.common.ArrayFormatInfo;
import org.apache.inlong.sort.formats.common.BooleanFormatInfo;
import org.apache.inlong.sort.formats.common.ByteFormatInfo;
import org.apache.inlong.sort.formats.common.DateFormatInfo;
@@ -41,6 +49,7 @@ import org.apache.inlong.sort.formats.common.DoubleFormatInfo;
import org.apache.inlong.sort.formats.common.FloatFormatInfo;
import org.apache.inlong.sort.formats.common.IntFormatInfo;
import org.apache.inlong.sort.formats.common.LongFormatInfo;
+import org.apache.inlong.sort.formats.common.MapFormatInfo;
import org.apache.inlong.sort.formats.common.ShortFormatInfo;
import org.apache.inlong.sort.formats.common.StringFormatInfo;
import org.apache.inlong.sort.formats.common.TimeFormatInfo;
@@ -65,9 +74,17 @@ public class ParquetBulkWriterTest {
@Test
public void testAllSupportedTypes() throws IOException {
File testFile = temporaryFolder.newFile("test.parquet");
- ParquetBulkWriter<Row> parquetBulkWriter = (ParquetBulkWriter<Row>)
HiveSinkHelper
+ final ParquetBulkWriter<Row> parquetBulkWriter =
(ParquetBulkWriter<Row>) HiveSinkHelper
.createBulkWriterFactory(prepareHiveSinkInfo(), new
Configuration())
.create(new LocalDataOutputStream(testFile));
+
+ final int[] testArray = new int[] {1, 2, 3};
+ Map<String, Double> testMap = new HashMap<>();
+ testMap.put("key1", 1.0);
+ testMap.put("key2", 2.0);
+ testMap.put("key3", 3.0);
+ final Integer[] testEmptyArray = new Integer[] {};
+ final Map<String, Double> testEmptyMap = new HashMap<>();
parquetBulkWriter.addElement(Row.of(
"string",
false,
@@ -80,7 +97,11 @@ public class ParquetBulkWriterTest {
new BigDecimal("123456789123456789"),
new Date(0),
new Time(0),
- new Timestamp(0)
+ new Timestamp(0),
+ testArray,
+ testMap,
+ testEmptyArray,
+ testEmptyMap
));
parquetBulkWriter.finish();
@@ -102,6 +123,22 @@ public class ParquetBulkWriterTest {
assertEquals(0, line.getInteger(9, 0));
assertEquals(0, line.getInteger(10, 0));
assertEquals(0, line.getLong(11, 0));
+
+ Group f13 = line.getGroup("f13", 0);
+ assertEquals(1, f13.getGroup(ARRAY_FIELD_NAME, 0).getInteger(0, 0));
+ assertEquals(2, f13.getGroup(ARRAY_FIELD_NAME, 1).getInteger(0, 0));
+ assertEquals(3, f13.getGroup(ARRAY_FIELD_NAME, 2).getInteger(0, 0));
+
+ Group f14 = line.getGroup("f14", 0);
+ assertEquals("key1", f14.getGroup(MAP_ENTITY_FIELD_NAME,
0).getString(MAP_KEY_FIELD_NAME, 0));
+ assertEquals(1.0, f14.getGroup(MAP_ENTITY_FIELD_NAME,
0).getDouble(MAP_VALUE_FIELD_NAME, 0), 0.01);
+ assertEquals("key2", f14.getGroup(MAP_ENTITY_FIELD_NAME,
1).getString(MAP_KEY_FIELD_NAME, 0));
+ assertEquals(2.0, f14.getGroup(MAP_ENTITY_FIELD_NAME,
1).getDouble(MAP_VALUE_FIELD_NAME, 0), 0.01);
+ assertEquals("key3", f14.getGroup(MAP_ENTITY_FIELD_NAME,
2).getString(MAP_KEY_FIELD_NAME, 0));
+ assertEquals(3.0, f14.getGroup(MAP_ENTITY_FIELD_NAME,
2).getDouble(MAP_VALUE_FIELD_NAME, 0), 0.01);
+
+ assertEquals("", line.getGroup("f15", 0).toString());
+ assertEquals("", line.getGroup("f16", 0).toString());
}
private HiveSinkInfo prepareHiveSinkInfo() {
@@ -118,7 +155,11 @@ public class ParquetBulkWriterTest {
new FieldInfo("f9", DecimalFormatInfo.INSTANCE),
new FieldInfo("f10", new DateFormatInfo()),
new FieldInfo("f11", new TimeFormatInfo()),
- new FieldInfo("f12", new TimestampFormatInfo())
+ new FieldInfo("f12", new TimestampFormatInfo()),
+ new FieldInfo("f13", new
ArrayFormatInfo(IntFormatInfo.INSTANCE)),
+ new FieldInfo("f14", new
MapFormatInfo(StringFormatInfo.INSTANCE, DoubleFormatInfo.INSTANCE)),
+ new FieldInfo("f15", new
ArrayFormatInfo(IntFormatInfo.INSTANCE)),
+ new FieldInfo("f16", new
MapFormatInfo(StringFormatInfo.INSTANCE, DoubleFormatInfo.INSTANCE))
},
"jdbc:mysql://127.0.0.1:3306/testDatabaseName",
"testDatabaseName",
@@ -126,7 +167,7 @@ public class ParquetBulkWriterTest {
"testUsername",
"testPassword",
"/path",
- new HivePartitionInfo[]{
+ new HivePartitionInfo[] {
new HiveFieldPartitionInfo("f13"),
},
new ParquetFileFormat()