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

Reply via email to