This is an automated email from the ASF dual-hosted git repository.
zirui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new ca4846d9b [INLONG-3828][Sort] Enhanced the node filter capabilities
(#4159)
ca4846d9b is described below
commit ca4846d9b5ff7b67affa824eb7680d6b4492e21b
Author: yunqingmoswu <[email protected]>
AuthorDate: Wed May 11 16:54:02 2022 +0800
[INLONG-3828][Sort] Enhanced the node filter capabilities (#4159)
---
.../service/sort/util/FilterFunctionUtils.java | 33 ++++
.../manager/service/sort/util/LoadNodeUtils.java | 2 +
.../service/sort/util/TransformNodeUtils.java | 2 +
.../inlong/sort/protocol/enums/FilterStrategy.java | 32 ++++
.../apache/inlong/sort/protocol/node/LoadNode.java | 9 +-
.../sort/protocol/node/load/HiveLoadNode.java | 4 +-
.../sort/protocol/node/load/KafkaLoadNode.java | 4 +-
.../sort/protocol/node/transform/DistinctNode.java | 4 +-
.../protocol/node/transform/TransformNode.java | 11 +-
.../apache/inlong/sort/protocol/GroupInfoTest.java | 2 +-
.../inlong/sort/protocol/StreamInfoTest.java | 4 +-
.../sort/protocol/node/load/HiveLoadNodeTest.java | 2 +-
.../sort/protocol/node/load/KafkaLoadNodeTest.java | 2 +-
.../protocol/node/transform/DistinctNodeTest.java | 2 +-
.../flink/parser/impl/FlinkSqlParser.java | 21 ++-
.../singletenant/flink/parser/AllMigrateTest.java | 48 +++---
.../flink/parser/DistinctNodeSqlParseTest.java | 12 +-
.../singletenant/flink/parser/FilterParseTest.java | 172 +++++++++++++++++++++
.../flink/parser/FlinkSqlParserTest.java | 4 +-
.../flink/parser/FullOuterJoinSqlParseTest.java | 7 +-
.../parser/InnerJoinRelationShipSqlParseTest.java | 9 +-
.../flink/parser/LeftOuterJoinSqlParseTest.java | 8 +-
.../flink/parser/MetaFieldSyncTest.java | 4 +-
.../flink/parser/RightOuterJoinSqlParseTest.java | 7 +-
24 files changed, 337 insertions(+), 68 deletions(-)
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/FilterFunctionUtils.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/FilterFunctionUtils.java
index 16f63bd69..b53f31917 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/FilterFunctionUtils.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/FilterFunctionUtils.java
@@ -31,6 +31,7 @@ import
org.apache.inlong.manager.common.pojo.transform.filter.FilterDefinition.T
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.common.util.StreamParseUtils;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.transformation.ConstantParam;
import org.apache.inlong.sort.protocol.transformation.FilterFunction;
import org.apache.inlong.sort.protocol.transformation.FunctionParam;
@@ -91,6 +92,38 @@ public class FilterFunctionUtils {
return filterFunctions;
}
+ /**
+ * Parse filter strategy from TransformResponse and convert to the filter
strategy of sort protocol
+ *
+ * @param transformResponse The transform response that may contains
filter operation
+ * @return The filter strategy, see {@link FilterStrategy}
+ */
+ public static FilterStrategy parseFilterStrategy(TransformResponse
transformResponse) {
+ TransformType transformType =
TransformType.forType(transformResponse.getTransformType());
+ TransformDefinition transformDefinition =
StreamParseUtils.parseTransformDefinition(
+ transformResponse.getTransformDefinition(), transformType);
+ switch (transformType) {
+ case FILTER:
+ FilterDefinition filterDefinition = (FilterDefinition)
transformDefinition;
+ switch (filterDefinition.getFilterStrategy()) {
+ case REMOVE:
+ return FilterStrategy.REMOVE;
+ case RETAIN:
+ return FilterStrategy.RETAIN;
+ default:
+ return FilterStrategy.RETAIN;
+ }
+ case DE_DUPLICATION:
+ case SPLITTER:
+ case JOINER:
+ case STRING_REPLACER:
+ return null;
+ default:
+ throw new UnsupportedOperationException(
+ String.format("Unsupported transformType=%s for
Inlong", transformType));
+ }
+ }
+
private static FilterFunction createFilterFunction(FilterRule filterRule,
String transformName) {
StreamField streamField = filterRule.getSourceField();
String fieldType = streamField.getFieldType().name();
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/LoadNodeUtils.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/LoadNodeUtils.java
index c4cae249c..8072f5d94 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/LoadNodeUtils.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/LoadNodeUtils.java
@@ -109,6 +109,7 @@ public class LoadNodeUtils {
fieldInfos,
fieldRelationShips,
Lists.newArrayList(),
+ null,
topicName,
bootstrapServers,
format,
@@ -145,6 +146,7 @@ public class LoadNodeUtils {
fieldRelationShips,
Lists.newArrayList(),
null,
+ null,
properties,
null,
database,
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/TransformNodeUtils.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/TransformNodeUtils.java
index 5221a300c..5d44a62dd 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/TransformNodeUtils.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/TransformNodeUtils.java
@@ -95,6 +95,7 @@ public class TransformNodeUtils {
transformNode.getFields(),
transformNode.getFieldRelationShips(),
transformNode.getFilters(),
+ transformNode.getFilterStrategy(),
distinctFields,
orderField,
orderDirection);
@@ -117,6 +118,7 @@ public class TransformNodeUtils {
transformNode.setFieldRelationShips(FieldRelationShipUtils.createFieldRelationShips(transformResponse));
transformNode.setFilters(
FilterFunctionUtils.createFilterFunctions(transformResponse));
+
transformNode.setFilterStrategy(FilterFunctionUtils.parseFilterStrategy(transformResponse));
return transformNode;
}
}
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/enums/FilterStrategy.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/enums/FilterStrategy.java
new file mode 100644
index 000000000..fd2a45c38
--- /dev/null
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/enums/FilterStrategy.java
@@ -0,0 +1,32 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.protocol.enums;
+
+/**
+ * Filter strategy,it defines the strategy of filter
+ */
+public enum FilterStrategy {
+ /**
+ * The RETAIN strategy that is retain the data that satisfies the filter
+ */
+ RETAIN,
+ /**
+ * The REMOVE strategy that is remove the data that satisfies the filter
+ */
+ REMOVE
+}
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
index 332700a33..a87c5b7a4 100644
---
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/LoadNode.java
@@ -27,13 +27,13 @@ import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonPro
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonSubTypes;
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.node.load.HiveLoadNode;
import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
import org.apache.inlong.sort.protocol.transformation.FilterFunction;
import javax.annotation.Nullable;
-import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@@ -64,7 +64,10 @@ public abstract class LoadNode implements Node {
private Integer sinkParallelism;
@JsonProperty("filters")
@JsonInclude(Include.NON_NULL)
- private List<FilterFunction> filters = new ArrayList<>();
+ private List<FilterFunction> filters;
+ @JsonProperty("filterStrategy")
+ @JsonInclude(Include.NON_NULL)
+ private FilterStrategy filterStrategy;
@Nullable
@JsonInclude(Include.NON_NULL)
@JsonProperty("properties")
@@ -76,6 +79,7 @@ public abstract class LoadNode implements Node {
@JsonProperty("fields") List<FieldInfo> fields,
@JsonProperty("fieldRelationShips") List<FieldRelationShip>
fieldRelationShips,
@JsonProperty("filters") List<FilterFunction> filters,
+ @JsonProperty("filterStrategy") FilterStrategy filterStrategy,
@Nullable @JsonProperty("sinkParallelism") Integer sinkParallelism,
@JsonProperty("properties") Map<String, String> properties) {
this.id = Preconditions.checkNotNull(id, "id is null");
@@ -86,6 +90,7 @@ public abstract class LoadNode implements Node {
"fieldRelationShips is null");
Preconditions.checkState(!fieldRelationShips.isEmpty(),
"fieldRelationShips is empty");
this.filters = filters;
+ this.filterStrategy = filterStrategy;
this.sinkParallelism = sinkParallelism;
this.properties = properties;
}
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNode.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNode.java
index d4f114f89..8a9baae44 100644
---
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNode.java
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNode.java
@@ -27,6 +27,7 @@ import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTyp
import org.apache.inlong.sort.formats.common.LocalZonedTimestampFormatInfo;
import org.apache.inlong.sort.formats.common.TimestampFormatInfo;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
import org.apache.inlong.sort.protocol.transformation.FilterFunction;
@@ -79,6 +80,7 @@ public class HiveLoadNode extends LoadNode implements
Serializable {
@JsonProperty("fields") List<FieldInfo> fields,
@JsonProperty("fieldRelationShips") List<FieldRelationShip>
fieldRelationShips,
@JsonProperty("filters") List<FilterFunction> filters,
+ @JsonProperty("filterStrategy") FilterStrategy filterStrategy,
@JsonProperty("sinkParallelism") Integer sinkParallelism,
@JsonProperty("properties") Map<String, String> properties,
@JsonProperty("catalogName") String catalogName,
@@ -88,7 +90,7 @@ public class HiveLoadNode extends LoadNode implements
Serializable {
@JsonProperty("hiveVersion") String hiveVersion,
@JsonProperty("hadoopConfDir") String hadoopConfDir,
@JsonProperty("parFields") List<FieldInfo> partitionFields) {
- super(id, name, fields, fieldRelationShips, filters, sinkParallelism,
properties);
+ super(id, name, fields, fieldRelationShips, filters, filterStrategy,
sinkParallelism, properties);
this.database = Preconditions.checkNotNull(database, "database of hive
is null");
this.tableName = Preconditions.checkNotNull(tableName, "table of hive
is null");
this.hiveVersion = Preconditions.checkNotNull(hiveVersion, "version of
hive is null");
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNode.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNode.java
index 7ef0d8cb3..9c3c40585 100644
---
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNode.java
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNode.java
@@ -26,6 +26,7 @@ import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCre
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeName;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.format.AvroFormat;
import org.apache.inlong.sort.protocol.node.format.CanalJsonFormat;
@@ -70,13 +71,14 @@ public class KafkaLoadNode extends LoadNode implements
Serializable {
@JsonProperty("fields") List<FieldInfo> fields,
@JsonProperty("fieldRelationShips") List<FieldRelationShip>
fieldRelationShips,
@JsonProperty("filters") List<FilterFunction> filters,
+ @JsonProperty("filterStrategy") FilterStrategy filterStrategy,
@Nonnull @JsonProperty("topic") String topic,
@Nonnull @JsonProperty("bootstrapServers") String bootstrapServers,
@Nonnull @JsonProperty("format") Format format,
@Nullable @JsonProperty("sinkParallelism") Integer sinkParallelism,
@JsonProperty("properties") Map<String, String> properties,
@JsonProperty("primaryKey") String primaryKey) {
- super(id, name, fields, fieldRelationShips, filters, sinkParallelism,
properties);
+ super(id, name, fields, fieldRelationShips, filters, filterStrategy,
sinkParallelism, properties);
this.topic = Preconditions.checkNotNull(topic, "topic is null");
this.bootstrapServers = Preconditions.checkNotNull(bootstrapServers,
"bootstrapServers is null");
this.format = Preconditions.checkNotNull(format, "format is null");
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/DistinctNode.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/DistinctNode.java
index 66acdccdd..f1bcfa344 100644
---
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/DistinctNode.java
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/DistinctNode.java
@@ -25,6 +25,7 @@ import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCre
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeName;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
import org.apache.inlong.sort.protocol.transformation.FilterFunction;
import org.apache.inlong.sort.protocol.transformation.OrderDirection;
@@ -75,10 +76,11 @@ public class DistinctNode extends TransformNode {
@JsonProperty("fields") List<FieldInfo> fields,
@JsonProperty("fieldRelationShips") List<FieldRelationShip>
fieldRelationShips,
@JsonProperty("filters") List<FilterFunction> filters,
+ @JsonProperty("filterStrategy") FilterStrategy filterStrategy,
@JsonProperty("distinctFields") List<FieldInfo> distinctFields,
@JsonProperty("orderField") FieldInfo orderField,
@JsonProperty("orderDirection") OrderDirection orderDirection) {
- super(id, name, fields, fieldRelationShips, filters);
+ super(id, name, fields, fieldRelationShips, filters, filterStrategy);
this.distinctFields = Preconditions.checkNotNull(distinctFields,
"distinctFields is null");
Preconditions.checkState(!distinctFields.isEmpty(), "distinct fields
is empty");
this.orderField = Preconditions.checkNotNull(orderField, "orderField
is null");
diff --git
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/TransformNode.java
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/TransformNode.java
index da619dc26..c037b0e89 100644
---
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/TransformNode.java
+++
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/node/transform/TransformNode.java
@@ -28,12 +28,12 @@ import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonPro
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonSubTypes;
import
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.node.Node;
import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
import org.apache.inlong.sort.protocol.transformation.FilterFunction;
import java.io.Serializable;
-import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@@ -61,14 +61,18 @@ public class TransformNode implements Node, Serializable {
private List<FieldRelationShip> fieldRelationShips;
@JsonProperty("filters")
@JsonInclude(Include.NON_NULL)
- private List<FilterFunction> filters = new ArrayList<>();
+ private List<FilterFunction> filters;
+ @JsonProperty("filterStrategy")
+ @JsonInclude(Include.NON_NULL)
+ private FilterStrategy filterStrategy;
@JsonCreator
public TransformNode(@JsonProperty("id") String id,
@JsonProperty("name") String name,
@JsonProperty("fields") List<FieldInfo> fields,
@JsonProperty("fieldRelationShips") List<FieldRelationShip>
fieldRelationShips,
- @JsonProperty("filters") List<FilterFunction> filters) {
+ @JsonProperty("filters") List<FilterFunction> filters,
+ @JsonProperty("filterStrategy") FilterStrategy filterStrategy) {
this.id = Preconditions.checkNotNull(id, "id is null");
this.name = name;
this.fields = Preconditions.checkNotNull(fields, "fields is null");
@@ -77,6 +81,7 @@ public class TransformNode implements Node, Serializable {
"fieldRelationShips is null");
Preconditions.checkState(!fieldRelationShips.isEmpty(),
"fieldRelationShips is empty");
this.filters = filters;
+ this.filterStrategy = filterStrategy;
}
@JsonIgnore
diff --git
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/GroupInfoTest.java
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/GroupInfoTest.java
index 452cb20b6..38684674f 100644
---
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/GroupInfoTest.java
+++
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/GroupInfoTest.java
@@ -75,7 +75,7 @@ public class GroupInfoTest extends
SerializeBaseTest<GroupInfo> {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
+ return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
null,
"topic", "localhost:9092", new JsonFormat(),
1, null, "id");
}
diff --git
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/StreamInfoTest.java
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/StreamInfoTest.java
index 4df4f30ed..f9f8cf3e7 100644
---
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/StreamInfoTest.java
+++
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/StreamInfoTest.java
@@ -76,7 +76,7 @@ public class StreamInfoTest extends
SerializeBaseTest<StreamInfo> {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
+ return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
null,
"topic", "localhost:9092", new JsonFormat(),
1, null, "id");
}
@@ -97,7 +97,7 @@ public class StreamInfoTest extends
SerializeBaseTest<StreamInfo> {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new HiveLoadNode("2", "hive_output", fields, relations, null,
+ return new HiveLoadNode("2", "hive_output", fields, relations, null,
null,
1, null, "myHive", "default", "test", "/opt/hive-conf",
"3.1.2",
null, Arrays.asList(new FieldInfo("day", new
LongFormatInfo())));
}
diff --git
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNodeTest.java
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNodeTest.java
index e11c50ee1..012e3dd05 100644
---
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNodeTest.java
+++
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/HiveLoadNodeTest.java
@@ -36,7 +36,7 @@ public class HiveLoadNodeTest extends
SerializeBaseTest<HiveLoadNode> {
return new HiveLoadNode("1", "test_hive_node",
Arrays.asList(new FieldInfo("field", new StringFormatInfo())),
Arrays.asList(new FieldRelationShip(new FieldInfo("field", new
StringFormatInfo()),
- new FieldInfo("field", new StringFormatInfo()))), null,
+ new FieldInfo("field", new StringFormatInfo()))),
null, null,
1, new HashMap<>(), "myHive", "default",
"test", "/opt/hive-conf", "3.1.2",
null, Arrays.asList(new FieldInfo("day", new
LongFormatInfo())));
diff --git
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNodeTest.java
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNodeTest.java
index 085d9dbca..0ff7316c7 100644
---
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNodeTest.java
+++
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/load/KafkaLoadNodeTest.java
@@ -36,7 +36,7 @@ public class KafkaLoadNodeTest extends
SerializeBaseTest<KafkaLoadNode> {
return new KafkaLoadNode("1", null,
Arrays.asList(new FieldInfo("field", new StringFormatInfo())),
Arrays.asList(new FieldRelationShip(new FieldInfo("field", new
StringFormatInfo()),
- new FieldInfo("field", new StringFormatInfo()))), null,
+ new FieldInfo("field", new StringFormatInfo()))),
null, null,
"topic", "localhost:9092", new CanalJsonFormat(),
1, new TreeMap<>(), null);
}
diff --git
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/transform/DistinctNodeTest.java
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/transform/DistinctNodeTest.java
index 678526f9c..7f308271c 100644
---
a/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/transform/DistinctNodeTest.java
+++
b/inlong-sort/sort-common/src/test/java/org/apache/inlong/sort/protocol/node/transform/DistinctNodeTest.java
@@ -48,7 +48,7 @@ public class DistinctNodeTest extends
SerializeBaseTest<DistinctNode> {
new FieldRelationShip(new FieldInfo("ts", new
StringFormatInfo()),
new FieldInfo("ts", new StringFormatInfo()))
),
- null,
+ null, null,
Arrays.asList(new FieldInfo("f1", new StringFormatInfo()),
new FieldInfo("f2", new StringFormatInfo())),
new FieldInfo("ts", new StringFormatInfo()),
OrderDirection.ASC);
diff --git
a/inlong-sort/sort-single-tenant/src/main/java/org/apache/inlong/sort/singletenant/flink/parser/impl/FlinkSqlParser.java
b/inlong-sort/sort-single-tenant/src/main/java/org/apache/inlong/sort/singletenant/flink/parser/impl/FlinkSqlParser.java
index 9247498fe..07a7822ae 100644
---
a/inlong-sort/sort-single-tenant/src/main/java/org/apache/inlong/sort/singletenant/flink/parser/impl/FlinkSqlParser.java
+++
b/inlong-sort/sort-single-tenant/src/main/java/org/apache/inlong/sort/singletenant/flink/parser/impl/FlinkSqlParser.java
@@ -26,6 +26,7 @@ import
org.apache.inlong.sort.protocol.BuiltInFieldInfo.BuiltInField;
import org.apache.inlong.sort.protocol.FieldInfo;
import org.apache.inlong.sort.protocol.GroupInfo;
import org.apache.inlong.sort.protocol.StreamInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.node.ExtractNode;
import org.apache.inlong.sort.protocol.node.LoadNode;
import org.apache.inlong.sort.protocol.node.Node;
@@ -54,6 +55,7 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.stream.Collectors;
/**
* Flink sql parse handler
@@ -336,7 +338,7 @@ public class FlinkSqlParser implements Parser {
// Fill out the tablename alias for param
fillOutTableNameAlias(new ArrayList<>(node.getFilters()),
tableNameAliasMap);
// Parse filter fields to generate filter sql like 'WHERE 1=1...'
- parseFilterFields(node.getFilters(), sb);
+ parseFilterFields(node.getFilterStrategy(), node.getFilters(), sb);
}
if (node instanceof DistinctNode) {
// Generate distinct filter sql like 'WHERE row_num = 1'
@@ -428,7 +430,7 @@ public class FlinkSqlParser implements Parser {
genDistinctSql((DistinctNode) node, sb);
}
sb.append("\n FROM
`").append(nodeMap.get(relation.getInputs().get(0)).genTableName()).append("`
");
- parseFilterFields(node.getFilters(), sb);
+ parseFilterFields(node.getFilterStrategy(), node.getFilters(), sb);
if (node instanceof DistinctNode) {
sb = genDistinctFilterSql(node.getFields(), sb);
}
@@ -438,14 +440,19 @@ public class FlinkSqlParser implements Parser {
/**
* Parse filter fields to generate filter sql like 'where 1=1...'
*
+ * @param filterStrategy The filter strategy default[RETAIN], it decide
whether to retain or remove
* @param filters The filter functions
* @param sb Container for storing sql
*/
- private void parseFilterFields(List<FilterFunction> filters, StringBuilder
sb) {
+ private void parseFilterFields(FilterStrategy filterStrategy,
List<FilterFunction> filters, StringBuilder sb) {
if (filters != null && !filters.isEmpty()) {
- sb.append("\n WHERE");
- for (FilterFunction filter : filters) {
- sb.append(" ").append(filter.format());
+ sb.append("\n WHERE ");
+ String subSql = StringUtils
+
.join(filters.stream().map(FunctionParam::format).collect(Collectors.toList()),
" ");
+ if (filterStrategy == FilterStrategy.REMOVE) {
+ sb.append("not (").append(subSql).append(")");
+ } else {
+ sb.append(subSql);
}
}
}
@@ -492,7 +499,7 @@ public class FlinkSqlParser implements Parser {
});
parseFieldRelations(loadNode.getFields(), fieldRelationMap, sb);
sb.append("\n FROM `").append(inputNode.genTableName()).append("`");
- parseFilterFields(loadNode.getFilters(), sb);
+ parseFilterFields(loadNode.getFilterStrategy(), loadNode.getFilters(),
sb);
return sb.toString();
}
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/AllMigrateTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/AllMigrateTest.java
index 7347b2078..36980f13f 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/AllMigrateTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/AllMigrateTest.java
@@ -17,14 +17,6 @@
package org.apache.inlong.sort.singletenant.flink.parser;
-import static
org.apache.inlong.sort.protocol.BuiltInFieldInfo.BuiltInField.MYSQL_METADATA_DATA;
-
-import java.util.Arrays;
-import java.util.Collections;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.stream.Collectors;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
@@ -44,32 +36,40 @@ import
org.apache.inlong.sort.singletenant.flink.parser.result.FlinkSqlParseResu
import org.junit.Assert;
import org.junit.Test;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+import static
org.apache.inlong.sort.protocol.BuiltInFieldInfo.BuiltInField.MYSQL_METADATA_DATA;
+
public class AllMigrateTest {
private MySqlExtractNode buildAllMigrateExtractNode() {
List<FieldInfo> fields = Arrays.asList(
- new BuiltInFieldInfo("data", new StringFormatInfo(),
MYSQL_METADATA_DATA));
+ new BuiltInFieldInfo("data", new StringFormatInfo(),
MYSQL_METADATA_DATA));
Map<String, String> option = new HashMap<>();
option.put("append-mode", "true");
option.put("migrate-all", "true");
MySqlExtractNode node = new MySqlExtractNode("1", "mysql_input",
fields,
- null, option, null,
- Arrays.asList("[\\s\\S]*.*"), "localhost", "root", "password",
- "[\\s\\S]*.*", null, null, false, null);
+ null, option, null,
+ Arrays.asList("[\\s\\S]*.*"), "localhost", "root", "password",
+ "[\\s\\S]*.*", null, null, false, null);
return node;
}
private KafkaLoadNode buildAllMigrateKafkaNode() {
List<FieldInfo> fields = Arrays.asList(new FieldInfo("data", new
StringFormatInfo()));
List<FieldRelationShip> relations = Arrays
- .asList(new FieldRelationShip(new FieldInfo("data", new
StringFormatInfo()),
- new FieldInfo("data", new StringFormatInfo())));
+ .asList(new FieldRelationShip(new FieldInfo("data", new
StringFormatInfo()),
+ new FieldInfo("data", new StringFormatInfo())));
CsvFormat csvFormat = new CsvFormat();
csvFormat.setDisableQuoteCharacter(true);
- return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
- "topic", "localhost:9092",
- csvFormat, null,
- null, null);
+ return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
null,
+ "topic", "localhost:9092",
+ csvFormat, null,
+ null, null);
}
private NodeRelationShip buildNodeRelation(List<Node> inputs, List<Node>
outputs) {
@@ -86,10 +86,10 @@ public class AllMigrateTest {
@Test
public void testAllMigrate() throws Exception {
EnvironmentSettings settings = EnvironmentSettings
- .newInstance()
- .useBlinkPlanner()
- .inStreamingMode()
- .build();
+ .newInstance()
+ .useBlinkPlanner()
+ .inStreamingMode()
+ .build();
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
env.enableCheckpointing(10000);
@@ -97,8 +97,8 @@ public class AllMigrateTest {
Node inputNode = buildAllMigrateExtractNode();
Node outputNode = buildAllMigrateKafkaNode();
StreamInfo streamInfo = new StreamInfo("1", Arrays.asList(inputNode,
outputNode),
-
Collections.singletonList(buildNodeRelation(Collections.singletonList(inputNode),
- Collections.singletonList(outputNode))));
+
Collections.singletonList(buildNodeRelation(Collections.singletonList(inputNode),
+ Collections.singletonList(outputNode))));
GroupInfo groupInfo = new GroupInfo("1",
Collections.singletonList(streamInfo));
FlinkSqlParser parser = FlinkSqlParser.getInstance(tableEnv,
groupInfo);
FlinkSqlParseResult result = parser.parse();
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/DistinctNodeSqlParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/DistinctNodeSqlParseTest.java
index f18bd365d..61a5134b8 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/DistinctNodeSqlParseTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/DistinctNodeSqlParseTest.java
@@ -115,7 +115,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode("3", "kafka_output", fields, relations,
+ return new KafkaLoadNode("3", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, null);
@@ -137,7 +137,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode("3", "kafka_output", fields, relations,
+ return new KafkaLoadNode("3", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -159,7 +159,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode("3", "kafka_output", fields, relations,
+ return new KafkaLoadNode("3", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -183,7 +183,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
),
- null,
+ null, null,
Collections.singletonList(new FieldInfo("name", new
StringFormatInfo())),
new FieldInfo("proctime", new TimestampFormatInfo()),
OrderDirection.ASC);
@@ -207,7 +207,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
),
- null,
+ null, null,
Collections.singletonList(new FieldInfo("name", new
StringFormatInfo())),
new FieldInfo("ts", new TimestampFormatInfo()),
OrderDirection.ASC);
@@ -231,7 +231,7 @@ public class DistinctNodeSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
),
- null,
+ null, null,
Collections.singletonList(new FieldInfo("name", new
StringFormatInfo())),
new FieldInfo("ts", new TimestampFormatInfo()),
OrderDirection.ASC);
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FilterParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FilterParseTest.java
new file mode 100644
index 000000000..c3e14a1ce
--- /dev/null
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FilterParseTest.java
@@ -0,0 +1,172 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.singletenant.flink.parser;
+
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+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.sort.formats.common.FloatFormatInfo;
+import org.apache.inlong.sort.formats.common.IntFormatInfo;
+import org.apache.inlong.sort.formats.common.LongFormatInfo;
+import org.apache.inlong.sort.formats.common.StringFormatInfo;
+import org.apache.inlong.sort.formats.common.TimestampFormatInfo;
+import org.apache.inlong.sort.protocol.FieldInfo;
+import org.apache.inlong.sort.protocol.GroupInfo;
+import org.apache.inlong.sort.protocol.StreamInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
+import org.apache.inlong.sort.protocol.node.Node;
+import org.apache.inlong.sort.protocol.node.extract.MySqlExtractNode;
+import org.apache.inlong.sort.protocol.node.format.CanalJsonFormat;
+import org.apache.inlong.sort.protocol.node.load.KafkaLoadNode;
+import org.apache.inlong.sort.protocol.transformation.ConstantParam;
+import org.apache.inlong.sort.protocol.transformation.FieldRelationShip;
+import org.apache.inlong.sort.protocol.transformation.FilterFunction;
+import
org.apache.inlong.sort.protocol.transformation.function.SingleValueFilterFunction;
+import org.apache.inlong.sort.protocol.transformation.operator.AndOperator;
+import org.apache.inlong.sort.protocol.transformation.operator.EmptyOperator;
+import
org.apache.inlong.sort.protocol.transformation.operator.LessThanOperator;
+import
org.apache.inlong.sort.protocol.transformation.operator.MoreThanOrEqualOperator;
+import
org.apache.inlong.sort.protocol.transformation.relation.NodeRelationShip;
+import org.apache.inlong.sort.singletenant.flink.parser.impl.FlinkSqlParser;
+import
org.apache.inlong.sort.singletenant.flink.parser.result.FlinkSqlParseResult;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.stream.Collectors;
+
+/**
+ * Test for filter parse
+ */
+public class FilterParseTest extends AbstractTestBase {
+
+ private Node buildMySQLExtractNode() {
+ List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new
LongFormatInfo()),
+ new FieldInfo("name", new StringFormatInfo()),
+ new FieldInfo("age", new IntFormatInfo()),
+ new FieldInfo("salary", new FloatFormatInfo()),
+ new FieldInfo("ts", new TimestampFormatInfo())
+ );
+ return new MySqlExtractNode("1", "mysql_input", fields, null, null,
+ "id", Collections.singletonList("mysql_table"),
+ "localhost", "inlong", "inlong",
+ "inlong", null, null, null, null);
+ }
+
+ private Node buildKafkaLoadNode(FilterStrategy filterStrategy) {
+ List<FieldInfo> fields = Arrays.asList(new FieldInfo("id", new
LongFormatInfo()),
+ new FieldInfo("name", new StringFormatInfo()),
+ new FieldInfo("age", new IntFormatInfo()),
+ new FieldInfo("salary", new FloatFormatInfo()),
+ new FieldInfo("ts", new TimestampFormatInfo())
+ );
+ List<FieldRelationShip> relations = Arrays
+ .asList(new FieldRelationShip(new FieldInfo("id", new
LongFormatInfo()),
+ new FieldInfo("id", new LongFormatInfo())),
+ new FieldRelationShip(new FieldInfo("name", new
StringFormatInfo()),
+ new FieldInfo("name", new StringFormatInfo())),
+ new FieldRelationShip(new FieldInfo("age", new
IntFormatInfo()),
+ new FieldInfo("age", new IntFormatInfo())),
+ new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
+ new FieldInfo("ts", new TimestampFormatInfo()))
+ );
+ List<FilterFunction> filters = Arrays.asList(
+ new SingleValueFilterFunction(EmptyOperator.getInstance(),
+ new FieldInfo("age", new IntFormatInfo()),
+ LessThanOperator.getInstance(), new ConstantParam(25)),
+ new SingleValueFilterFunction(AndOperator.getInstance(),
+ new FieldInfo("age", new IntFormatInfo()),
+ MoreThanOrEqualOperator.getInstance(), new
ConstantParam(18))
+ );
+ return new KafkaLoadNode("2", "kafka_output", fields, relations,
filters,
+ filterStrategy, "topic1", "localhost:9092",
+ new CanalJsonFormat(), null,
+ null, "id");
+ }
+
+ public NodeRelationShip buildNodeRelation(List<Node> inputs, List<Node>
outputs) {
+ List<String> inputIds =
inputs.stream().map(Node::getId).collect(Collectors.toList());
+ List<String> outputIds =
outputs.stream().map(Node::getId).collect(Collectors.toList());
+ return new NodeRelationShip(inputIds, outputIds);
+ }
+
+ /**
+ * Test filter with the strategy {@link FilterStrategy#RETAIN}
+ *
+ * @throws Exception The exception may throws when execute the case
+ */
+ @Test
+ public void testFilterWithRetainParse() 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 mysqlInputNode = buildMySQLExtractNode();
+ Node kafkaOutputNode = buildKafkaLoadNode(FilterStrategy.RETAIN);
+ StreamInfo streamInfo = new StreamInfo("1",
+ Arrays.asList(mysqlInputNode, kafkaOutputNode),
+ Collections.singletonList(
+
buildNodeRelation(Collections.singletonList(mysqlInputNode),
+ Collections.singletonList(kafkaOutputNode))
+ )
+ );
+ GroupInfo groupInfo = new GroupInfo("1",
Collections.singletonList(streamInfo));
+ FlinkSqlParser parser = FlinkSqlParser.getInstance(tableEnv,
groupInfo);
+ FlinkSqlParseResult result = parser.parse();
+ Assert.assertTrue(result.tryExecute());
+ }
+
+ /**
+ * Test filter with the strategy {@link FilterStrategy#REMOVE}
+ *
+ * @throws Exception The exception may throws when execute the case
+ */
+ @Test
+ public void testFilterWithRemoveParse() 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 mysqlInputNode = buildMySQLExtractNode();
+ Node kafkaOutputNode = buildKafkaLoadNode(FilterStrategy.REMOVE);
+ StreamInfo streamInfo = new StreamInfo("1",
+ Arrays.asList(mysqlInputNode, kafkaOutputNode),
+ Collections.singletonList(
+
buildNodeRelation(Collections.singletonList(mysqlInputNode),
+ Collections.singletonList(kafkaOutputNode))
+ )
+ );
+ GroupInfo groupInfo = new GroupInfo("1",
Collections.singletonList(streamInfo));
+ FlinkSqlParser parser = FlinkSqlParser.getInstance(tableEnv,
groupInfo);
+ FlinkSqlParseResult result = parser.parse();
+ Assert.assertTrue(result.tryExecute());
+ }
+}
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FlinkSqlParserTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FlinkSqlParserTest.java
index c238b7450..3942e5fae 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FlinkSqlParserTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FlinkSqlParserTest.java
@@ -105,7 +105,7 @@ public class FlinkSqlParserTest extends AbstractTestBase {
new FieldRelationShip(new FieldInfo("ts", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
);
- return new KafkaLoadNode(id, "kafka_output", fields, relations, null,
+ return new KafkaLoadNode(id, "kafka_output", fields, relations, null,
null,
"workerJson", "localhost:9092",
new JsonFormat(), null,
null, null);
@@ -134,7 +134,7 @@ public class FlinkSqlParserTest extends AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo()))
);
return new HiveLoadNode(id, "hive_output",
- fields, relations, null, 1,
+ fields, relations, null, null, 1,
null, "myCatalog", "default", "work2",
"/opt/hive/conf", "3.1.2",
null, null);
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FullOuterJoinSqlParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FullOuterJoinSqlParseTest.java
index aa08bce11..cdf785ef6 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FullOuterJoinSqlParseTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/FullOuterJoinSqlParseTest.java
@@ -130,7 +130,7 @@ public class FullOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("5", "kafka_output", fields, relations,
+ return new KafkaLoadNode("5", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -159,7 +159,7 @@ public class FullOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("ts", "3", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
- ), null);
+ ), null, null);
}
/**
@@ -195,7 +195,8 @@ public class FullOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new TimestampFormatInfo()))
- ), filters, Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
+ ), filters, null,
+ Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
new FieldInfo("ts", "3", new TimestampFormatInfo()),
OrderDirection.ASC);
}
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/InnerJoinRelationShipSqlParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/InnerJoinRelationShipSqlParseTest.java
index b96c1887b..1eab9ff2c 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/InnerJoinRelationShipSqlParseTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/InnerJoinRelationShipSqlParseTest.java
@@ -130,7 +130,7 @@ public class InnerJoinRelationShipSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("5", "kafka_output", fields, relations,
+ return new KafkaLoadNode("5", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, null);
@@ -159,7 +159,7 @@ public class InnerJoinRelationShipSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("5", "kafka_output", fields, relations,
+ return new KafkaLoadNode("5", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -188,7 +188,7 @@ public class InnerJoinRelationShipSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("ts", "3", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
- ), null);
+ ), null, null);
}
/**
@@ -224,7 +224,8 @@ public class InnerJoinRelationShipSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new TimestampFormatInfo()))
- ), filters, Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
+ ), filters, null,
+ Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
new FieldInfo("ts", "3", new TimestampFormatInfo()),
OrderDirection.ASC);
}
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/LeftOuterJoinSqlParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/LeftOuterJoinSqlParseTest.java
index 920934b14..645e63915 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/LeftOuterJoinSqlParseTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/LeftOuterJoinSqlParseTest.java
@@ -29,6 +29,7 @@ import
org.apache.inlong.sort.formats.common.TimestampFormatInfo;
import org.apache.inlong.sort.protocol.FieldInfo;
import org.apache.inlong.sort.protocol.GroupInfo;
import org.apache.inlong.sort.protocol.StreamInfo;
+import org.apache.inlong.sort.protocol.enums.FilterStrategy;
import org.apache.inlong.sort.protocol.enums.ScanStartupMode;
import org.apache.inlong.sort.protocol.node.Node;
import org.apache.inlong.sort.protocol.node.extract.KafkaExtractNode;
@@ -130,7 +131,7 @@ public class LeftOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("5", "kafka_output", fields, relations,
+ return new KafkaLoadNode("5", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -159,7 +160,7 @@ public class LeftOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("ts", "3", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
- ), null);
+ ), null, null);
}
/**
@@ -195,7 +196,8 @@ public class LeftOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new TimestampFormatInfo()))
- ), filters, Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
+ ), filters, FilterStrategy.RETAIN,
+ Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
new FieldInfo("ts", "3", new TimestampFormatInfo()),
OrderDirection.ASC);
}
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/MetaFieldSyncTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/MetaFieldSyncTest.java
index 2237cb67e..e1e2bf574 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/MetaFieldSyncTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/MetaFieldSyncTest.java
@@ -151,7 +151,7 @@ public class MetaFieldSyncTest extends AbstractTestBase {
new FieldRelationShip(new FieldInfo("up_before", new
TimestampFormatInfo()),
new FieldInfo("up_before", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("2", "kafka_output", fields, relations,
+ return new KafkaLoadNode("2", "kafka_output", fields, relations, null,
null, "topic1", "localhost:9092",
new CanalJsonFormat(), null,
null, "id");
@@ -251,7 +251,7 @@ public class MetaFieldSyncTest extends AbstractTestBase {
new FieldRelationShip(new FieldInfo("up_before", new
TimestampFormatInfo()),
new FieldInfo("up_before", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("4", "kafka_output2", fields, relations,
+ return new KafkaLoadNode("4", "kafka_output2", fields, relations, null,
null, "topic2", "localhost:9092",
new CanalJsonFormat(), null,
null, "id");
diff --git
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/RightOuterJoinSqlParseTest.java
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/RightOuterJoinSqlParseTest.java
index 49cb773c6..2314ba4f5 100644
---
a/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/RightOuterJoinSqlParseTest.java
+++
b/inlong-sort/sort-single-tenant/src/test/java/org/apache/inlong/sort/singletenant/flink/parser/RightOuterJoinSqlParseTest.java
@@ -130,7 +130,7 @@ public class RightOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new
TimestampFormatInfo()))
);
- return new KafkaLoadNode("5", "kafka_output", fields, relations,
+ return new KafkaLoadNode("5", "kafka_output", fields, relations, null,
null, "topic_output", "localhost:9092",
new JsonFormat(), null,
null, "id");
@@ -159,7 +159,7 @@ public class RightOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("ts", "3", new
TimestampFormatInfo()),
new FieldInfo("ts", new TimestampFormatInfo()))
- ), null);
+ ), null, null);
}
/**
@@ -195,7 +195,8 @@ public class RightOuterJoinSqlParseTest extends
AbstractTestBase {
new FieldInfo("ts", new TimestampFormatInfo())),
new FieldRelationShip(new FieldInfo("salary", "3", new
TimestampFormatInfo()),
new FieldInfo("salary", new TimestampFormatInfo()))
- ), filters, Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
+ ), filters, null,
+ Collections.singletonList(new FieldInfo("name", "1", new
StringFormatInfo())),
new FieldInfo("ts", "3", new TimestampFormatInfo()),
OrderDirection.ASC);
}