This is an automated email from the ASF dual-hosted git repository.

healchow 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 84644f21d [INLONG-4678][Sort] Add util class of parsing meta info for 
data node (#4680)
84644f21d is described below

commit 84644f21d11dee8ff748e1d7ecd7c789f8fcbf2f
Author: pacino <[email protected]>
AuthorDate: Thu Jun 16 23:57:20 2022 +0800

    [INLONG-4678][Sort] Add util class of parsing meta info for data node 
(#4680)
---
 .../org/apache/inlong/common/enums/MetaField.java  |   9 +-
 .../inlong/sort/parser/impl/FlinkSqlParser.java    | 188 +------------
 .../apache/inlong/sort/util/MetaInfoParseUtil.java | 294 +++++++++++++++++++++
 .../inlong/sort/util/MetaInfoParseUtilTest.java    | 189 +++++++++++++
 4 files changed, 493 insertions(+), 187 deletions(-)

diff --git 
a/inlong-common/src/main/java/org/apache/inlong/common/enums/MetaField.java 
b/inlong-common/src/main/java/org/apache/inlong/common/enums/MetaField.java
index 017e01219..fbfd2773b 100644
--- a/inlong-common/src/main/java/org/apache/inlong/common/enums/MetaField.java
+++ b/inlong-common/src/main/java/org/apache/inlong/common/enums/MetaField.java
@@ -28,7 +28,7 @@ public enum MetaField {
      */
     PROCESS_TIME,
     /**
-     * Name of the schema that contain the row, currently used for Oracle 
database
+     * Name of the schema that contain the row, currently used for Oracle, 
PostgreSQL, SQLSERVER
      */
     SCHEMA_NAME,
     /**
@@ -79,7 +79,12 @@ public enum MetaField {
     /**
      * Primary key field name. Currently, it is used for MySQL database.
      */
-    PK_NAMES;
+    PK_NAMES,
+
+    /**
+     * Name of the collection that contain the row. For MongoDB
+     */
+    COLLECTION_NAME;
 
     public static MetaField forName(String name) {
         for (MetaField metaField : values()) {
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 4bab1d790..fa6be00e4 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
@@ -20,7 +20,6 @@ package org.apache.inlong.sort.parser.impl;
 import com.google.common.base.Preconditions;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.flink.table.api.TableEnvironment;
-import org.apache.inlong.common.enums.MetaField;
 import org.apache.inlong.sort.formats.base.TableFormatUtils;
 import org.apache.inlong.sort.function.RegexpReplaceFirstFunction;
 import org.apache.inlong.sort.parser.Parser;
@@ -34,11 +33,7 @@ import org.apache.inlong.sort.protocol.enums.FilterStrategy;
 import org.apache.inlong.sort.protocol.node.ExtractNode;
 import org.apache.inlong.sort.protocol.node.LoadNode;
 import org.apache.inlong.sort.protocol.node.Node;
-import org.apache.inlong.sort.protocol.node.extract.KafkaExtractNode;
-import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
-import org.apache.inlong.sort.protocol.node.extract.OracleExtractNode;
 import org.apache.inlong.sort.protocol.node.load.HbaseLoadNode;
-import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
 import org.apache.inlong.sort.protocol.node.transform.DistinctNode;
 import org.apache.inlong.sort.protocol.node.transform.TransformNode;
 import org.apache.inlong.sort.protocol.transformation.FieldRelation;
@@ -48,6 +43,7 @@ import 
org.apache.inlong.sort.protocol.transformation.FunctionParam;
 import org.apache.inlong.sort.protocol.transformation.relation.JoinRelation;
 import org.apache.inlong.sort.protocol.transformation.relation.NodeRelation;
 import 
org.apache.inlong.sort.protocol.transformation.relation.UnionNodeRelation;
+import org.apache.inlong.sort.util.MetaInfoParseUtil;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -62,7 +58,7 @@ import java.util.stream.Collectors;
 
 /**
  * Flink sql parse handler
- * It accepts a Tableenv and GroupInfo, and outputs the parsed 
FlinkSqlParseResult
+ * It accepts a TableEnvironment and GroupInfo, and outputs the parsed 
FlinkSqlParseResult
  */
 public class FlinkSqlParser implements Parser {
 
@@ -668,7 +664,7 @@ public class FlinkSqlParser implements Parser {
             sb.append("    `").append(field.getName()).append("` ");
             if (field instanceof MetaFieldInfo) {
                 MetaFieldInfo metaFieldInfo = (MetaFieldInfo) field;
-                parseMetaField(node, metaFieldInfo, sb);
+                MetaInfoParseUtil.parseMetaField(node, metaFieldInfo, sb);
             } else {
                 
sb.append(TableFormatUtils.deriveLogicalType(field.getFormatInfo()).asSummaryString());
             }
@@ -680,184 +676,6 @@ public class FlinkSqlParser implements Parser {
         return sb.toString();
     }
 
-    private void parseMetaField(Node node, MetaFieldInfo metaFieldInfo, 
StringBuilder sb) {
-        if (metaFieldInfo.getMetaField() == MetaField.PROCESS_TIME) {
-            sb.append(" AS PROCTIME()");
-            return;
-        }
-        if (node instanceof MySqlExtractNode) {
-            sb.append(parseMySqlExtractNodeMetaField(metaFieldInfo));
-        } else if (node instanceof OracleExtractNode) {
-            sb.append(parseOracleExtractNodeMetaField(metaFieldInfo));
-        } else if (node instanceof KafkaExtractNode) {
-            sb.append(parseKafkaExtractNodeMetaField(metaFieldInfo));
-        } else if (node instanceof KafkaLoadNode) {
-            sb.append(parseKafkaLoadNodeMetaField(metaFieldInfo));
-        } else {
-            throw new UnsupportedOperationException(
-                    String.format("This node:%s does not currently support 
metadata fields",
-                            node.getClass().getName()));
-        }
-    }
-
-    private String parseKafkaLoadNodeMetaField(MetaFieldInfo metaFieldInfo) {
-        String metaType;
-        switch (metaFieldInfo.getMetaField()) {
-            case TABLE_NAME:
-                metaType = "STRING METADATA FROM 'value.table'";
-                break;
-            case DATABASE_NAME:
-                metaType = "STRING METADATA FROM 'value.database'";
-                break;
-            case OP_TS:
-                metaType = "TIMESTAMP(3) METADATA FROM 
'value.event-timestamp'";
-                break;
-            case OP_TYPE:
-                metaType = "STRING METADATA FROM 'value.op-type'";
-                break;
-            case DATA:
-                metaType = "STRING METADATA FROM 'value.data'";
-                break;
-            case IS_DDL:
-                metaType = "BOOLEAN METADATA FROM 'value.is-ddl'";
-                break;
-            case TS:
-                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.ingestion-timestamp'";
-                break;
-            case SQL_TYPE:
-                metaType = "MAP<STRING, INT> METADATA FROM 'value.sql-type'";
-                break;
-            case MYSQL_TYPE:
-                metaType = "MAP<STRING, STRING> METADATA FROM 
'value.mysql-type'";
-                break;
-            case PK_NAMES:
-                metaType = "ARRAY<STRING> METADATA FROM 'value.pk-names'";
-                break;
-            case BATCH_ID:
-                metaType = "BIGINT METADATA FROM 'value.batch-id'";
-                break;
-            case UPDATE_BEFORE:
-                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'value.update-before'";
-                break;
-            default:
-                throw new 
UnsupportedOperationException(String.format("Unsupport meta field: %s",
-                        metaFieldInfo.getMetaField()));
-        }
-        return metaType;
-    }
-
-    private String parseKafkaExtractNodeMetaField(MetaFieldInfo metaFieldInfo) 
{
-        String metaType;
-        switch (metaFieldInfo.getMetaField()) {
-            case TABLE_NAME:
-                metaType = "STRING METADATA FROM 'value.table'";
-                break;
-            case DATABASE_NAME:
-                metaType = "STRING METADATA FROM 'value.database'";
-                break;
-            case SQL_TYPE:
-                metaType = "MAP<STRING, INT> METADATA FROM 'value.sql-type'";
-                break;
-            case PK_NAMES:
-                metaType = "ARRAY<STRING> METADATA FROM 'value.pk-names'";
-                break;
-            case TS:
-                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.ingestion-timestamp'";
-                break;
-            case OP_TS:
-                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.event-timestamp'";
-                break;
-            // additional metadata
-            case OP_TYPE:
-                metaType = "STRING METADATA FROM 'value.op-type'";
-                break;
-            case IS_DDL:
-                metaType = "BOOLEAN METADATA FROM 'value.is-ddl'";
-                break;
-            case MYSQL_TYPE:
-                metaType = "MAP<STRING, STRING> METADATA FROM 
'value.mysql-type'";
-                break;
-            case BATCH_ID:
-                metaType = "BIGINT METADATA FROM 'value.batch-id'";
-                break;
-            case UPDATE_BEFORE:
-                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'value.update-before'";
-                break;
-            default:
-                throw new 
UnsupportedOperationException(String.format("Unsupport meta field: %s",
-                        metaFieldInfo.getMetaField()));
-        }
-        return metaType;
-    }
-
-    private String parseMySqlExtractNodeMetaField(MetaFieldInfo metaFieldInfo) 
{
-        String metaType;
-        switch (metaFieldInfo.getMetaField()) {
-            case TABLE_NAME:
-                metaType = "STRING METADATA FROM 'meta.table_name' VIRTUAL";
-                break;
-            case DATABASE_NAME:
-                metaType = "STRING METADATA FROM 'meta.database_name' VIRTUAL";
-                break;
-            case OP_TS:
-                metaType = "TIMESTAMP(3) METADATA FROM 'meta.op_ts' VIRTUAL";
-                break;
-            case OP_TYPE:
-                metaType = "STRING METADATA FROM 'meta.op_type' VIRTUAL";
-                break;
-            case DATA:
-                metaType = "STRING METADATA FROM 'meta.data' VIRTUAL";
-                break;
-            case IS_DDL:
-                metaType = "BOOLEAN METADATA FROM 'meta.is_ddl' VIRTUAL";
-                break;
-            case TS:
-                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'meta.ts' VIRTUAL";
-                break;
-            case SQL_TYPE:
-                metaType = "MAP<STRING, INT> METADATA FROM 'meta.sql_type' 
VIRTUAL";
-                break;
-            case MYSQL_TYPE:
-                metaType = "MAP<STRING, STRING> METADATA FROM 
'meta.mysql_type' VIRTUAL";
-                break;
-            case PK_NAMES:
-                metaType = "ARRAY<STRING> METADATA FROM 'meta.pk_names' 
VIRTUAL";
-                break;
-            case BATCH_ID:
-                metaType = "BIGINT METADATA FROM 'meta.batch_id' VIRTUAL";
-                break;
-            case UPDATE_BEFORE:
-                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'meta.update_before' VIRTUAL";
-                break;
-            default:
-                throw new 
UnsupportedOperationException(String.format("Unsupport meta field: %s",
-                        metaFieldInfo.getMetaField()));
-        }
-        return metaType;
-    }
-
-    private String parseOracleExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
-        String metaType;
-        switch (metaFieldInfo.getMetaField()) {
-            case TABLE_NAME:
-                metaType = "STRING METADATA FROM 'table_name' VIRTUAL";
-                break;
-            case SCHEMA_NAME:
-                metaType = "STRING METADATA FROM 'schema_name' VIRTUAL";
-                break;
-            case DATABASE_NAME:
-                metaType = "STRING METADATA FROM 'database_name' VIRTUAL";
-                break;
-            case OP_TS:
-                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL";
-                break;
-            default:
-                throw new 
UnsupportedOperationException(String.format("Unsupport meta field: %s",
-                        metaFieldInfo.getMetaField()));
-        }
-        return metaType;
-    }
-
     /**
      * Generate primary key format in sql
      *
diff --git 
a/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/util/MetaInfoParseUtil.java
 
b/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/util/MetaInfoParseUtil.java
new file mode 100644
index 000000000..8c9447bbc
--- /dev/null
+++ 
b/inlong-sort/sort-core/src/main/java/org/apache/inlong/sort/util/MetaInfoParseUtil.java
@@ -0,0 +1,294 @@
+/*
+ *   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.util;
+
+import org.apache.inlong.common.enums.MetaField;
+import org.apache.inlong.sort.protocol.MetaFieldInfo;
+import org.apache.inlong.sort.protocol.node.Node;
+import org.apache.inlong.sort.protocol.node.extract.KafkaExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.MongoExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.OracleExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.PostgresExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.SqlServerExtractNode;
+import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
+
+/**
+ * Tool class of parsing meta field for different database
+ */
+public class MetaInfoParseUtil {
+
+    /**
+     * parse meta field for different database
+     *
+     * @param node data node info
+     * @param metaFieldInfo meta field info
+     * @param sb sql string append
+     */
+    public static void parseMetaField(Node node, MetaFieldInfo metaFieldInfo, 
StringBuilder sb) {
+        if (metaFieldInfo.getMetaField() == MetaField.PROCESS_TIME) {
+            sb.append(" AS PROCTIME()");
+            return;
+        }
+        if (node instanceof MySqlExtractNode) {
+            sb.append(parseMySqlExtractNodeMetaField(metaFieldInfo));
+        } else if (node instanceof OracleExtractNode) {
+            sb.append(parseOracleExtractNodeMetaField(metaFieldInfo));
+        } else if (node instanceof KafkaExtractNode) {
+            sb.append(parseKafkaExtractNodeMetaField(metaFieldInfo));
+        } else if (node instanceof KafkaLoadNode) {
+            sb.append(parseKafkaLoadNodeMetaField(metaFieldInfo));
+        } else if (node instanceof PostgresExtractNode) {
+            sb.append(parsePostgresExtractNodeMetaField(metaFieldInfo));
+        } else if (node instanceof SqlServerExtractNode) {
+            sb.append(parseSQLServerExtractNodeMetaField(metaFieldInfo));
+        } else if (node instanceof MongoExtractNode) {
+            sb.append(parseMongoExtractNodeMetaField(metaFieldInfo));
+        } else {
+            throw new UnsupportedOperationException(
+                    String.format("This node:%s does not currently support 
metadata fields",
+                            node.getClass().getName()));
+        }
+    }
+
+    private static String parseKafkaLoadNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'value.table'";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'value.database'";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP(3) METADATA FROM 
'value.event-timestamp'";
+                break;
+            case OP_TYPE:
+                metaType = "STRING METADATA FROM 'value.op-type'";
+                break;
+            case DATA:
+                metaType = "STRING METADATA FROM 'value.data'";
+                break;
+            case IS_DDL:
+                metaType = "BOOLEAN METADATA FROM 'value.is-ddl'";
+                break;
+            case TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.ingestion-timestamp'";
+                break;
+            case SQL_TYPE:
+                metaType = "MAP<STRING, INT> METADATA FROM 'value.sql-type'";
+                break;
+            case MYSQL_TYPE:
+                metaType = "MAP<STRING, STRING> METADATA FROM 
'value.mysql-type'";
+                break;
+            case PK_NAMES:
+                metaType = "ARRAY<STRING> METADATA FROM 'value.pk-names'";
+                break;
+            case BATCH_ID:
+                metaType = "BIGINT METADATA FROM 'value.batch-id'";
+                break;
+            case UPDATE_BEFORE:
+                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'value.update-before'";
+                break;
+            default:
+                throw new 
UnsupportedOperationException(String.format("Unsupport meta field for Kafka 
Load Node: %s",
+                        metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parseKafkaExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'value.table'";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'value.database'";
+                break;
+            case SQL_TYPE:
+                metaType = "MAP<STRING, INT> METADATA FROM 'value.sql-type'";
+                break;
+            case PK_NAMES:
+                metaType = "ARRAY<STRING> METADATA FROM 'value.pk-names'";
+                break;
+            case TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.ingestion-timestamp'";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 
'value.event-timestamp'";
+                break;
+            // additional metadata
+            case OP_TYPE:
+                metaType = "STRING METADATA FROM 'value.op-type'";
+                break;
+            case IS_DDL:
+                metaType = "BOOLEAN METADATA FROM 'value.is-ddl'";
+                break;
+            case MYSQL_TYPE:
+                metaType = "MAP<STRING, STRING> METADATA FROM 
'value.mysql-type'";
+                break;
+            case BATCH_ID:
+                metaType = "BIGINT METADATA FROM 'value.batch-id'";
+                break;
+            case UPDATE_BEFORE:
+                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'value.update-before'";
+                break;
+            default:
+                throw new 
UnsupportedOperationException(String.format("Unsupport meta field for Kafka 
Extract Node: %s",
+                        metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parseMySqlExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'meta.table_name' VIRTUAL";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'meta.database_name' VIRTUAL";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP(3) METADATA FROM 'meta.op_ts' VIRTUAL";
+                break;
+            case OP_TYPE:
+                metaType = "STRING METADATA FROM 'meta.op_type' VIRTUAL";
+                break;
+            case DATA:
+                metaType = "STRING METADATA FROM 'meta.data' VIRTUAL";
+                break;
+            case IS_DDL:
+                metaType = "BOOLEAN METADATA FROM 'meta.is_ddl' VIRTUAL";
+                break;
+            case TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'meta.ts' VIRTUAL";
+                break;
+            case SQL_TYPE:
+                metaType = "MAP<STRING, INT> METADATA FROM 'meta.sql_type' 
VIRTUAL";
+                break;
+            case MYSQL_TYPE:
+                metaType = "MAP<STRING, STRING> METADATA FROM 
'meta.mysql_type' VIRTUAL";
+                break;
+            case PK_NAMES:
+                metaType = "ARRAY<STRING> METADATA FROM 'meta.pk_names' 
VIRTUAL";
+                break;
+            case BATCH_ID:
+                metaType = "BIGINT METADATA FROM 'meta.batch_id' VIRTUAL";
+                break;
+            case UPDATE_BEFORE:
+                metaType = "ARRAY<MAP<STRING, STRING>> METADATA FROM 
'meta.update_before' VIRTUAL";
+                break;
+            default:
+                throw new 
UnsupportedOperationException(String.format("Unsupport meta field for MySQL 
Extract Node: %s",
+                        metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parseOracleExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'table_name' VIRTUAL";
+                break;
+            case SCHEMA_NAME:
+                metaType = "STRING METADATA FROM 'schema_name' VIRTUAL";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'database_name' VIRTUAL";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL";
+                break;
+            default:
+                throw new 
UnsupportedOperationException(String.format("Unsupport meta field for Oracle 
Extract Node: "
+                                + "%s",
+                        metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parsePostgresExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'table_name' VIRTUAL";
+                break;
+            case SCHEMA_NAME:
+                metaType = "STRING METADATA FROM 'schema_name' VIRTUAL";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'database_name' VIRTUAL";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL";
+                break;
+            default:
+                throw new UnsupportedOperationException(
+                        String.format("Unsupport meta field for Postgres 
ExtractNode : %s",
+                                metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parseSQLServerExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case TABLE_NAME:
+                metaType = "STRING METADATA FROM 'table_name' VIRTUAL";
+                break;
+            case SCHEMA_NAME:
+                metaType = "STRING METADATA FROM 'schema_name' VIRTUAL";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'database_name' VIRTUAL";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL";
+                break;
+            default:
+                throw new UnsupportedOperationException(
+                        String.format("Unsupport meta field for SQLServer 
ExtractNode: %s",
+                                metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+
+    private static String parseMongoExtractNodeMetaField(MetaFieldInfo 
metaFieldInfo) {
+        String metaType;
+        switch (metaFieldInfo.getMetaField()) {
+            case COLLECTION_NAME:
+                metaType = "STRING METADATA FROM 'collection_name' VIRTUAL";
+                break;
+            case DATABASE_NAME:
+                metaType = "STRING METADATA FROM 'database_name' VIRTUAL";
+                break;
+            case OP_TS:
+                metaType = "TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL";
+                break;
+            default:
+                throw new 
UnsupportedOperationException(String.format("Unsupport meta field for MongoDB 
ExtractNode :"
+                                + " %s",
+                        metaFieldInfo.getMetaField()));
+        }
+        return metaType;
+    }
+}
diff --git 
a/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/util/MetaInfoParseUtilTest.java
 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/util/MetaInfoParseUtilTest.java
new file mode 100644
index 000000000..f14c82fbf
--- /dev/null
+++ 
b/inlong-sort/sort-core/src/test/java/org/apache/inlong/sort/util/MetaInfoParseUtilTest.java
@@ -0,0 +1,189 @@
+/*
+ *   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.util;
+
+import org.apache.inlong.common.enums.MetaField;
+import org.apache.inlong.sort.formats.common.IntFormatInfo;
+import org.apache.inlong.sort.formats.common.LongFormatInfo;
+import org.apache.inlong.sort.formats.common.StringFormatInfo;
+import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.MetaFieldInfo;
+import org.apache.inlong.sort.protocol.enums.KafkaScanStartupMode;
+import org.apache.inlong.sort.protocol.node.extract.KafkaExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.MongoExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.OracleExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.PostgresExtractNode;
+import org.apache.inlong.sort.protocol.node.extract.SqlServerExtractNode;
+import org.apache.inlong.sort.protocol.node.format.CanalJsonFormat;
+import org.apache.inlong.sort.protocol.node.format.JsonFormat;
+import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
+import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+/**
+ * Test for {@link MetaInfoParseUtil}
+ */
+public class MetaInfoParseUtilTest {
+
+    /**
+     * Test Mongo meta field
+     */
+    @Test
+    public void testMongoExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("collection_name", 
MetaField.COLLECTION_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        MongoExtractNode mongoExtractNode = new MongoExtractNode("1", "mongo", 
fields,
+                null, null, "test", "localhost:27017",
+                "root", "inlong", "test");
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(mongoExtractNode, metaFieldInfo, sb);
+        Assert.assertEquals("STRING METADATA FROM 'collection_name' VIRTUAL", 
sb.toString());
+    }
+
+    /**
+     * Test Postgres meta field
+     */
+    @Test
+    public void testPostgresExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("database", 
MetaField.DATABASE_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        PostgresExtractNode postgresExtractNode = new PostgresExtractNode("1", 
"postgres_input", fields, null, null,
+                null, Arrays.asList("user"), "localhost", "postgres", 
"inlong", "postgres", "public", 5432, null);
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(postgresExtractNode, metaFieldInfo, 
sb);
+        Assert.assertEquals("STRING METADATA FROM 'database_name' VIRTUAL", 
sb.toString());
+    }
+
+    /**
+     * Test Oracle meta field
+     */
+    @Test
+    public void testOracleExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("schema_name", 
MetaField.SCHEMA_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        OracleExtractNode oracleExtractNode = new OracleExtractNode("1", 
"oracle_input", fields,
+                null, null, "ID", "localhost",
+                "flinkuser", "flinkpw", "xE",
+                "flinkuser", "table", 1521, null);
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(oracleExtractNode, metaFieldInfo, sb);
+        Assert.assertEquals("STRING METADATA FROM 'schema_name' VIRTUAL", 
sb.toString());
+    }
+
+    /**
+     * Test SQLServer meta field
+     */
+    @Test
+    public void testSQLServerExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("database", 
MetaField.DATABASE_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        SqlServerExtractNode sqlServerExtractNode = new 
SqlServerExtractNode("1", "sqlserver_out", fields, null, null,
+                null, "localhost", 1433, "SA", "INLONG*123",
+                "column_type_test", "dbo", "full_types", null);
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(sqlServerExtractNode, metaFieldInfo, 
sb);
+        Assert.assertEquals("STRING METADATA FROM 'database_name' VIRTUAL", 
sb.toString());
+    }
+
+    /**
+     * Test Kafka extract node meta field
+     */
+    @Test
+    public void testKafkaExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("value.database", 
MetaField.DATABASE_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        KafkaExtractNode kafkaExtractNode = new KafkaExtractNode("1", 
"kafka_input", fields, null, null, "workerJson",
+                "localhost:9092", new CanalJsonFormat(), 
KafkaScanStartupMode.EARLIEST_OFFSET, null, "groupId");
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(kafkaExtractNode, metaFieldInfo, sb);
+        Assert.assertEquals("STRING METADATA FROM 'value.database'", 
sb.toString());
+    }
+
+    /**
+     * Test Kafka load node meta field
+     */
+    @Test
+    public void testKafkaLoadNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("database", 
MetaField.DATABASE_NAME);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        List<FieldRelation> relations = Arrays
+                .asList(new FieldRelation(new FieldInfo("name", new 
StringFormatInfo()),
+                                new FieldInfo("name", new StringFormatInfo())),
+                        new FieldRelation(new FieldInfo("age", new 
IntFormatInfo()),
+                                new FieldInfo("age", new IntFormatInfo())));
+        KafkaLoadNode kafkaLoadNode = new KafkaLoadNode("", "kafka_output", 
fields, relations, null, null,
+                "workerJson", "localhost:9092",
+                new JsonFormat(), null,
+                null, null);
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(kafkaLoadNode, metaFieldInfo, sb);
+        Assert.assertEquals("STRING METADATA FROM 'value.database'", 
sb.toString());
+    }
+
+    /**
+     * Test MySQL meta field
+     */
+    @Test
+    public void testMysqlExtractNodeMetaField() {
+        MetaFieldInfo metaFieldInfo = new MetaFieldInfo("meta.op_type", 
MetaField.OP_TYPE);
+        List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new 
LongFormatInfo()),
+                new FieldInfo("name", new StringFormatInfo()),
+                new FieldInfo("age", new IntFormatInfo()),
+                metaFieldInfo
+        );
+        MySqlExtractNode mySqlExtractNode = new MySqlExtractNode("1", 
"mysql_input", fields, null, null,
+                "id", Collections.singletonList("mysql_table"),
+                "localhost", "inlong", "inlong",
+                "inlong", null, null, null, null);
+        StringBuilder sb = new StringBuilder();
+        MetaInfoParseUtil.parseMetaField(mySqlExtractNode, metaFieldInfo, sb);
+        Assert.assertEquals("STRING METADATA FROM 'meta.op_type' VIRTUAL", 
sb.toString());
+    }
+
+}

Reply via email to