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