This is an automated email from the ASF dual-hosted git repository.
nzw921rx pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 857d138ef7 [Improve][Transform-V2] Migrate FilterRowKind validation to
OptionRule (#11763)
857d138ef7 is described below
commit 857d138ef7accc2eb6429a31f9f6df10d6d6f899
Author: Goutam Adwant <[email protected]>
AuthorDate: Sun Sep 13 00:39:15 2026 -0700
[Improve][Transform-V2] Migrate FilterRowKind validation to OptionRule
(#11763)
Signed-off-by: Goutam Adwant <[email protected]>
Signed-off-by: goutamadwant <[email protected]>
---
.../filterrowkind/FilterRowKindTransform.java | 12 +--
.../FilterRowKindTransformFactory.java | 9 +-
.../FilterRowKindTransformFactoryTest.java | 98 ++++++++++++++++++++--
3 files changed, 102 insertions(+), 17 deletions(-)
diff --git
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransform.java
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransform.java
index 9df88d5061..7dccec5d11 100644
---
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransform.java
+++
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransform.java
@@ -56,15 +56,6 @@ public class FilterRowKindTransform extends
FilterRowTransform {
} else {
includeKinds = new
HashSet<>(config.get(FilterRowKinkTransformConfig.INCLUDE_KINDS));
}
- if ((includeKinds.isEmpty() && excludeKinds.isEmpty())
- || (!includeKinds.isEmpty() && !excludeKinds.isEmpty())) {
- throw new SeaTunnelRuntimeException(
- CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,
- String.format(
- "These options(%s,%s) are mutually exclusive,
allowing only one set of options to be configured.",
- FilterRowKinkTransformConfig.INCLUDE_KINDS.key(),
- FilterRowKinkTransformConfig.EXCLUDE_KINDS.key()));
- }
}
@Override
@@ -73,8 +64,7 @@ public class FilterRowKindTransform extends
FilterRowTransform {
return this.excludeKinds.contains(inputRow.getRowKind()) ? null :
inputRow;
}
if (!this.includeKinds.isEmpty()) {
- Set<RowKind> includeKinds = this.includeKinds;
- return includeKinds.contains(inputRow.getRowKind()) ? inputRow :
null;
+ return this.includeKinds.contains(inputRow.getRowKind()) ?
inputRow : null;
}
throw new SeaTunnelRuntimeException(
CommonErrorCodeDeprecated.UNSUPPORTED_OPERATION,
diff --git
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransformFactory.java
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransformFactory.java
index 5023f37e60..80eb85695e 100644
---
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransformFactory.java
+++
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/filterrowkind/FilterRowKindTransformFactory.java
@@ -17,6 +17,7 @@
package org.apache.seatunnel.transform.filterrowkind;
+import org.apache.seatunnel.api.configuration.util.Conditions;
import org.apache.seatunnel.api.configuration.util.OptionRule;
import org.apache.seatunnel.api.table.connector.TableTransform;
import org.apache.seatunnel.api.table.factory.Factory;
@@ -36,9 +37,15 @@ public class FilterRowKindTransformFactory implements
TableTransformFactory {
@Override
public OptionRule optionRule() {
return OptionRule.builder()
+ .exclusive(
+ FilterRowKinkTransformConfig.INCLUDE_KINDS,
+ FilterRowKinkTransformConfig.EXCLUDE_KINDS)
+ .optional(
+ FilterRowKinkTransformConfig.INCLUDE_KINDS,
+
Conditions.notEmpty(FilterRowKinkTransformConfig.INCLUDE_KINDS))
.optional(
FilterRowKinkTransformConfig.EXCLUDE_KINDS,
- FilterRowKinkTransformConfig.INCLUDE_KINDS)
+
Conditions.notEmpty(FilterRowKinkTransformConfig.EXCLUDE_KINDS))
.optional(TransformCommonOptions.MULTI_TABLES)
.optional(TransformCommonOptions.TABLE_MATCH_REGEX)
.optional(TransformCommonOptions.RULE_MATCH_MODE)
diff --git
a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/FilterRowKindTransformFactoryTest.java
b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/FilterRowKindTransformFactoryTest.java
index 0637a8a707..ca16321bb0 100644
---
a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/FilterRowKindTransformFactoryTest.java
+++
b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/FilterRowKindTransformFactoryTest.java
@@ -17,17 +17,105 @@
package org.apache.seatunnel.transform;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
+import org.apache.seatunnel.transform.filterrowkind.FilterRowKindTransform;
import
org.apache.seatunnel.transform.filterrowkind.FilterRowKindTransformFactory;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-public class FilterRowKindTransformFactoryTest {
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+class FilterRowKindTransformFactoryTest {
+
+ private final OptionRule rule = new
FilterRowKindTransformFactory().optionRule();
+
+ private void validate(Map<String, Object> config) {
+ ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(rule);
+ }
+
+ private void assertValidationFails(Map<String, Object> config, String...
optionKeys) {
+ OptionValidationException exception =
+ Assertions.assertThrows(OptionValidationException.class, () ->
validate(config));
+ for (String optionKey : optionKeys) {
+ Assertions.assertTrue(
+ exception.getMessage().contains(optionKey),
+ "Should mention " + optionKey + ": " +
exception.getMessage());
+ }
+ }
@Test
- public void testOptionRule() throws Exception {
- FilterRowKindTransformFactory filterRowKindTransformFactory =
- new FilterRowKindTransformFactory();
- Assertions.assertNotNull(filterRowKindTransformFactory.optionRule());
+ void testIncludeKindsValid() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("include_kinds", Arrays.asList("INSERT", "UPDATE_AFTER"));
+ Assertions.assertDoesNotThrow(() -> validate(cfg));
+ }
+
+ @Test
+ void testExcludeKindsValid() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("exclude_kinds", Arrays.asList("DELETE", "UPDATE_BEFORE"));
+ Assertions.assertDoesNotThrow(() -> validate(cfg));
+ }
+
+ @Test
+ void testNeitherKindsFails() {
+ Map<String, Object> cfg = new HashMap<>();
+ assertValidationFails(cfg, "include_kinds", "exclude_kinds");
+ }
+
+ @Test
+ void testBothKindsFail() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("include_kinds", Arrays.asList("INSERT"));
+ cfg.put("exclude_kinds", Arrays.asList("DELETE"));
+ assertValidationFails(cfg, "include_kinds", "exclude_kinds");
+ }
+
+ @Test
+ void testIncludeKindsEmptyFails() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("include_kinds", Collections.emptyList());
+ assertValidationFails(cfg, "include_kinds");
+ }
+
+ @Test
+ void testExcludeKindsEmptyFails() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("exclude_kinds", Collections.emptyList());
+ assertValidationFails(cfg, "exclude_kinds");
+ }
+
+ @Test
+ void testDirectConstructionWithEmptyKindsFailsOnTransform() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("exclude_kinds", Collections.emptyList());
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"id"}, new SeaTunnelDataType[]
{BasicType.INT_TYPE});
+ CatalogTable catalogTable = CatalogTableUtil.getCatalogTable("test",
rowType);
+ FilterRowKindTransform transform =
+ new FilterRowKindTransform(ReadonlyConfig.fromMap(cfg),
catalogTable);
+
+ SeaTunnelRuntimeException exception =
+ Assertions.assertThrows(
+ SeaTunnelRuntimeException.class,
+ () -> transform.map(new SeaTunnelRow(new Object[]
{1})));
+ Assertions.assertTrue(
+ exception.getMessage().contains("Either excludeKinds or
includeKinds"),
+ "Should explain the required options: " +
exception.getMessage());
}
}