This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] 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 5dbfb374f9 [Improve][Connector-V2][Pulsar] Migrate validation to 
declarative OptionRule (#11985)
5dbfb374f9 is described below

commit 5dbfb374f985349aeefde9cf84169aa98b3ac5ca
Author: Johan Lin <[email protected]>
AuthorDate: Mon Aug 31 17:38:10 2026 +0000

    [Improve][Connector-V2][Pulsar] Migrate validation to declarative 
OptionRule (#11985)
    
    Co-authored-by: David Zollo <[email protected]>
---
 .../seatunnel/pulsar/config/PulsarAdminConfig.java |   5 -
 .../pulsar/config/PulsarClientConfig.java          |   5 -
 .../pulsar/config/PulsarConsumerConfig.java        |   6 -
 .../seatunnel/pulsar/sink/PulsarSinkFactory.java   |   8 +-
 .../pulsar/source/PulsarSourceFactory.java         |   8 +-
 .../pulsar/sink/PulsarSinkFactoryTest.java         |  72 ++++++++++
 .../pulsar/source/PulsarSourceFactoryTest.java     | 152 +++++++++++++++++++++
 7 files changed, 238 insertions(+), 18 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarAdminConfig.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarAdminConfig.java
index 07f3737e49..b8e127ddf6 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarAdminConfig.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarAdminConfig.java
@@ -17,9 +17,6 @@
 
 package org.apache.seatunnel.connectors.seatunnel.pulsar.config;
 
-import org.apache.pulsar.shade.com.google.common.base.Preconditions;
-import org.apache.pulsar.shade.org.apache.commons.lang3.StringUtils;
-
 // TODO: more field
 
 public class PulsarAdminConfig extends BasePulsarConfig {
@@ -65,8 +62,6 @@ public class PulsarAdminConfig extends BasePulsarConfig {
         }
 
         public PulsarAdminConfig build() {
-            Preconditions.checkArgument(
-                    StringUtils.isNotBlank(adminUrl), "Pulsar admin URL is 
required.");
             return new PulsarAdminConfig(authPluginClassName, authParams, 
adminUrl);
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarClientConfig.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarClientConfig.java
index d69870ff74..642c2655eb 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarClientConfig.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarClientConfig.java
@@ -17,9 +17,6 @@
 
 package org.apache.seatunnel.connectors.seatunnel.pulsar.config;
 
-import org.apache.pulsar.shade.com.google.common.base.Preconditions;
-import org.apache.pulsar.shade.org.apache.commons.lang3.StringUtils;
-
 // TODO: more field
 
 public class PulsarClientConfig extends BasePulsarConfig {
@@ -66,8 +63,6 @@ public class PulsarClientConfig extends BasePulsarConfig {
         }
 
         public PulsarClientConfig build() {
-            Preconditions.checkArgument(
-                    StringUtils.isNotBlank(serviceUrl), "Pulsar service URL is 
required.");
             return new PulsarClientConfig(authPluginClassName, authParams, 
serviceUrl);
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarConsumerConfig.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarConsumerConfig.java
index 1563082efa..8bb62bac0f 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarConsumerConfig.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/config/PulsarConsumerConfig.java
@@ -19,9 +19,6 @@ package 
org.apache.seatunnel.connectors.seatunnel.pulsar.config;
 
 // TODO: more field
 
-import org.apache.pulsar.shade.com.google.common.base.Preconditions;
-import org.apache.pulsar.shade.org.apache.commons.lang3.StringUtils;
-
 import java.io.Serializable;
 
 public class PulsarConsumerConfig implements Serializable {
@@ -52,9 +49,6 @@ public class PulsarConsumerConfig implements Serializable {
         }
 
         public PulsarConsumerConfig build() {
-            Preconditions.checkArgument(
-                    StringUtils.isNotBlank(subscriptionName),
-                    "Pulsar subscription name is required.");
             return new PulsarConsumerConfig(subscriptionName);
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactory.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactory.java
index caf550218e..a1fb885db6 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactory.java
@@ -18,6 +18,7 @@
 package org.apache.seatunnel.connectors.seatunnel.pulsar.sink;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.Conditions;
 import org.apache.seatunnel.api.configuration.util.OptionRule;
 import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
 import org.apache.seatunnel.api.table.connector.TableSink;
@@ -40,7 +41,12 @@ public class PulsarSinkFactory implements TableSinkFactory {
     @Override
     public OptionRule optionRule() {
         return OptionRule.builder()
-                .required(PulsarSinkOptions.CLIENT_SERVICE_URL, 
PulsarSinkOptions.ADMIN_SERVICE_URL)
+                .required(
+                        PulsarSinkOptions.CLIENT_SERVICE_URL,
+                        
Conditions.notBlank(PulsarSinkOptions.CLIENT_SERVICE_URL))
+                .required(
+                        PulsarSinkOptions.ADMIN_SERVICE_URL,
+                        
Conditions.notBlank(PulsarSinkOptions.ADMIN_SERVICE_URL))
                 .optional(
                         PulsarSinkOptions.TOPIC,
                         PulsarSinkOptions.FORMAT,
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactory.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactory.java
index 80b2f15b52..a26dba8cd5 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactory.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactory.java
@@ -19,6 +19,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.pulsar.source;
 
 import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.Conditions;
 import org.apache.seatunnel.api.configuration.util.OptionRule;
 import org.apache.seatunnel.api.options.table.TableSchemaOptions;
 import org.apache.seatunnel.api.source.SeaTunnelSource;
@@ -48,9 +49,14 @@ public class PulsarSourceFactory implements 
TableSourceFactory {
         return OptionRule.builder()
                 .required(
                         PulsarSourceOptions.CLIENT_SERVICE_URL,
-                        PulsarSourceOptions.ADMIN_SERVICE_URL)
+                        
Conditions.notBlank(PulsarSourceOptions.CLIENT_SERVICE_URL))
+                .required(
+                        PulsarSourceOptions.ADMIN_SERVICE_URL,
+                        
Conditions.notBlank(PulsarSourceOptions.ADMIN_SERVICE_URL))
                 .optional(
                         PulsarSourceOptions.SUBSCRIPTION_NAME,
+                        
Conditions.notBlank(PulsarSourceOptions.SUBSCRIPTION_NAME))
+                .optional(
                         PulsarSourceOptions.CURSOR_STARTUP_MODE,
                         PulsarSourceOptions.CURSOR_STOP_MODE,
                         PulsarSourceOptions.TOPIC_DISCOVERY_INTERVAL,
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactoryTest.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactoryTest.java
index 4859f2aa39..9afe50d6c5 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkFactoryTest.java
@@ -18,7 +18,9 @@
 package org.apache.seatunnel.connectors.seatunnel.pulsar.sink;
 
 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.options.SinkConnectorCommonOptions;
 import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.catalog.Column;
@@ -82,6 +84,76 @@ public class PulsarSinkFactoryTest {
                         
.contains(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA));
     }
 
+    @Test
+    void testValidSinkConfig() {
+        Map<String, Object> options = validSinkOptions();
+        Assertions.assertDoesNotThrow(() -> validate(options));
+    }
+
+    @Test
+    void testMissingClientServiceUrlFails() {
+        Map<String, Object> options = validSinkOptions();
+        options.remove(PulsarSinkOptions.CLIENT_SERVICE_URL.key());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(options));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSinkOptions.CLIENT_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testMissingAdminServiceUrlFails() {
+        Map<String, Object> options = validSinkOptions();
+        options.remove(PulsarSinkOptions.ADMIN_SERVICE_URL.key());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(options));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSinkOptions.ADMIN_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testAuthOptionsMustBeBundled() {
+        Map<String, Object> options = validSinkOptions();
+        options.put(
+                PulsarSinkOptions.AUTH_PLUGIN_CLASS.key(),
+                "org.apache.pulsar.client.impl.auth.AuthenticationToken");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(options));
+        Assertions.assertTrue(exception.getMessage().contains("bundled"));
+    }
+
+    @Test
+    void testBlankClientServiceUrlFails() {
+        Map<String, Object> options = validSinkOptions();
+        options.put(PulsarSinkOptions.CLIENT_SERVICE_URL.key(), "");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(options));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSinkOptions.CLIENT_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testBlankAdminServiceUrlFails() {
+        Map<String, Object> options = validSinkOptions();
+        options.put(PulsarSinkOptions.ADMIN_SERVICE_URL.key(), "   ");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(options));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSinkOptions.ADMIN_SERVICE_URL.key()));
+    }
+
+    private Map<String, Object> validSinkOptions() {
+        Map<String, Object> options = new HashMap<>();
+        options.put(PulsarSinkOptions.CLIENT_SERVICE_URL.key(), 
"pulsar://localhost:6650");
+        options.put(PulsarSinkOptions.ADMIN_SERVICE_URL.key(), 
"http://localhost:8080";);
+        options.put(PulsarSinkOptions.TOPIC.key(), "test-topic");
+        return options;
+    }
+
+    private void validate(Map<String, Object> options) {
+        ConfigValidator.of(ReadonlyConfig.fromMap(options))
+                .validate(new PulsarSinkFactory().optionRule());
+    }
+
     private ReadonlyConfig config() {
         Map<String, Object> options = new HashMap<>();
         options.put("client.service-url", "pulsar://localhost:6650");
diff --git 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactoryTest.java
 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactoryTest.java
index 6a0773298e..870289a12f 100644
--- 
a/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-pulsar/src/test/java/org/apache/seatunnel/connectors/seatunnel/pulsar/source/PulsarSourceFactoryTest.java
@@ -17,12 +17,18 @@
 
 package org.apache.seatunnel.connectors.seatunnel.pulsar.source;
 
+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.connectors.seatunnel.pulsar.config.PulsarSourceOptions;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.util.HashMap;
+import java.util.Map;
+
 public class PulsarSourceFactoryTest {
 
     @Test
@@ -38,4 +44,150 @@ public class PulsarSourceFactoryTest {
         OptionRule optionRule = pulsarSourceFactory.optionRule();
         Assertions.assertNotNull(optionRule);
     }
+
+    @Test
+    void testValidSourceConfig() {
+        Assertions.assertDoesNotThrow(() -> validate(validSourceConfig()));
+    }
+
+    @Test
+    void testMissingClientServiceUrlFails() {
+        Map<String, Object> config = validSourceConfig();
+        config.remove(PulsarSourceOptions.CLIENT_SERVICE_URL.key());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.CLIENT_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testMissingAdminServiceUrlFails() {
+        Map<String, Object> config = validSourceConfig();
+        config.remove(PulsarSourceOptions.ADMIN_SERVICE_URL.key());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.ADMIN_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testTopicAndTopicPatternAreExclusive() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(PulsarSourceOptions.TOPIC_PATTERN.key(), "test-topic-.*");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(exception.getMessage().contains("mutually 
exclusive"));
+    }
+
+    @Test
+    void testExactlyOneOfTopicSourceMustBeSet() {
+        Map<String, Object> config = validSourceConfig();
+        config.remove(PulsarSourceOptions.TOPIC.key());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(exception.getMessage().contains("exactly one 
option must be set"));
+    }
+
+    @Test
+    void testStartupModeTimestampRequiresTimestamp() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(
+                PulsarSourceOptions.CURSOR_STARTUP_MODE.key(),
+                PulsarSourceOptions.StartMode.TIMESTAMP.name());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                exception
+                        .getMessage()
+                        
.contains(PulsarSourceOptions.CURSOR_STARTUP_TIMESTAMP.key()));
+    }
+
+    @Test
+    void testStartupModeSubscriptionRequiresResetMode() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(
+                PulsarSourceOptions.CURSOR_STARTUP_MODE.key(),
+                PulsarSourceOptions.StartMode.SUBSCRIPTION.name());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.CURSOR_RESET_MODE.key()));
+    }
+
+    @Test
+    void testStopModeTimestampRequiresTimestamp() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(
+                PulsarSourceOptions.CURSOR_STOP_MODE.key(),
+                PulsarSourceOptions.StopMode.TIMESTAMP.name());
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.CURSOR_STOP_TIMESTAMP.key()));
+    }
+
+    @Test
+    void testAuthOptionsMustBeBundled() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(
+                PulsarSourceOptions.AUTH_PLUGIN_CLASS.key(),
+                "org.apache.pulsar.client.impl.auth.AuthenticationToken");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(exception.getMessage().contains("bundled"));
+    }
+
+    @Test
+    void testValidSourceConfigWithBundledAuthOptions() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(
+                PulsarSourceOptions.AUTH_PLUGIN_CLASS.key(),
+                "org.apache.pulsar.client.impl.auth.AuthenticationToken");
+        config.put(PulsarSourceOptions.AUTH_PARAMS.key(), "token:dummy");
+        Assertions.assertDoesNotThrow(() -> validate(config));
+    }
+
+    @Test
+    void testBlankClientServiceUrlFails() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(PulsarSourceOptions.CLIENT_SERVICE_URL.key(), "");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.CLIENT_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testBlankAdminServiceUrlFails() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(PulsarSourceOptions.ADMIN_SERVICE_URL.key(), "   ");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.ADMIN_SERVICE_URL.key()));
+    }
+
+    @Test
+    void testBlankSubscriptionNameFails() {
+        Map<String, Object> config = validSourceConfig();
+        config.put(PulsarSourceOptions.SUBSCRIPTION_NAME.key(), "");
+        OptionValidationException exception =
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        Assertions.assertTrue(
+                
exception.getMessage().contains(PulsarSourceOptions.SUBSCRIPTION_NAME.key()));
+    }
+
+    private Map<String, Object> validSourceConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put(PulsarSourceOptions.CLIENT_SERVICE_URL.key(), 
"pulsar://localhost:6650");
+        config.put(PulsarSourceOptions.ADMIN_SERVICE_URL.key(), 
"http://localhost:8080";);
+        config.put(PulsarSourceOptions.SUBSCRIPTION_NAME.key(), 
"seatunnel-subscription");
+        config.put(PulsarSourceOptions.TOPIC.key(), "test-topic");
+        return config;
+    }
+
+    private void validate(Map<String, Object> config) {
+        ConfigValidator.of(ReadonlyConfig.fromMap(config))
+                .validate(new PulsarSourceFactory().optionRule());
+    }
 }

Reply via email to