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