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/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 2695d2b64 [INLONG-7331][Sort][Manager] Support complex type field 
(#7343)
2695d2b64 is described below

commit 2695d2b64b1c73a1e241a349f301bb7df54b799e
Author: hannah <[email protected]>
AuthorDate: Mon Feb 13 10:37:46 2023 +0800

    [INLONG-7331][Sort][Manager] Support complex type field (#7343)
    
    Co-authored-by: healchow <[email protected]>
---
 .../manager/pojo/sort/util/FieldInfoUtils.java     | 17 +++++-
 .../org/apache/inlong/sort/protocol/FieldInfo.java | 12 ++--
 .../inlong/sort/parser/impl/FlinkSqlParser.java    | 13 +++-
 .../sort/parser/MongoExtractFlinkSqlParseTest.java | 71 ++++++++++++++++++++++
 4 files changed, 104 insertions(+), 9 deletions(-)

diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/util/FieldInfoUtils.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/util/FieldInfoUtils.java
index 8dbe3362b..b1d297030 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/util/FieldInfoUtils.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/util/FieldInfoUtils.java
@@ -148,7 +148,22 @@ public class FieldInfoUtils {
      * @return Sort field format instance
      */
     public static FormatInfo convertFieldFormat(String type) {
-        return convertFieldFormat(type, null);
+        FieldType fieldType = FieldType.forName(type);
+        FormatInfo formatInfo;
+        switch (fieldType) {
+            case ARRAY:
+                formatInfo = new ArrayFormatInfo(new StringFormatInfo());
+                break;
+            case MAP:
+                formatInfo = new MapFormatInfo(new StringFormatInfo(), new 
StringFormatInfo());
+                break;
+            case STRUCT:
+                formatInfo = new RowFormatInfo(new String[]{}, new 
FormatInfo[]{});
+                break;
+            default:
+                formatInfo = convertFieldFormat(type, null);
+        }
+        return formatInfo;
     }
 
     /**
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/FieldInfo.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/FieldInfo.java
index 623229fda..0ab1a495d 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/FieldInfo.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/FieldInfo.java
@@ -79,11 +79,13 @@ public class FieldInfo implements FunctionParam, 
Serializable {
     @Override
     public String format() {
         String formatName = name.trim();
-        if (!formatName.startsWith("`")) {
-            formatName = String.format("`%s", formatName);
-        }
-        if (!formatName.endsWith("`")) {
-            formatName = String.format("%s`", formatName);
+        if (!formatName.contains(".")) {
+            if (!formatName.startsWith("`")) {
+                formatName = String.format("`%s", formatName);
+            }
+            if (!formatName.endsWith("`")) {
+                formatName = String.format("%s`", formatName);
+            }
         }
         if (StringUtils.isNotBlank(tableNameAlias)) {
             return String.format("%s.%s", tableNameAlias, formatName);
diff --git 
a/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/parser/impl/FlinkSqlParser.java
 
b/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/parser/impl/FlinkSqlParser.java
index 22f3b968f..64fe0dd98 100644
--- 
a/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/parser/impl/FlinkSqlParser.java
+++ 
b/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/parser/impl/FlinkSqlParser.java
@@ -22,7 +22,10 @@ import org.apache.commons.lang3.StringUtils;
 import org.apache.flink.table.api.TableEnvironment;
 import org.apache.inlong.sort.configuration.Constants;
 import org.apache.inlong.sort.formats.base.TableFormatUtils;
+import org.apache.inlong.sort.formats.common.ArrayFormatInfo;
 import org.apache.inlong.sort.formats.common.FormatInfo;
+import org.apache.inlong.sort.formats.common.MapFormatInfo;
+import org.apache.inlong.sort.formats.common.RowFormatInfo;
 import org.apache.inlong.sort.function.EncryptFunction;
 import org.apache.inlong.sort.function.JsonGetterFunction;
 import org.apache.inlong.sort.function.RegexpReplaceFirstFunction;
@@ -635,11 +638,15 @@ public class FlinkSqlParser implements Parser {
             Map<String, FieldRelation> fieldRelationMap, StringBuilder sb) {
         for (FieldInfo field : fields) {
             FieldRelation fieldRelation = 
fieldRelationMap.get(field.getName());
+            FormatInfo fieldFormatInfo = field.getFormatInfo();
             if (fieldRelation == null) {
-                String targetType = 
TableFormatUtils.deriveLogicalType(field.getFormatInfo()).asSummaryString();
+                String targetType = 
TableFormatUtils.deriveLogicalType(fieldFormatInfo).asSummaryString();
                 sb.append("\n    CAST(NULL as ").append(targetType).append(") 
AS ").append(field.format()).append(",");
                 continue;
             }
+            boolean complexType = fieldFormatInfo instanceof RowFormatInfo
+                    || fieldFormatInfo instanceof ArrayFormatInfo
+                    || fieldFormatInfo instanceof MapFormatInfo;
             FunctionParam inputField = fieldRelation.getInputField();
             if (inputField instanceof FieldInfo) {
                 FieldInfo fieldInfo = (FieldInfo) inputField;
@@ -649,10 +656,10 @@ public class FlinkSqlParser implements Parser {
                         && outputField != null
                         && outputField.getFormatInfo() != null
                         && 
outputField.getFormatInfo().getTypeInfo().equals(formatInfo.getTypeInfo());
-                if (sameType || field.getFormatInfo() == null) {
+                if (complexType || sameType || fieldFormatInfo == null) {
                     sb.append("\n    ").append(inputField.format()).append(" 
AS ").append(field.format()).append(",");
                 } else {
-                    String targetType = 
TableFormatUtils.deriveLogicalType(field.getFormatInfo()).asSummaryString();
+                    String targetType = 
TableFormatUtils.deriveLogicalType(fieldFormatInfo).asSummaryString();
                     sb.append("\n    
CAST(").append(inputField.format()).append(" as ")
                             .append(targetType).append(") AS 
").append(field.format()).append(",");
                 }
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
index 2b59e7fb7..240b11a7b 100644
--- 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/parser/MongoExtractFlinkSqlParseTest.java
@@ -22,6 +22,8 @@ import org.apache.flink.table.api.EnvironmentSettings;
 import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
 import org.apache.flink.test.util.AbstractTestBase;
 import org.apache.inlong.common.enums.MetaField;
+import org.apache.inlong.sort.formats.common.FormatInfo;
+import org.apache.inlong.sort.formats.common.RowFormatInfo;
 import org.apache.inlong.sort.formats.common.StringFormatInfo;
 import org.apache.inlong.sort.formats.common.TimestampFormatInfo;
 import org.apache.inlong.sort.parser.impl.FlinkSqlParser;
@@ -123,4 +125,73 @@ public class MongoExtractFlinkSqlParseTest extends 
AbstractTestBase {
         Assert.assertTrue(result.tryExecute());
     }
 
+    private MongoExtractNode buildComplexTypeMongoNode() {
+        String[] fieldNames = new String[]{"name", "age"};
+        FormatInfo[] fieldFormatInfos = new FormatInfo[]{new 
StringFormatInfo(), new StringFormatInfo()};
+        List<FieldInfo> fields = Arrays.asList(
+                new FieldInfo("info", new RowFormatInfo(fieldNames, 
fieldFormatInfos)),
+                new MetaFieldInfo("proctime", MetaField.PROCESS_TIME),
+                new MetaFieldInfo("database_name", MetaField.DATABASE_NAME),
+                new MetaFieldInfo("collection_name", 
MetaField.COLLECTION_NAME),
+                new MetaFieldInfo("op_ts", MetaField.OP_TS));
+        return new MongoExtractNode("1", "mysql_input", fields,
+                null, null, "test", "localhost:27017",
+                "root", "inlong", "test");
+    }
+
+    private KafkaLoadNode buildComplexTypeKafkaLoadNode() {
+        List<FieldInfo> fields = Arrays.asList(
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("_id", new StringFormatInfo()),
+                new FieldInfo("proctime", new TimestampFormatInfo()),
+                new FieldInfo("database_name", new StringFormatInfo()),
+                new FieldInfo("collection_name", new StringFormatInfo()),
+                new FieldInfo("op_ts", new TimestampFormatInfo()));
+        List<FieldRelation> relations = Arrays.asList(
+                new FieldRelation(new FieldInfo("info.name", new 
StringFormatInfo()),
+                        new FieldInfo("name", new StringFormatInfo())),
+                new FieldRelation(new FieldInfo("_id", new StringFormatInfo()),
+                        new FieldInfo("_id", new StringFormatInfo())),
+                new FieldRelation(new FieldInfo("proctime", new 
TimestampFormatInfo()),
+                        new FieldInfo("proctime", new TimestampFormatInfo())),
+                new FieldRelation(new FieldInfo("database_name", new 
StringFormatInfo()),
+                        new FieldInfo("database_name", new 
StringFormatInfo())),
+                new FieldRelation(new FieldInfo("collection_name", new 
StringFormatInfo()),
+                        new FieldInfo("collection_name", new 
StringFormatInfo())),
+                new FieldRelation(new FieldInfo("op_ts", new 
TimestampFormatInfo()),
+                        new FieldInfo("op_ts", new TimestampFormatInfo())));
+        CsvFormat csvFormat = new CsvFormat();
+        csvFormat.setDisableQuoteCharacter(true);
+        return new KafkaLoadNode("2", "kafka_output", fields, relations, null, 
null,
+                "test", "localhost:9092",
+                csvFormat, null,
+                null, "_id");
+    }
+
+    /**
+     * Test mongodb to kafka
+     *
+     * @throws Exception The exception may throws when execute the case
+     */
+    @Test
+    public void testMongoDbComplexTypeToKafka() throws Exception {
+        EnvironmentSettings settings = EnvironmentSettings
+                .newInstance()
+                .useBlinkPlanner()
+                .inStreamingMode()
+                .build();
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+        env.setParallelism(1);
+        env.enableCheckpointing(10000);
+        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, 
settings);
+        Node inputNode = buildComplexTypeMongoNode();
+        Node outputNode = buildComplexTypeKafkaLoadNode();
+        StreamInfo streamInfo = new StreamInfo("1", Arrays.asList(inputNode, 
outputNode),
+                
Collections.singletonList(buildNodeRelation(Collections.singletonList(inputNode),
+                        Collections.singletonList(outputNode))));
+        GroupInfo groupInfo = new GroupInfo("1", 
Collections.singletonList(streamInfo));
+        FlinkSqlParser parser = FlinkSqlParser.getInstance(tableEnv, 
groupInfo);
+        ParseResult result = parser.parse();
+        Assert.assertTrue(result.tryExecute());
+    }
 }

Reply via email to