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