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-12422-6934fd00626e6347f86320b4f97d7c4cddc0fb59 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 398fdfac5c6f48e894ce94c7b0ba7eeb74c96f7c Author: hutiefang76 <[email protected]> AuthorDate: Tue Sep 22 02:12:30 2026 +0000 [Improve][Connector-V2] Validate Fluss sink bootstrap servers (#12422) --- docs/en/connectors/sink/Fluss.md | 6 ++ docs/zh/connectors/sink/Fluss.md | 6 ++ .../seatunnel/fluss/sink/FlussSinkFactory.java | 5 +- .../seatunnel/fluss/sink/FlussSinkFactoryTest.java | 67 ++++++++++++++++++++++ 4 files changed, 83 insertions(+), 1 deletion(-) diff --git a/docs/en/connectors/sink/Fluss.md b/docs/en/connectors/sink/Fluss.md index 186269b8cf..4dc168fc02 100644 --- a/docs/en/connectors/sink/Fluss.md +++ b/docs/en/connectors/sink/Fluss.md @@ -48,6 +48,12 @@ The target table schema should match the upstream SeaTunnel row schema by field | multi_table_sink_replica | int | no | 1 | Number of writer replicas for multi-table sink mode. | | common-options | - | no | - | Sink common options. See [Sink Common Options](../common-options/sink-common-options.md). | +### bootstrap.servers + +The Fluss coordinator address. + +This value must not be empty or contain only whitespace. + ### database When `database` is not configured, the sink uses the upstream database name from the input table identifier. diff --git a/docs/zh/connectors/sink/Fluss.md b/docs/zh/connectors/sink/Fluss.md index 55fa9651de..3ea5832c41 100644 --- a/docs/zh/connectors/sink/Fluss.md +++ b/docs/zh/connectors/sink/Fluss.md @@ -48,6 +48,12 @@ Fluss Sink 用于在批处理或流处理作业中,将 SeaTunnel 数据写入 | multi_table_sink_replica | int | 否 | 1 | 多表写入模式下的 Sink writer 副本数。 | | common-options | - | 否 | - | Sink 通用参数,详见 [Sink Common Options](../common-options/sink-common-options.md)。 | +### bootstrap.servers + +Fluss coordinator 地址。 + +该值不能为空字符串或仅包含空白字符。 + ### database 未配置 `database` 时,Sink 会使用输入表标识中的上游数据库名。 diff --git a/seatunnel-connectors-v2/connector-fluss/src/main/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactory.java b/seatunnel-connectors-v2/connector-fluss/src/main/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactory.java index 13ee142468..570682c85d 100644 --- a/seatunnel-connectors-v2/connector-fluss/src/main/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactory.java +++ b/seatunnel-connectors-v2/connector-fluss/src/main/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactory.java @@ -26,6 +26,7 @@ import org.apache.seatunnel.connectors.seatunnel.fluss.config.FlussSinkOptions; import com.google.auto.service.AutoService; +import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank; import static org.apache.seatunnel.api.options.SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA; @AutoService(Factory.class) @@ -38,7 +39,9 @@ public class FlussSinkFactory implements TableSinkFactory { @Override public OptionRule optionRule() { return OptionRule.builder() - .required(FlussSinkOptions.BOOTSTRAP_SERVERS) + .required( + FlussSinkOptions.BOOTSTRAP_SERVERS, + notBlank(FlussSinkOptions.BOOTSTRAP_SERVERS)) .optional(FlussSinkOptions.DATABASE) .optional(FlussSinkOptions.TABLE) .optional(FlussSinkOptions.CLIENT_CONFIG) diff --git a/seatunnel-connectors-v2/connector-fluss/src/test/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactoryTest.java b/seatunnel-connectors-v2/connector-fluss/src/test/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactoryTest.java new file mode 100644 index 0000000000..bc6576dde3 --- /dev/null +++ b/seatunnel-connectors-v2/connector-fluss/src/test/java/org/apache/seatunnel/connectors/seatunnel/fluss/sink/FlussSinkFactoryTest.java @@ -0,0 +1,67 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.seatunnel.connectors.seatunnel.fluss.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.connectors.seatunnel.fluss.config.FlussSinkOptions; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +class FlussSinkFactoryTest { + + private final OptionRule sinkRule = new FlussSinkFactory().optionRule(); + + @Test + void testValidConfig() { + Assertions.assertDoesNotThrow(() -> validate(validConfig())); + } + + @Test + void testMissingBootstrapServersRejected() { + Map<String, Object> cfg = validConfig(); + cfg.remove(FlussSinkOptions.BOOTSTRAP_SERVERS.key()); + Assertions.assertThrows(OptionValidationException.class, () -> validate(cfg)); + } + + @Test + void testBlankBootstrapServersRejected() { + Map<String, Object> emptyCfg = validConfig(); + emptyCfg.put(FlussSinkOptions.BOOTSTRAP_SERVERS.key(), ""); + Assertions.assertThrows(OptionValidationException.class, () -> validate(emptyCfg)); + + Map<String, Object> whitespaceCfg = validConfig(); + whitespaceCfg.put(FlussSinkOptions.BOOTSTRAP_SERVERS.key(), " "); + Assertions.assertThrows(OptionValidationException.class, () -> validate(whitespaceCfg)); + } + + private void validate(Map<String, Object> config) { + ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(sinkRule); + } + + private Map<String, Object> validConfig() { + Map<String, Object> cfg = new HashMap<>(); + cfg.put(FlussSinkOptions.BOOTSTRAP_SERVERS.key(), "localhost:9123"); + return cfg; + } +}
