KevinWen007 commented on a change in pull request #2615:
URL: https://github.com/apache/incubator-inlong/pull/2615#discussion_r810884719
##########
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);
+ }
+ }
Review comment:
Same as above
--
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]