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 aef0b9dbee [INLONG-10328][Manager] Support automatic synchronization
of stream fields to sink (#10329)
aef0b9dbee is described below
commit aef0b9dbee08ff42f97ded3f5be4b427d5df0269
Author: fuweng11 <[email protected]>
AuthorDate: Wed Jun 5 14:46:43 2024 +0800
[INLONG-10328][Manager] Support automatic synchronization of stream fields
to sink (#10329)
---
.../common/fieldtype/FieldTypeMappingReader.java | 13 +-
.../strategy/ClickHouseFieldTypeStrategy.java | 13 +-
.../strategy/DefaultFieldTypeStrategy.java | 43 +++-
.../strategy/FieldTypeMappingStrategy.java | 13 +-
...Strategy.java => FieldTypeStrategyFactory.java} | 28 ++-
.../strategy/IcebergFieldTypeStrategy.java | 15 +-
.../strategy/MongoDBFieldTypeStrategy.java | 15 +-
.../fieldtype/strategy/MySQLFieldTypeStrategy.java | 15 +-
.../strategy/OracleFieldTypeStrategy.java | 15 +-
.../strategy/PostgreSQLFieldTypeStrategy.java | 15 +-
.../strategy/SQLServerFieldTypeStrategy.java | 15 +-
.../resources/starrocks-field-type-mapping.yaml | 231 +++++++++++++++++++++
.../pojo/sort/node/ExtractNodeProviderFactory.java | 40 +---
.../pojo/sort/node/LoadNodeProviderFactory.java | 53 +----
.../inlong/manager/pojo/sort/node/NodeFactory.java | 30 ++-
.../sort/node/provider/ClickHouseProvider.java | 15 +-
.../pojo/sort/node/provider/DorisProvider.java | 2 +
.../sort/node/provider/ElasticsearchProvider.java | 3 +
.../pojo/sort/node/provider/GreenplumProvider.java | 3 +
.../pojo/sort/node/provider/HBaseProvider.java | 2 +
.../pojo/sort/node/provider/HDFSProvider.java | 2 +
.../pojo/sort/node/provider/HiveProvider.java | 2 +
.../pojo/sort/node/provider/HudiProvider.java | 3 +
.../pojo/sort/node/provider/IcebergProvider.java | 16 +-
.../pojo/sort/node/provider/KafkaProvider.java | 2 +
.../pojo/sort/node/provider/KuduProvider.java | 3 +
.../pojo/sort/node/provider/MongoDBProvider.java | 14 +-
.../sort/node/provider/MySQLBinlogProvider.java | 12 +-
.../pojo/sort/node/provider/MySQLProvider.java | 12 +-
.../pojo/sort/node/provider/OracleProvider.java | 16 +-
.../sort/node/provider/PostgreSQLProvider.java | 17 +-
.../pojo/sort/node/provider/PulsarProvider.java | 7 +-
.../pojo/sort/node/provider/RedisProvider.java | 2 +
.../pojo/sort/node/provider/SQLServerProvider.java | 17 +-
.../pojo/sort/node/provider/StarRocksProvider.java | 2 +
.../node/provider/TDSQLPostgreSQLProvider.java | 3 +
.../pojo/sort/node/provider/TubeMqProvider.java | 3 +
.../manager/pojo/sort/util/FieldInfoUtils.java | 6 +-
.../manager/pojo/stream/InlongStreamInfo.java | 3 +
.../manager/pojo/stream/InlongStreamRequest.java | 3 +
.../resource/sort/SortFlinkConfigOperator.java | 11 +-
.../manager/service/sink/AbstractSinkOperator.java | 31 +++
.../manager/service/sink/StreamSinkOperator.java | 9 +
.../service/stream/InlongStreamServiceImpl.java | 14 +-
44 files changed, 592 insertions(+), 197 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/FieldTypeMappingReader.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/FieldTypeMappingReader.java
index 24ae532df4..6db9d3c8aa 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/FieldTypeMappingReader.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/FieldTypeMappingReader.java
@@ -46,6 +46,11 @@ public class FieldTypeMappingReader implements Serializable {
*/
private static final String SOURCE_TO_TARGET_KEY =
"source.type.to.target.type.converter";
+ /**
+ * Stream type to target type key in converter file.
+ */
+ private static final String STREAM_TO_TARGET_KEY =
"stream.type.to.target.type.converter";
+
/**
* Field type mapping source type key
*/
@@ -59,7 +64,10 @@ public class FieldTypeMappingReader implements Serializable {
@Getter
protected final String streamType;
@Getter
- protected final Map<String, String> FIELD_TYPE_MAPPING_MAP =
Maps.newHashMap();
+ protected final Map<String, String> streamToSinkFieldTypeMap =
Maps.newHashMap();
+
+ @Getter
+ protected final Map<String, String> sourceToSinkFieldTypeMap =
Maps.newHashMap();
public FieldTypeMappingReader(String streamType) {
this.streamType = streamType;
@@ -76,7 +84,8 @@ public class FieldTypeMappingReader implements Serializable {
Yaml yamlReader = new Yaml();
Map<?, ?> converterConf = yamlReader.loadAs(new InputStreamReader(
resource.openStream()), Map.class);
- readerOption(converterConf, SOURCE_TO_TARGET_KEY,
FIELD_TYPE_MAPPING_MAP);
+ readerOption(converterConf, SOURCE_TO_TARGET_KEY,
sourceToSinkFieldTypeMap);
+ readerOption(converterConf, STREAM_TO_TARGET_KEY,
streamToSinkFieldTypeMap);
} catch (Exception e) {
log.error("Yaml reader read option error", e);
throw new RuntimeException(e);
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/ClickHouseFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/ClickHouseFieldTypeStrategy.java
index b464367b04..44e3faa03a 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/ClickHouseFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/ClickHouseFieldTypeStrategy.java
@@ -21,6 +21,7 @@ import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -30,7 +31,8 @@ import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACK
/**
* The ClickHouse field type mapping strategy
*/
-public class ClickHouseFieldTypeStrategy implements FieldTypeMappingStrategy {
+@Service
+public class ClickHouseFieldTypeStrategy extends DefaultFieldTypeStrategy {
private final FieldTypeMappingReader reader;
@@ -43,7 +45,12 @@ public class ClickHouseFieldTypeStrategy implements
FieldTypeMappingStrategy {
}
@Override
- public String getFieldTypeMapping(String sourceType) {
+ public Boolean accept(String type) {
+ return DataNodeType.CLICKHOUSE.equals(type);
+ }
+
+ @Override
+ public String getSourceToSinkFieldTypeMapping(String sourceType) {
// support clickHouse field type special modifier Nullable
if (StringUtils.isNotBlank(sourceType)) {
Matcher matcher = PATTERN.matcher(sourceType.toUpperCase());
@@ -53,6 +60,6 @@ public class ClickHouseFieldTypeStrategy implements
FieldTypeMappingStrategy {
}
}
String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ return reader.getSourceToSinkFieldTypeMap().getOrDefault(dataType,
sourceType.toUpperCase());
}
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/DefaultFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/DefaultFieldTypeStrategy.java
index 21716952d6..752b74f579 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/DefaultFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/DefaultFieldTypeStrategy.java
@@ -17,13 +17,50 @@
package org.apache.inlong.manager.common.fieldtype.strategy;
+import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
+
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
+
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+
/**
* The default field type mapping strategy
*/
-public class DefaultFieldTypeStrategy implements FieldTypeMappingStrategy {
+@Service
+public abstract class DefaultFieldTypeStrategy implements
FieldTypeMappingStrategy {
+
+ private final static String NULLABLE_PATTERN = "^NULLABLE\\((.*)\\)$";
+
+ private static final Pattern PATTERN = Pattern.compile(NULLABLE_PATTERN);
+
+ protected FieldTypeMappingReader reader = null;
+
+ @Override
+ public String getSourceToSinkFieldTypeMapping(String sourceType) {
+ if (reader == null) {
+ return sourceType;
+ }
+ String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
+ return reader.getSourceToSinkFieldTypeMap().getOrDefault(dataType,
sourceType.toUpperCase());
+ }
@Override
- public String getFieldTypeMapping(String sourceType) {
- return sourceType;
+ public String getStreamToSinkFieldTypeMapping(String sourceType) {
+ if (reader == null) {
+ return sourceType;
+ }
+ if (StringUtils.isNotBlank(sourceType)) {
+ Matcher matcher = PATTERN.matcher(sourceType.toUpperCase());
+ if (matcher.matches()) {
+ // obtain the field type modified by Nullable, for example,
uint8(12) in Nullable(uint8(12))
+ sourceType = matcher.group(1);
+ }
+ }
+ String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
+ return reader.getStreamToSinkFieldTypeMap().getOrDefault(dataType,
sourceType.toUpperCase());
}
}
\ No newline at end of file
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
index b456170bb6..a4442e7738 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
@@ -22,11 +22,22 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
*/
public interface FieldTypeMappingStrategy {
+ Boolean accept(String type);
+
/**
* Get the field type of inlong field type mapping by the source field
type.
*
* @param sourceType the source field type
* @return the target field type of inlong field type mapping
*/
- String getFieldTypeMapping(String sourceType);
+ String getSourceToSinkFieldTypeMapping(String sourceType);
+
+ /**
+ * Get the field type of inlong field type mapping by the stream field
type.
+ *
+ * @param streamType the stream field type
+ * @return the target field type of inlong field type mapping
+ */
+ String getStreamToSinkFieldTypeMapping(String streamType);
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeStrategyFactory.java
similarity index 58%
copy from
inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
copy to
inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeStrategyFactory.java
index b456170bb6..c072310108 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeMappingStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/FieldTypeStrategyFactory.java
@@ -17,16 +17,32 @@
package org.apache.inlong.manager.common.fieldtype.strategy;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.List;
+
/**
- * The interface of base field type mapping strategy operation.
+ * Factory for {@link FieldTypeMappingStrategy}.
*/
-public interface FieldTypeMappingStrategy {
+@Service
+@Slf4j
+public class FieldTypeStrategyFactory {
+
+ @Autowired
+ private List<FieldTypeMappingStrategy> strategyList;
/**
- * Get the field type of inlong field type mapping by the source field
type.
+ * Get a field type mapping strategy instance.
*
- * @param sourceType the source field type
- * @return the target field type of inlong field type mapping
+ * @param type type
*/
- String getFieldTypeMapping(String sourceType);
+ public FieldTypeMappingStrategy getInstance(String type) {
+ return strategyList.stream()
+ .filter(inst -> inst.accept(type))
+ .findFirst()
+ .orElse(null);
+ }
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/IcebergFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/IcebergFieldTypeStrategy.java
index 2555fc3e5e..383fe99a33 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/IcebergFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/IcebergFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The iceberg field type mapping strategy
*/
-public class IcebergFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class IcebergFieldTypeStrategy extends DefaultFieldTypeStrategy {
public IcebergFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.ICEBERG);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.ICEBERG.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MongoDBFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MongoDBFieldTypeStrategy.java
index 0e760b87c0..3b0e510f00 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MongoDBFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MongoDBFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The mongoDB field type mapping strategy
*/
-public class MongoDBFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class MongoDBFieldTypeStrategy extends DefaultFieldTypeStrategy {
public MongoDBFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.MONGODB);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.MONGODB.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MySQLFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MySQLFieldTypeStrategy.java
index 114519a6ac..fcf249e7a0 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MySQLFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/MySQLFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The mysql field type mapping strategy
*/
-public class MySQLFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class MySQLFieldTypeStrategy extends DefaultFieldTypeStrategy {
public MySQLFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.MYSQL);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.MYSQL.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/OracleFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/OracleFieldTypeStrategy.java
index f41cc2af29..7de753da44 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/OracleFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/OracleFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The oracle field type mapping strategy
*/
-public class OracleFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class OracleFieldTypeStrategy extends DefaultFieldTypeStrategy {
public OracleFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.ORACLE);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.ORACLE.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/PostgreSQLFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/PostgreSQLFieldTypeStrategy.java
index 17e9c98d7b..1606050e6a 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/PostgreSQLFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/PostgreSQLFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The postgresql field type mapping strategy
*/
-public class PostgreSQLFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class PostgreSQLFieldTypeStrategy extends DefaultFieldTypeStrategy {
public PostgreSQLFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.POSTGRESQL);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.POSTGRESQL.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/SQLServerFieldTypeStrategy.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/SQLServerFieldTypeStrategy.java
index 8fcecd94c6..725d497bde 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/SQLServerFieldTypeStrategy.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/fieldtype/strategy/SQLServerFieldTypeStrategy.java
@@ -20,24 +20,21 @@ package org.apache.inlong.manager.common.fieldtype.strategy;
import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.fieldtype.FieldTypeMappingReader;
-import org.apache.commons.lang3.StringUtils;
-
-import static
org.apache.inlong.manager.common.consts.InlongConstants.LEFT_BRACKET;
+import org.springframework.stereotype.Service;
/**
* The sqlServer field type mapping strategy
*/
-public class SQLServerFieldTypeStrategy implements FieldTypeMappingStrategy {
-
- private final FieldTypeMappingReader reader;
+@Service
+public class SQLServerFieldTypeStrategy extends DefaultFieldTypeStrategy {
public SQLServerFieldTypeStrategy() {
this.reader = new FieldTypeMappingReader(DataNodeType.SQLSERVER);
}
@Override
- public String getFieldTypeMapping(String sourceType) {
- String dataType = StringUtils.substringBefore(sourceType,
LEFT_BRACKET).toUpperCase();
- return reader.getFIELD_TYPE_MAPPING_MAP().getOrDefault(dataType,
sourceType.toUpperCase());
+ public Boolean accept(String type) {
+ return DataNodeType.SQLSERVER.equals(type);
}
+
}
diff --git
a/inlong-manager/manager-common/src/main/resources/starrocks-field-type-mapping.yaml
b/inlong-manager/manager-common/src/main/resources/starrocks-field-type-mapping.yaml
new file mode 100644
index 0000000000..d868b93fd6
--- /dev/null
+++
b/inlong-manager/manager-common/src/main/resources/starrocks-field-type-mapping.yaml
@@ -0,0 +1,231 @@
+#
+# 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.
+#
+
+stream.type.to.target.type.converter:
+
+ - source.type: TINYINT
+ target.type: BIGINT
+
+ - source.type: SMALLINT
+ target.type: BIGINT
+
+ - source.type: TINYINT UNSIGNED
+ target.type: BIGINT
+
+ - source.type: TINYINT UNSIGNED ZEROFILL
+ target.type: BIGINT
+
+ - source.type: INT
+ target.type: BIGINT
+
+ - source.type: INTEGER
+ target.type: BIGINT
+
+ - source.type: YEAR
+ target.type: BIGINT
+
+ - source.type: SHORT
+ target.type: BIGINT
+
+ - source.type: MEDIUMINT
+ target.type: BIGINT
+
+ - source.type: SMALLINT UNSIGNED
+ target.type: BIGINT
+
+ - source.type: SMALLINT UNSIGNED ZEROFILL
+ target.type: BIGINT
+
+ - source.type: BIGINT
+ target.type: BIGINT
+
+ - source.type: INT UNSIGNED
+ target.type: BIGINT
+
+ - source.type: MEDIUMINT UNSIGNED
+ target.type: BIGINT
+
+ - source.type: MEDIUMINT UNSIGNED ZEROFILL
+ target.type: BIGINT
+
+ - source.type: INT UNSIGNED ZEROFILL
+ target.type: BIGINT
+
+ - source.type: BIGINT UNSIGNED
+ target.type: BIGINT
+
+ - source.type: LONG
+ target.type: BIGINT
+
+ - source.type: BIGINT UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: SERIAL
+ target.type: DOUBLE
+
+ - source.type: FLOAT
+ target.type: DOUBLE
+
+ - source.type: FLOAT UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: FLOAT UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: DOUBLE
+ target.type: DOUBLE
+
+ - source.type: DOUBLE UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: DOUBLE UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: DOUBLE PRECISION
+ target.type: DOUBLE
+
+ - source.type: DOUBLE PRECISION UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: REAL
+ target.type: DOUBLE
+
+ - source.type: REAL UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: REAL UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: NUMERIC
+ target.type: DOUBLE
+
+ - source.type: NUMERIC UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: NUMERIC UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: DECIMAL
+ target.type: DOUBLE
+
+ - source.type: DECIMAL UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: DECIMAL UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: FIXED
+ target.type: DOUBLE
+
+ - source.type: FIXED UNSIGNED
+ target.type: DOUBLE
+
+ - source.type: FIXED UNSIGNED ZEROFILL
+ target.type: DOUBLE
+
+ - source.type: BOOLEAN
+ target.type: DOUBLE
+
+ - source.type: DATE
+ target.type: VARCHAR
+
+ - source.type: TIME
+ target.type: VARCHAR
+
+ - source.type: DATETIME
+ target.type: VARCHAR
+
+ - source.type: TIMESTAMP
+ target.type: VARCHAR
+
+ - source.type: CHAR
+ target.type: VARCHAR
+
+ - source.type: JSON
+ target.type: VARCHAR
+
+ - source.type: BIT
+ target.type: VARCHAR
+
+ - source.type: VARCHAR
+ target.type: VARCHAR
+
+ - source.type: TEXT
+ target.type: VARCHAR
+
+ - source.type: BLOB
+ target.type: VARCHAR
+
+ - source.type: TINYBLOB
+ target.type: VARCHAR
+
+ - source.type: TINYTEXT
+ target.type: VARCHAR
+
+ - source.type: MEDIUMBLOB
+ target.type: VARCHAR
+
+ - source.type: MEDIUMTEXT
+ target.type: VARCHAR
+
+ - source.type: LONGBLOB
+ target.type: VARCHAR
+
+ - source.type: LONGTEXT
+ target.type: VARCHAR
+
+ - source.type: VARBINARY
+ target.type: VARCHAR
+
+ - source.type: GEOMETRY
+ target.type: VARCHAR
+
+ - source.type: POINT
+ target.type: VARCHAR
+
+ - source.type: LINESTRING
+ target.type: VARCHAR
+
+ - source.type: POLYGON
+ target.type: VARCHAR
+
+ - source.type: MULTIPOINT
+ target.type: VARCHAR
+
+ - source.type: MULTILINESTRING
+ target.type: VARCHAR
+
+ - source.type: MULTIPOLYGON
+ target.type: VARCHAR
+
+ - source.type: GEOMETRYCOLLECTION
+ target.type: VARCHAR
+
+ - source.type: ENUM
+ target.type: VARCHAR
+
+ - source.type: STRING
+ target.type: VARCHAR
+
+ - source.type: BINARY
+ target.type: VARCHAR
+
+ - source.type: BYTE
+ target.type: VARCHAR
\ No newline at end of file
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/ExtractNodeProviderFactory.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/ExtractNodeProviderFactory.java
index 07da9d3b5a..ac3f457f0a 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/ExtractNodeProviderFactory.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/ExtractNodeProviderFactory.java
@@ -20,17 +20,10 @@ package org.apache.inlong.manager.pojo.sort.node;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.HudiProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.IcebergProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.KafkaProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.MongoDBProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.MySQLBinlogProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.OracleProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.PostgreSQLProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.PulsarProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.RedisProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.SQLServerProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.TubeMqProvider;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -38,28 +31,15 @@ import java.util.List;
/**
* Factory of the extract node provider.
*/
+@Service
+@Slf4j
public class ExtractNodeProviderFactory {
/**
* The extract node provider collection
*/
- private static final List<ExtractNodeProvider> EXTRACT_NODE_PROVIDER_LIST
= new ArrayList<>();
-
- static {
- // The Providers Parsing SourceInfo to ExtractNode which sort needed
- EXTRACT_NODE_PROVIDER_LIST.add(new HudiProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new KafkaProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new MongoDBProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new OracleProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new PulsarProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new RedisProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new TubeMqProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new SQLServerProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new PostgreSQLProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new MySQLBinlogProvider());
- EXTRACT_NODE_PROVIDER_LIST.add(new IcebergProvider());
-
- }
+ @Autowired
+ private List<ExtractNodeProvider> extractNodeProviderList = new
ArrayList<>();
/**
* Get extract node provider
@@ -67,8 +47,8 @@ public class ExtractNodeProviderFactory {
* @param sourceType the specified source type
* @return the extract node provider
*/
- public static ExtractNodeProvider getExtractNodeProvider(String
sourceType) {
- return EXTRACT_NODE_PROVIDER_LIST.stream()
+ public ExtractNodeProvider getExtractNodeProvider(String sourceType) {
+ return extractNodeProviderList.stream()
.filter(inst -> inst.accept(sourceType))
.findFirst()
.orElseThrow(() -> new
BusinessException(ErrorCodeEnum.SOURCE_TYPE_NOT_SUPPORT,
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/LoadNodeProviderFactory.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/LoadNodeProviderFactory.java
index 1df97b77fb..742b346fcc 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/LoadNodeProviderFactory.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/LoadNodeProviderFactory.java
@@ -20,24 +20,10 @@ package org.apache.inlong.manager.pojo.sort.node;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.ClickHouseProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.DorisProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.ElasticsearchProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.GreenplumProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.HBaseProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.HDFSProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.HiveProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.HudiProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.IcebergProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.KafkaProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.KuduProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.MySQLProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.OracleProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.PostgreSQLProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.RedisProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.SQLServerProvider;
-import org.apache.inlong.manager.pojo.sort.node.provider.StarRocksProvider;
-import
org.apache.inlong.manager.pojo.sort.node.provider.TDSQLPostgreSQLProvider;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -45,34 +31,15 @@ import java.util.List;
/**
* Factory of the load node provider.
*/
+@Service
+@Slf4j
public class LoadNodeProviderFactory {
/**
* The load node provider collection
*/
- private static final List<LoadNodeProvider> LOAD_NODE_PROVIDER_LIST = new
ArrayList<>();
-
- static {
- // The Providers Parsing SinkInfo to LoadNode which sort needed
- LOAD_NODE_PROVIDER_LIST.add(new KafkaProvider());
- LOAD_NODE_PROVIDER_LIST.add(new ClickHouseProvider());
- LOAD_NODE_PROVIDER_LIST.add(new DorisProvider());
- LOAD_NODE_PROVIDER_LIST.add(new ElasticsearchProvider());
- LOAD_NODE_PROVIDER_LIST.add(new GreenplumProvider());
- LOAD_NODE_PROVIDER_LIST.add(new HBaseProvider());
- LOAD_NODE_PROVIDER_LIST.add(new HDFSProvider());
- LOAD_NODE_PROVIDER_LIST.add(new HiveProvider());
- LOAD_NODE_PROVIDER_LIST.add(new HudiProvider());
- LOAD_NODE_PROVIDER_LIST.add(new IcebergProvider());
- LOAD_NODE_PROVIDER_LIST.add(new KuduProvider());
- LOAD_NODE_PROVIDER_LIST.add(new MySQLProvider());
- LOAD_NODE_PROVIDER_LIST.add(new OracleProvider());
- LOAD_NODE_PROVIDER_LIST.add(new PostgreSQLProvider());
- LOAD_NODE_PROVIDER_LIST.add(new RedisProvider());
- LOAD_NODE_PROVIDER_LIST.add(new SQLServerProvider());
- LOAD_NODE_PROVIDER_LIST.add(new StarRocksProvider());
- LOAD_NODE_PROVIDER_LIST.add(new TDSQLPostgreSQLProvider());
- }
+ @Autowired
+ private List<LoadNodeProvider> loadNodeProviderList = new ArrayList<>();
/**
* Get load node provider
@@ -80,8 +47,8 @@ public class LoadNodeProviderFactory {
* @param sinkType the specified sink type
* @return the load node provider
*/
- public static LoadNodeProvider getLoadNodeProvider(String sinkType) {
- return LOAD_NODE_PROVIDER_LIST.stream()
+ public LoadNodeProvider getLoadNodeProvider(String sinkType) {
+ return loadNodeProviderList.stream()
.filter(inst -> inst.accept(sinkType))
.findFirst()
.orElseThrow(() -> new
BusinessException(ErrorCodeEnum.SINK_TYPE_NOT_SUPPORT,
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/NodeFactory.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/NodeFactory.java
index 88a7992f78..3cda291dbe 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/NodeFactory.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/NodeFactory.java
@@ -33,6 +33,8 @@ import
org.apache.inlong.sort.protocol.node.transform.TransformNode;
import com.google.common.collect.Lists;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -43,65 +45,71 @@ import java.util.stream.Collectors;
* The node factory
*/
@Slf4j
+@Service
public class NodeFactory {
+ @Autowired
+ private LoadNodeProviderFactory loadNodeProviderFactory;
+ @Autowired
+ private ExtractNodeProviderFactory extractNodeProviderFactory;
+
/**
* Create extract nodes from the given sources.
*/
- public static List<ExtractNode> createExtractNodes(List<StreamSource>
sourceInfos) {
+ public List<ExtractNode> createExtractNodes(List<StreamSource>
sourceInfos) {
if (CollectionUtils.isEmpty(sourceInfos)) {
return Lists.newArrayList();
}
return sourceInfos.stream().map(v -> {
String sourceType = v.getSourceType();
- return
ExtractNodeProviderFactory.getExtractNodeProvider(sourceType).createExtractNode(v);
+ return
extractNodeProviderFactory.getExtractNodeProvider(sourceType).createExtractNode(v);
}).collect(Collectors.toList());
}
/**
* Create load nodes from the given sinks.
*/
- public static List<LoadNode> createLoadNodes(List<StreamSink> sinkInfos,
+ public List<LoadNode> createLoadNodes(List<StreamSink> sinkInfos,
Map<String, StreamField> constantFieldMap) {
if (CollectionUtils.isEmpty(sinkInfos)) {
return Lists.newArrayList();
}
return sinkInfos.stream().map(v -> {
String sinkType = v.getSinkType();
- return
LoadNodeProviderFactory.getLoadNodeProvider(sinkType).createLoadNode(v,
constantFieldMap);
+ return
loadNodeProviderFactory.getLoadNodeProvider(sinkType).createLoadNode(v,
constantFieldMap);
}).collect(Collectors.toList());
}
/**
* Create extract node from the given source.
*/
- public static ExtractNode createExtractNode(StreamSource sourceInfo) {
+ public ExtractNode createExtractNode(StreamSource sourceInfo) {
if (sourceInfo == null) {
return null;
}
String sourceType = sourceInfo.getSourceType();
- return
ExtractNodeProviderFactory.getExtractNodeProvider(sourceType).createExtractNode(sourceInfo);
+ return
extractNodeProviderFactory.getExtractNodeProvider(sourceType).createExtractNode(sourceInfo);
}
/**
* Create load node from the given sink.
*/
- public static LoadNode createLoadNode(StreamSink sinkInfo, Map<String,
StreamField> constantFieldMap) {
+ public LoadNode createLoadNode(StreamSink sinkInfo, Map<String,
StreamField> constantFieldMap) {
if (sinkInfo == null) {
return null;
}
String sinkType = sinkInfo.getSinkType();
- return
LoadNodeProviderFactory.getLoadNodeProvider(sinkType).createLoadNode(sinkInfo,
constantFieldMap);
+ return
loadNodeProviderFactory.getLoadNodeProvider(sinkType).createLoadNode(sinkInfo,
constantFieldMap);
}
/**
* Add built-in field for extra node and load node
*/
- public static List<Node> addBuiltInField(StreamSource sourceInfo,
StreamSink sinkInfo,
+ public List<Node> addBuiltInField(StreamSource sourceInfo, StreamSink
sinkInfo,
List<TransformResponse> transformResponses, Map<String,
StreamField> constantFieldMap) {
- ExtractNodeProvider extractNodeProvider =
ExtractNodeProviderFactory.getExtractNodeProvider(
+ ExtractNodeProvider extractNodeProvider =
extractNodeProviderFactory.getExtractNodeProvider(
sourceInfo.getSourceType());
- LoadNodeProvider loadNodeProvider =
LoadNodeProviderFactory.getLoadNodeProvider(sinkInfo.getSinkType());
+ LoadNodeProvider loadNodeProvider =
loadNodeProviderFactory.getLoadNodeProvider(sinkInfo.getSinkType());
if (loadNodeProvider.isSinkMultiple(sinkInfo)) {
sourceInfo.setFieldList(loadNodeProvider.addStreamFieldsForSinkMultiple(sourceInfo.getFieldList()));
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ClickHouseProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ClickHouseProvider.java
index 8bab226fb8..6b1e7cebe3 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ClickHouseProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ClickHouseProvider.java
@@ -18,8 +18,8 @@
package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.SinkType;
-import
org.apache.inlong.manager.common.fieldtype.strategy.ClickHouseFieldTypeStrategy;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.ck.ClickHouseSink;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
import org.apache.inlong.manager.pojo.stream.StreamField;
@@ -29,15 +29,20 @@ import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.load.ClickHouseLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating ClickHouse load nodes.
*/
+@Service
public class ClickHouseProvider implements LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new ClickHouseFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String sinkType) {
@@ -48,8 +53,10 @@ public class ClickHouseProvider implements LoadNodeProvider {
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
ClickHouseSink streamSink = (ClickHouseSink) nodeInfo;
Map<String, String> properties =
parseProperties(streamSink.getProperties());
- List<FieldInfo> fieldInfos =
parseSinkFieldInfos(streamSink.getSinkFieldList(), streamSink.getSinkName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance((streamSink.getSinkType()));
+ List<FieldInfo> fieldInfos =
+ parseSinkFieldInfos(streamSink.getSinkFieldList(),
streamSink.getSinkName(), fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(streamSink.getSinkFieldList(), constantFieldMap);
return new ClickHouseLoadNode(
streamSink.getSinkName(),
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/DorisProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/DorisProvider.java
index a87db38690..515ffb0936 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/DorisProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/DorisProvider.java
@@ -32,6 +32,7 @@ import
org.apache.inlong.sort.protocol.transformation.FieldRelation;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -41,6 +42,7 @@ import java.util.Map;
* The Provider for creating Doris load nodes.
*/
@Slf4j
+@Service
public class DorisProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ElasticsearchProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ElasticsearchProvider.java
index 30360cf21d..6ac8642cc1 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ElasticsearchProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/ElasticsearchProvider.java
@@ -29,12 +29,15 @@ import org.apache.inlong.sort.protocol.node.format.Format;
import org.apache.inlong.sort.protocol.node.load.ElasticsearchLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating Elasticsearch load nodes.
*/
+@Service
public class ElasticsearchProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/GreenplumProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/GreenplumProvider.java
index df5618bce0..3e86bdba80 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/GreenplumProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/GreenplumProvider.java
@@ -27,12 +27,15 @@ import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.load.GreenplumLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating Greenplum load nodes.
*/
+@Service
public class GreenplumProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HBaseProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HBaseProvider.java
index b00ae617cf..c5f86d9802 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HBaseProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HBaseProvider.java
@@ -28,6 +28,7 @@ import
org.apache.inlong.sort.protocol.node.load.HbaseLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import com.google.common.collect.Lists;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -35,6 +36,7 @@ import java.util.Map;
/**
* The Provider for creating HBase load nodes.
*/
+@Service
public class HBaseProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HDFSProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HDFSProvider.java
index 94d5415043..29561bed2b 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HDFSProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HDFSProvider.java
@@ -30,6 +30,7 @@ import
org.apache.inlong.sort.protocol.transformation.FieldRelation;
import com.google.common.collect.Lists;
import org.apache.commons.collections.CollectionUtils;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -38,6 +39,7 @@ import java.util.stream.Collectors;
/**
* The Provider for creating HDFS load nodes.
*/
+@Service
public class HDFSProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HiveProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HiveProvider.java
index 6595409def..295dd7c647 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HiveProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HiveProvider.java
@@ -30,6 +30,7 @@ import
org.apache.inlong.sort.protocol.transformation.FieldRelation;
import com.google.common.collect.Lists;
import org.apache.commons.collections.CollectionUtils;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -38,6 +39,7 @@ import java.util.stream.Collectors;
/**
* The Provider for creating Hive load nodes.
*/
+@Service
public class HiveProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HudiProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HudiProvider.java
index 24dd749e41..2b2d291f54 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HudiProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/HudiProvider.java
@@ -33,12 +33,15 @@ import
org.apache.inlong.sort.protocol.node.extract.HudiExtractNode;
import org.apache.inlong.sort.protocol.node.load.HudiLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating Hudi extract or load nodes.
*/
+@Service
public class HudiProvider implements ExtractNodeProvider, LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/IcebergProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/IcebergProvider.java
index 669f7c2217..e6004db45f 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/IcebergProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/IcebergProvider.java
@@ -20,7 +20,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.common.enums.MetaField;
import org.apache.inlong.manager.common.consts.StreamType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.IcebergFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.SinkField;
import org.apache.inlong.manager.pojo.sink.iceberg.IcebergSink;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
@@ -41,6 +41,8 @@ import
org.apache.inlong.sort.protocol.transformation.FieldRelation;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -51,9 +53,11 @@ import java.util.stream.Collectors;
* The Provider for creating Iceberg load nodes.
*/
@Slf4j
+@Service
public class IcebergProvider implements ExtractNodeProvider, LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new IcebergFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String sinkType) {
@@ -63,8 +67,10 @@ public class IcebergProvider implements ExtractNodeProvider,
LoadNodeProvider {
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
IcebergSource icebergSource = (IcebergSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(icebergSource.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(icebergSource.getFieldList(),
icebergSource.getSourceName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(icebergSource.getProperties());
return new IcebergExtractNode(icebergSource.getSourceName(),
@@ -86,8 +92,10 @@ public class IcebergProvider implements ExtractNodeProvider,
LoadNodeProvider {
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
IcebergSink icebergSink = (IcebergSink) nodeInfo;
Map<String, String> properties =
parseProperties(icebergSink.getProperties());
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(icebergSink.getSinkType());
List<FieldInfo> fieldInfos =
parseSinkFieldInfos(icebergSink.getSinkFieldList(), icebergSink.getSinkName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(icebergSink.getSinkFieldList(), constantFieldMap);
IcebergConstant.CatalogType catalogType =
CatalogType.forName(icebergSink.getCatalogType());
Format format =
parsingSinkMultipleFormat(icebergSink.getSinkMultipleEnable(),
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KafkaProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KafkaProvider.java
index ecfa418f2a..9e602c0293 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KafkaProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KafkaProvider.java
@@ -42,6 +42,7 @@ import
org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import com.google.common.collect.Lists;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -49,6 +50,7 @@ import java.util.Map;
/**
* The Provider for creating Kafka extract or load nodes.
*/
+@Service
public class KafkaProvider implements ExtractNodeProvider, LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KuduProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KuduProvider.java
index 3ddd3e5e40..7d27be81e6 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KuduProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/KuduProvider.java
@@ -27,12 +27,15 @@ import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.load.KuduLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating Kudu load nodes.
*/
+@Service
public class KuduProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MongoDBProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MongoDBProvider.java
index f93db2b064..70768bf6d5 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MongoDBProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MongoDBProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.SourceType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.MongoDBFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
import org.apache.inlong.manager.pojo.source.mongodb.MongoDBSource;
import org.apache.inlong.manager.pojo.stream.StreamNode;
@@ -27,16 +27,20 @@ import org.apache.inlong.sort.protocol.FieldInfo;
import org.apache.inlong.sort.protocol.node.ExtractNode;
import org.apache.inlong.sort.protocol.node.extract.MongoExtractNode;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating MongoDB extract nodes.
*/
+@Service
public class MongoDBProvider implements ExtractNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new MongoDBFieldTypeStrategy();
-
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String sourceType) {
return SourceType.MONGODB.equals(sourceType);
@@ -45,8 +49,10 @@ public class MongoDBProvider implements ExtractNodeProvider {
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
MongoDBSource source = (MongoDBSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+ fieldTypeStrategyFactory.getInstance(source.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(source.getFieldList(), source.getSourceName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(source.getProperties());
return new MongoExtractNode(
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLBinlogProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLBinlogProvider.java
index 4016d3ff71..2c444eaf92 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLBinlogProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLBinlogProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.SourceType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.MySQLFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
import org.apache.inlong.manager.pojo.source.mysql.MySQLBinlogSource;
import org.apache.inlong.manager.pojo.stream.StreamNode;
@@ -28,6 +28,8 @@ import org.apache.inlong.sort.protocol.node.ExtractNode;
import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
import com.google.common.base.Splitter;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -35,9 +37,11 @@ import java.util.Map;
/**
* The Provider for creating MySQLBinlog extract nodes.
*/
+@Service
public class MySQLBinlogProvider implements ExtractNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new MySQLFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String sourceType) {
@@ -47,8 +51,10 @@ public class MySQLBinlogProvider implements
ExtractNodeProvider {
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
MySQLBinlogSource binlogSource = (MySQLBinlogSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(binlogSource.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(binlogSource.getFieldList(), binlogSource.getSourceName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(binlogSource.getProperties());
final String database = binlogSource.getDatabaseWhiteList();
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLProvider.java
index a366688bc9..4509570c69 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/MySQLProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.SinkType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.MySQLFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.mysql.MySQLSink;
import org.apache.inlong.manager.pojo.sink.mysql.MySQLSinkDTO;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
@@ -31,6 +31,8 @@ import
org.apache.inlong.sort.protocol.node.load.MySqlLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import com.google.common.collect.Lists;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -38,9 +40,11 @@ import java.util.Map;
/**
* The Provider for creating MySQL load nodes.
*/
+@Service
public class MySQLProvider implements LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new MySQLFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String sinkType) {
@@ -51,8 +55,10 @@ public class MySQLProvider implements LoadNodeProvider {
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
MySQLSink mysqlSink = (MySQLSink) nodeInfo;
Map<String, String> properties =
parseProperties(mysqlSink.getProperties());
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+ fieldTypeStrategyFactory.getInstance(mysqlSink.getSinkType());
List<FieldInfo> fieldInfos =
parseSinkFieldInfos(mysqlSink.getSinkFieldList(), mysqlSink.getSinkName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(mysqlSink.getSinkFieldList(), constantFieldMap);
return new MySqlLoadNode(
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/OracleProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/OracleProvider.java
index ea5712c15f..3824657171 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/OracleProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/OracleProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.StreamType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.OracleFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.oracle.OracleSink;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
@@ -35,6 +35,8 @@ import
org.apache.inlong.sort.protocol.node.load.OracleLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import org.apache.commons.lang3.StringUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -42,9 +44,11 @@ import java.util.Map;
/**
* The Provider for creating Oracle extract or load nodes.
*/
+@Service
public class OracleProvider implements ExtractNodeProvider, LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new OracleFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String streamType) {
@@ -54,8 +58,10 @@ public class OracleProvider implements ExtractNodeProvider,
LoadNodeProvider {
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
OracleSource source = (OracleSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+ fieldTypeStrategyFactory.getInstance(source.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(source.getFieldList(), source.getSourceName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(source.getProperties());
ScanStartUpMode scanStartupMode =
StringUtils.isBlank(source.getScanStartupMode())
@@ -82,8 +88,10 @@ public class OracleProvider implements ExtractNodeProvider,
LoadNodeProvider {
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
OracleSink oracleSink = (OracleSink) nodeInfo;
Map<String, String> properties =
parseProperties(oracleSink.getProperties());
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+ fieldTypeStrategyFactory.getInstance(oracleSink.getSinkType());
List<FieldInfo> fieldInfos =
parseSinkFieldInfos(oracleSink.getSinkFieldList(), oracleSink.getSinkName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(oracleSink.getSinkFieldList(), constantFieldMap);
return new OracleLoadNode(
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PostgreSQLProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PostgreSQLProvider.java
index 17c15e88b0..f36f2aeb21 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PostgreSQLProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PostgreSQLProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.StreamType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.PostgreSQLFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.postgresql.PostgreSQLSink;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
@@ -33,15 +33,20 @@ import
org.apache.inlong.sort.protocol.node.extract.PostgresExtractNode;
import org.apache.inlong.sort.protocol.node.load.PostgresLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating PostgreSQL extract or load nodes.
*/
+@Service
public class PostgreSQLProvider implements ExtractNodeProvider,
LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new PostgreSQLFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String streamType) {
@@ -51,8 +56,10 @@ public class PostgreSQLProvider implements
ExtractNodeProvider, LoadNodeProvider
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
PostgreSQLSource postgreSQLSource = (PostgreSQLSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(postgreSQLSource.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(postgreSQLSource.getFieldList(),
- postgreSQLSource.getSourceName(), FIELD_TYPE_MAPPING_STRATEGY);
+ postgreSQLSource.getSourceName(), fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(postgreSQLSource.getProperties());
return new PostgresExtractNode(postgreSQLSource.getSourceName(),
@@ -77,8 +84,10 @@ public class PostgreSQLProvider implements
ExtractNodeProvider, LoadNodeProvider
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
PostgreSQLSink postgreSQLSink = (PostgreSQLSink) nodeInfo;
Map<String, String> properties =
parseProperties(postgreSQLSink.getProperties());
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(postgreSQLSink.getSinkType());
List<FieldInfo> fieldInfos =
parseSinkFieldInfos(postgreSQLSink.getSinkFieldList(),
- postgreSQLSink.getSinkName(), FIELD_TYPE_MAPPING_STRATEGY);
+ postgreSQLSink.getSinkName(), fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(postgreSQLSink.getSinkFieldList(), constantFieldMap);
return new PostgresLoadNode(
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PulsarProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PulsarProvider.java
index 815a2c945d..1d8ace1480 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PulsarProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/PulsarProvider.java
@@ -30,7 +30,7 @@ import org.apache.inlong.sort.protocol.node.ExtractNode;
import org.apache.inlong.sort.protocol.node.extract.PulsarExtractNode;
import org.apache.inlong.sort.protocol.node.format.Format;
-import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -40,6 +40,7 @@ import java.util.stream.Collectors;
/**
* The Provider for creating Pulsar extract nodes.
*/
+@Service
public class PulsarProvider implements ExtractNodeProvider {
@Override
@@ -64,9 +65,7 @@ public class PulsarProvider implements ExtractNodeProvider {
final String primaryKey = pulsarSource.getPrimaryKey();
final String serviceUrl = pulsarSource.getServiceUrl();
final String adminUrl = pulsarSource.getAdminUrl();
- final String scanStartupSubStartOffset =
- StringUtils.isNotBlank(pulsarSource.getSubscription()) ?
PulsarScanStartupMode.EARLIEST.getValue()
- : null;
+ final String scanStartupSubStartOffset = null;
return new PulsarExtractNode(pulsarSource.getSourceName(),
pulsarSource.getSourceName(),
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/RedisProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/RedisProvider.java
index 3fa9ccf2be..4566f593b6 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/RedisProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/RedisProvider.java
@@ -45,6 +45,7 @@ import
org.apache.inlong.sort.protocol.node.load.RedisLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Map;
@@ -52,6 +53,7 @@ import java.util.Map;
/**
* The Provider for creating Redis extract or load nodes.
*/
+@Service
public class RedisProvider implements ExtractNodeProvider, LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/SQLServerProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/SQLServerProvider.java
index c73f2f8507..31759ea589 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/SQLServerProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/SQLServerProvider.java
@@ -19,7 +19,7 @@ package org.apache.inlong.manager.pojo.sort.node.provider;
import org.apache.inlong.manager.common.consts.StreamType;
import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
-import
org.apache.inlong.manager.common.fieldtype.strategy.SQLServerFieldTypeStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.pojo.sink.sqlserver.SQLServerSink;
import org.apache.inlong.manager.pojo.sort.node.base.ExtractNodeProvider;
import org.apache.inlong.manager.pojo.sort.node.base.LoadNodeProvider;
@@ -33,15 +33,20 @@ import
org.apache.inlong.sort.protocol.node.extract.SqlServerExtractNode;
import org.apache.inlong.sort.protocol.node.load.SqlServerLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating SQLServer extract or load nodes.
*/
+@Service
public class SQLServerProvider implements ExtractNodeProvider,
LoadNodeProvider {
- private static final FieldTypeMappingStrategy FIELD_TYPE_MAPPING_STRATEGY
= new SQLServerFieldTypeStrategy();
+ @Autowired
+ private FieldTypeStrategyFactory fieldTypeStrategyFactory;
@Override
public Boolean accept(String streamType) {
@@ -51,8 +56,10 @@ public class SQLServerProvider implements
ExtractNodeProvider, LoadNodeProvider
@Override
public ExtractNode createExtractNode(StreamNode streamNodeInfo) {
SQLServerSource source = (SQLServerSource) streamNodeInfo;
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+ fieldTypeStrategyFactory.getInstance(source.getSourceType());
List<FieldInfo> fieldInfos =
parseStreamFieldInfos(source.getFieldList(), source.getSourceName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
Map<String, String> properties =
parseProperties(source.getProperties());
return new SqlServerExtractNode(
@@ -76,8 +83,10 @@ public class SQLServerProvider implements
ExtractNodeProvider, LoadNodeProvider
public LoadNode createLoadNode(StreamNode nodeInfo, Map<String,
StreamField> constantFieldMap) {
SQLServerSink sqlServerSink = (SQLServerSink) nodeInfo;
Map<String, String> properties =
parseProperties(sqlServerSink.getProperties());
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
+
fieldTypeStrategyFactory.getInstance(sqlServerSink.getSinkType());
List<FieldInfo> fieldInfos =
parseSinkFieldInfos(sqlServerSink.getSinkFieldList(),
sqlServerSink.getSinkName(),
- FIELD_TYPE_MAPPING_STRATEGY);
+ fieldTypeMappingStrategy);
List<FieldRelation> fieldRelations =
parseSinkFields(sqlServerSink.getSinkFieldList(), constantFieldMap);
return new SqlServerLoadNode(
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/StarRocksProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/StarRocksProvider.java
index 3f0c7e2165..57d3656cac 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/StarRocksProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/StarRocksProvider.java
@@ -32,6 +32,7 @@ import
org.apache.inlong.sort.protocol.node.load.StarRocksLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@@ -42,6 +43,7 @@ import java.util.stream.Collectors;
* The Provider for creating StarRocks load nodes.
*/
@Slf4j
+@Service
public class StarRocksProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TDSQLPostgreSQLProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TDSQLPostgreSQLProvider.java
index 0f9f797199..22c1d9b09d 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TDSQLPostgreSQLProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TDSQLPostgreSQLProvider.java
@@ -28,12 +28,15 @@ import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.load.TDSQLPostgresLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelation;
+import org.springframework.stereotype.Service;
+
import java.util.List;
import java.util.Map;
/**
* The Provider for creating TDSQLPostgreSQL load nodes.
*/
+@Service
public class TDSQLPostgreSQLProvider implements LoadNodeProvider {
@Override
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TubeMqProvider.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TubeMqProvider.java
index 546363fe7e..6b80d4735e 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TubeMqProvider.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/node/provider/TubeMqProvider.java
@@ -29,6 +29,8 @@ import org.apache.inlong.sort.protocol.node.ExtractNode;
import org.apache.inlong.sort.protocol.node.extract.TubeMQExtractNode;
import org.apache.inlong.sort.protocol.node.format.Format;
+import org.springframework.stereotype.Service;
+
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@@ -37,6 +39,7 @@ import java.util.stream.Collectors;
/**
* The Provider for creating TubeMQ extract nodes.
*/
+@Service
public class TubeMqProvider implements ExtractNodeProvider {
@Override
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 7ada7dcab0..05a541770e 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
@@ -73,7 +73,7 @@ public class FieldInfoUtils {
boolean isMetaField = sinkField.getIsMetaField() == 1;
String fieldType = sinkField.getFieldType();
if (Objects.nonNull(fieldTypeMappingStrategy)) {
- fieldType =
fieldTypeMappingStrategy.getFieldTypeMapping(fieldType);
+ fieldType =
fieldTypeMappingStrategy.getSourceToSinkFieldTypeMapping(fieldType);
}
FieldInfo fieldInfo = getFieldInfo(sinkField.getFieldName(),
@@ -88,7 +88,7 @@ public class FieldInfoUtils {
boolean isMetaField = streamField.getIsMetaField() == 1;
String fieldType = streamField.getFieldType();
if (Objects.nonNull(fieldTypeMappingStrategy)) {
- fieldType =
fieldTypeMappingStrategy.getFieldTypeMapping(fieldType);
+ fieldType =
fieldTypeMappingStrategy.getSourceToSinkFieldTypeMapping(fieldType);
}
FieldInfo fieldInfo = getFieldInfo(streamField.getFieldName(),
fieldType,
@@ -106,7 +106,7 @@ public class FieldInfoUtils {
boolean isMetaField = streamField.getIsMetaField() == 1;
String fieldType = streamField.getFieldType();
if (Objects.nonNull(fieldTypeMappingStrategy)) {
- fieldType =
fieldTypeMappingStrategy.getFieldTypeMapping(fieldType);
+ fieldType =
fieldTypeMappingStrategy.getSourceToSinkFieldTypeMapping(fieldType);
}
FieldInfo fieldInfo = getFieldInfo(streamField.getFieldName(),
fieldType,
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamInfo.java
index 6c89e04a1d..0cf85ab044 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamInfo.java
@@ -142,6 +142,9 @@ public class InlongStreamInfo extends BaseInlongStream {
@ApiModelProperty("The multiple enable of sink")
private Boolean sinkMultipleEnable;
+ @ApiModelProperty("Whether to sync field")
+ private Boolean syncField = false;
+
@ApiModelProperty(value = "Whether to ignore the parse errors of field
value")
private Boolean ignoreParseError;
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamRequest.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamRequest.java
index d4a28bd83e..b60fc72bc2 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamRequest.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/stream/InlongStreamRequest.java
@@ -130,6 +130,9 @@ public class InlongStreamRequest extends BaseInlongStream {
@ApiModelProperty("The multiple enable of sink")
private Boolean sinkMultipleEnable;
+ @ApiModelProperty("Whether to sync field")
+ private Boolean syncField = false;
+
@ApiModelProperty(value = "The message body wrap type, including: RAW,
INLONG_MSG_V0, INLONG_MSG_V1, PB, etc")
private String wrapType;
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sort/SortFlinkConfigOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sort/SortFlinkConfigOperator.java
index c236e13c5e..9216ff8c27 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sort/SortFlinkConfigOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sort/SortFlinkConfigOperator.java
@@ -32,7 +32,6 @@ import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
import org.apache.inlong.manager.pojo.stream.StreamField;
import org.apache.inlong.manager.pojo.transform.TransformResponse;
import org.apache.inlong.manager.service.core.AuditService;
-import org.apache.inlong.manager.service.sink.StreamSinkService;
import org.apache.inlong.manager.service.source.StreamSourceService;
import org.apache.inlong.manager.service.transform.StreamTransformService;
import org.apache.inlong.sort.protocol.GroupInfo;
@@ -77,9 +76,9 @@ public class SortFlinkConfigOperator implements
SortConfigOperator {
@Autowired
private StreamTransformService transformService;
@Autowired
- private StreamSinkService sinkService;
- @Autowired
private AuditService auditService;
+ @Autowired
+ private NodeFactory nodeFactory;
@Override
public Boolean accept(List<String> sinkTypeList) {
@@ -249,13 +248,13 @@ public class SortFlinkConfigOperator implements
SortConfigOperator {
List<StreamSink> sinks, Map<String, StreamField> constantFieldMap)
{
List<Node> nodes = new ArrayList<>();
if (Objects.equals(sources.size(), sinks.size()) &&
Objects.equals(sources.size(), 1)) {
- return NodeFactory.addBuiltInField(sources.get(0), sinks.get(0),
transformResponses, constantFieldMap);
+ return nodeFactory.addBuiltInField(sources.get(0), sinks.get(0),
transformResponses, constantFieldMap);
}
List<TransformNode> transformNodes =
TransformNodeUtils.createTransformNodes(transformResponses,
constantFieldMap);
- nodes.addAll(NodeFactory.createExtractNodes(sources));
+ nodes.addAll(nodeFactory.createExtractNodes(sources));
nodes.addAll(transformNodes);
- nodes.addAll(NodeFactory.createLoadNodes(sinks, constantFieldMap));
+ nodes.addAll(nodeFactory.createLoadNodes(sinks, constantFieldMap));
return nodes;
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
index 506a8a5b23..9ace6bff7f 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
@@ -22,6 +22,8 @@ import
org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.enums.SinkStatus;
import org.apache.inlong.manager.common.exceptions.BusinessException;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeMappingStrategy;
+import
org.apache.inlong.manager.common.fieldtype.strategy.FieldTypeStrategyFactory;
import org.apache.inlong.manager.common.util.CommonBeanUtils;
import org.apache.inlong.manager.common.util.JsonUtils;
import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
@@ -35,6 +37,7 @@ import org.apache.inlong.manager.pojo.node.DataNodeInfo;
import org.apache.inlong.manager.pojo.sink.SinkField;
import org.apache.inlong.manager.pojo.sink.SinkRequest;
import org.apache.inlong.manager.pojo.sink.StreamSink;
+import org.apache.inlong.manager.pojo.stream.StreamField;
import org.apache.inlong.manager.service.node.DataNodeOperateHelper;
import com.github.pagehelper.Page;
@@ -47,6 +50,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Objects;
@@ -70,6 +74,8 @@ public abstract class AbstractSinkOperator implements
StreamSinkOperator {
protected InlongStreamEntityMapper inlongStreamEntityMapper;
@Autowired
protected SortConfigEntityMapper sortConfigEntityMapper;
+ @Autowired
+ protected FieldTypeStrategyFactory fieldTypeStrategyFactory;
/**
* Setting the parameters of the latest entity.
@@ -206,6 +212,31 @@ public abstract class AbstractSinkOperator implements
StreamSinkOperator {
LOGGER.debug("success to save sink fields");
}
+ @Override
+ public void syncField(SinkRequest request, List<StreamField> streamFields)
{
+ FieldTypeMappingStrategy fieldTypeMappingStrategy =
fieldTypeStrategyFactory.getInstance(request.getSinkType());
+ if (fieldTypeMappingStrategy == null) {
+ LOGGER.info("current sink type ={} not support sync field",
request.getSinkType());
+ return;
+ }
+ List<SinkField> sinkFields = request.getSinkFieldList();
+ if (sinkFields.size() >= streamFields.size()) {
+ return;
+ }
+ for (int i = sinkFields.size(); i < streamFields.size(); i++) {
+ StreamField streamField = streamFields.get(i);
+ SinkField sinkField = CommonBeanUtils.copyProperties(streamField,
SinkField::new);
+ sinkField.setSourceFieldName(streamField.getFieldName());
+ sinkField.setSourceFieldType(streamField.getFieldType());
+ sinkField.setFieldComment(streamField.getFieldComment());
+ sinkField.setFieldName(streamField.getFieldName());
+
sinkField.setFieldType(fieldTypeMappingStrategy.getStreamToSinkFieldTypeMapping(streamField.getFieldType())
+ .toLowerCase(Locale.ROOT));
+ sinkFields.add(sinkField);
+ }
+ updateFieldOpt(true, request);
+ }
+
@Override
public void deleteOpt(StreamSinkEntity entity, String operator) {
sortConfigEntityMapper.logicDeleteBySinkId(entity.getId());
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
index e0541cb22b..87e3924b38 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
@@ -25,6 +25,7 @@ import org.apache.inlong.manager.pojo.node.DataNodeInfo;
import org.apache.inlong.manager.pojo.sink.SinkField;
import org.apache.inlong.manager.pojo.sink.SinkRequest;
import org.apache.inlong.manager.pojo.sink.StreamSink;
+import org.apache.inlong.manager.pojo.stream.StreamField;
import com.github.pagehelper.Page;
@@ -103,6 +104,14 @@ public interface StreamSinkOperator {
*/
void saveFieldOpt(SinkRequest request);
+ /**
+ * Sync the sink fields.
+ *
+ * @param request sink request info needs to save
+ * @param streamFields stream field list
+ */
+ void syncField(SinkRequest request, List<StreamField> streamFields);
+
/**
* Delete the sink info.
*
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
index d11c37b2a1..281858abb2 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamServiceImpl.java
@@ -62,6 +62,8 @@ import
org.apache.inlong.manager.service.group.InlongGroupOperator;
import org.apache.inlong.manager.service.group.InlongGroupOperatorFactory;
import org.apache.inlong.manager.service.resource.queue.QueueResourceOperator;
import
org.apache.inlong.manager.service.resource.queue.QueueResourceOperatorFactory;
+import org.apache.inlong.manager.service.sink.SinkOperatorFactory;
+import org.apache.inlong.manager.service.sink.StreamSinkOperator;
import org.apache.inlong.manager.service.sink.StreamSinkService;
import org.apache.inlong.manager.service.source.StreamSourceService;
@@ -144,6 +146,9 @@ public class InlongStreamServiceImpl implements
InlongStreamService {
@Autowired
@Lazy
private InlongGroupOperatorFactory groupOperatorFactory;
+ @Autowired
+ @Lazy
+ private SinkOperatorFactory sinkOperatorFactory;
@Transactional(rollbackFor = Throwable.class)
@Override
@@ -600,7 +605,14 @@ public class InlongStreamServiceImpl implements
InlongStreamService {
// update stream extension infos
List<InlongStreamExtInfo> extList = request.getExtList();
saveOrUpdateExt(groupId, streamId, extList);
-
+ if (request.getSyncField()) {
+ LOGGER.info("test begin sync field={}", request);
+ List<StreamSinkEntity> sinkEntityList =
sinkMapper.selectByRelatedId(groupId, streamId);
+ for (StreamSinkEntity sinkEntity : sinkEntityList) {
+ StreamSinkOperator sinkOperator =
sinkOperatorFactory.getInstance(sinkEntity.getSinkType());
+
sinkOperator.syncField(sinkOperator.getFromEntity(sinkEntity).genSinkRequest(),
request.getFieldList());
+ }
+ }
LOGGER.info("success to update inlong stream without check for
groupId={} streamId={}", groupId, streamId);
return true;
}