KevinWen007 commented on a change in pull request #2615:
URL: https://github.com/apache/incubator-inlong/pull/2615#discussion_r810914500



##########
File path: 
inlong-sort/sort-connectors/src/main/java/org/apache/inlong/sort/flink/hive/formats/parquet/ParquetRowWriter.java
##########
@@ -184,4 +202,281 @@ public void write(Row row, int ordinal) {
             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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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(), 0);
+                        FieldWriter writer = createWriter(elementTypeFlink, 
elementTypeParquet);
+                        writer.write(Row.of(ele), 0);
+                        endGroupAndField(elementTypeParquet.getName(), 0);
+                    }
+                    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) {
+                return;
+            }
+
+            Map<?, ?> inputMap = (Map<?, ?>) row.getField(ordinal);

Review comment:
       done




-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


Reply via email to