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()

Reply via email to