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

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12250-7e468d60f67832e3e5fa00bf3035fbf193e74238
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 313cf73196c6a2d0828fb945726ffe2d20e813d1
Author: smoggy666 <[email protected]>
AuthorDate: Thu Sep 10 15:13:40 2026 +0000

    [Improve][Connector-V2] Add declarative validation for Druid sink options 
(#12250)
---
 docs/en/connectors/sink/Druid.md                   |  4 +-
 docs/zh/connectors/sink/Druid.md                   |  4 +-
 .../connectors/druid/sink/DruidSinkFactory.java    |  7 +-
 .../seatunnel/druid/DruidFactoryTest.java          | 90 +++++++++++++++++++++-
 4 files changed, 99 insertions(+), 6 deletions(-)

diff --git a/docs/en/connectors/sink/Druid.md b/docs/en/connectors/sink/Druid.md
index 229dff70ad..8fe4d19499 100644
--- a/docs/en/connectors/sink/Druid.md
+++ b/docs/en/connectors/sink/Druid.md
@@ -36,7 +36,7 @@ Write data to Apache Druid through the Druid indexing task 
API.
 |----------------|--------|----------|---------------|-------------|
 | coordinatorUrl | string | yes      | -             | Druid coordinator or 
router host and port. |
 | datasource     | string | yes      | -             | Druid datasource name. 
Supports placeholders such as `${table_name}`. |
-| batchSize      | int    | no       | 10000         | Number of rows buffered 
before one indexing task is submitted. |
+| batchSize      | int    | no       | 10000         | Number of rows buffered 
before one indexing task is submitted. Must be greater than `0`. |
 | common-options |        | no       | -             | Sink common options. |
 
 ### coordinatorUrl [string]
@@ -53,7 +53,7 @@ When the upstream source has multiple tables, you can use 
placeholders such as `
 
 ### batchSize [int]
 
-The number of rows buffered before SeaTunnel sends one indexing task to Druid. 
The default value is `10000`.
+The number of rows buffered before SeaTunnel sends one indexing task to Druid. 
The default value is `10000`, and the configured value must be greater than `0`.
 
 SeaTunnel also flushes the remaining buffered rows when the writer closes.
 
diff --git a/docs/zh/connectors/sink/Druid.md b/docs/zh/connectors/sink/Druid.md
index 3fa9a7dce0..8e82ba40eb 100644
--- a/docs/zh/connectors/sink/Druid.md
+++ b/docs/zh/connectors/sink/Druid.md
@@ -36,7 +36,7 @@ import ChangeLog from '../changelog/connector-druid.md';
 |----------------|--------|------|--------|------|
 | coordinatorUrl | string | 是   | -      | Druid 协调器或路由节点的主机和端口。 |
 | datasource     | string | 是   | -      | Druid datasource 名称,支持 
`${table_name}` 这类占位符。 |
-| batchSize      | int    | 否   | 10000  | 缓存多少行后提交一次索引任务。 |
+| batchSize      | int    | 否   | 10000  | 缓存多少行后提交一次索引任务,取值必须大于 `0`。 |
 | common-options |        | 否   | -      | Sink 通用参数。 |
 
 ### coordinatorUrl [string]
@@ -53,7 +53,7 @@ SeaTunnel 会向 `http://{coordinatorUrl}/druid/indexer/v1/task` 
提交索引任
 
 ### batchSize [int]
 
-SeaTunnel 缓存多少行之后向 Druid 提交一次索引任务。默认值为 `10000`。
+SeaTunnel 缓存多少行之后向 Druid 提交一次索引任务。默认值为 `10000`,配置值必须大于 `0`。
 
 写入器关闭时,SeaTunnel 也会把剩余缓存数据提交到 Druid。
 
diff --git 
a/seatunnel-connectors-v2/connector-druid/src/main/java/org/apache/seatunnel/connectors/druid/sink/DruidSinkFactory.java
 
b/seatunnel-connectors-v2/connector-druid/src/main/java/org/apache/seatunnel/connectors/druid/sink/DruidSinkFactory.java
index fd4b1783e1..74bfbcb36e 100644
--- 
a/seatunnel-connectors-v2/connector-druid/src/main/java/org/apache/seatunnel/connectors/druid/sink/DruidSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-druid/src/main/java/org/apache/seatunnel/connectors/druid/sink/DruidSinkFactory.java
@@ -29,6 +29,9 @@ import 
org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
 
 import com.google.auto.service.AutoService;
 
+import static 
org.apache.seatunnel.api.configuration.util.Conditions.greaterThan;
+import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank;
+import static 
org.apache.seatunnel.connectors.druid.config.DruidSinkOptions.BATCH_SIZE;
 import static 
org.apache.seatunnel.connectors.druid.config.DruidSinkOptions.COORDINATOR_URL;
 import static 
org.apache.seatunnel.connectors.druid.config.DruidSinkOptions.DATASOURCE;
 
@@ -42,7 +45,9 @@ public class DruidSinkFactory implements TableSinkFactory {
     @Override
     public OptionRule optionRule() {
         return OptionRule.builder()
-                .required(COORDINATOR_URL, DATASOURCE)
+                .required(COORDINATOR_URL, notBlank(COORDINATOR_URL))
+                .required(DATASOURCE, notBlank(DATASOURCE))
+                .optional(BATCH_SIZE, greaterThan(BATCH_SIZE, 0))
                 .optional(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA)
                 .build();
     }
diff --git 
a/seatunnel-connectors-v2/connector-druid/src/test/java/org/apache/seatunnel/connectors/seatunnel/druid/DruidFactoryTest.java
 
b/seatunnel-connectors-v2/connector-druid/src/test/java/org/apache/seatunnel/connectors/seatunnel/druid/DruidFactoryTest.java
index b5b40bb3ec..c826747c83 100644
--- 
a/seatunnel-connectors-v2/connector-druid/src/test/java/org/apache/seatunnel/connectors/seatunnel/druid/DruidFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-druid/src/test/java/org/apache/seatunnel/connectors/seatunnel/druid/DruidFactoryTest.java
@@ -18,14 +18,102 @@
 
 package org.apache.seatunnel.connectors.seatunnel.druid;
 
+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.druid.config.DruidSinkOptions;
 import org.apache.seatunnel.connectors.druid.sink.DruidSinkFactory;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.util.HashMap;
+import java.util.Map;
+
 public class DruidFactoryTest {
+
+    private final OptionRule optionRule = new DruidSinkFactory().optionRule();
+
     @Test
     public void optionRuleTest() {
-        Assertions.assertNotNull((new DruidSinkFactory()).optionRule());
+        Assertions.assertNotNull(optionRule);
+    }
+
+    @Test
+    void testValidRequiredOptions() {
+        Assertions.assertDoesNotThrow(() -> validate(requiredConfig()));
+    }
+
+    @Test
+    void testMissingRequiredOptionsRejected() {
+        for (String key :
+                new String[] {
+                    DruidSinkOptions.COORDINATOR_URL.key(), 
DruidSinkOptions.DATASOURCE.key()
+                }) {
+            Map<String, Object> config = requiredConfig();
+            config.remove(key);
+            Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+        }
+    }
+
+    @Test
+    void testBlankRequiredOptionsRejected() {
+        for (String key :
+                new String[] {
+                    DruidSinkOptions.COORDINATOR_URL.key(), 
DruidSinkOptions.DATASOURCE.key()
+                }) {
+            for (String value : new String[] {"", " \t\r\n "}) {
+                Map<String, Object> config = requiredConfig();
+                config.put(key, value);
+                Assertions.assertThrows(OptionValidationException.class, () -> 
validate(config));
+            }
+        }
+    }
+
+    @Test
+    void testDefaultBatchSize() {
+        ReadonlyConfig config = validate(requiredConfig());
+        Assertions.assertEquals(
+                DruidSinkOptions.BATCH_SIZE_DEFAULT, 
config.get(DruidSinkOptions.BATCH_SIZE));
+    }
+
+    @Test
+    void testPositiveBatchSize() {
+        Assertions.assertDoesNotThrow(() -> validateBatchSize(1));
+        Assertions.assertDoesNotThrow(() -> validateBatchSize(100));
+    }
+
+    @Test
+    void testNonPositiveBatchSizeRejected() {
+        assertInvalidBatchSize(0);
+        assertInvalidBatchSize(-1);
+    }
+
+    private void assertInvalidBatchSize(int batchSize) {
+        OptionValidationException exception =
+                Assertions.assertThrows(
+                        OptionValidationException.class, () -> 
validateBatchSize(batchSize));
+        
Assertions.assertTrue(exception.getMessage().contains(DruidSinkOptions.BATCH_SIZE.key()));
+    }
+
+    private ReadonlyConfig validateBatchSize(int batchSize) {
+        Map<String, Object> config = requiredConfig();
+        config.put(DruidSinkOptions.BATCH_SIZE.key(), batchSize);
+        return validate(config);
+    }
+
+    private ReadonlyConfig validate(Map<String, Object> config) {
+        ReadonlyConfig readonlyConfig = ReadonlyConfig.fromMap(config);
+        ConfigValidator.validateUnknownKeys(readonlyConfig, optionRule, 
"DruidSink");
+        ConfigValidator.of(readonlyConfig).validate(optionRule);
+        return readonlyConfig;
+    }
+
+    private Map<String, Object> requiredConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put(DruidSinkOptions.COORDINATOR_URL.key(), "localhost:8888");
+        config.put(DruidSinkOptions.DATASOURCE.key(), "seatunnel");
+        return config;
     }
 }

Reply via email to