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());
+ }
+
+}