From 303b2eb47e64431abaff60fc12d92832a2720d45 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 10:38:17 +0800 Subject: [PATCH] refactor: enforce kafka mq sink type --- .../ingest/sink/mq/SinkMqAutoConfiguration.java | 3 +-- .../com/lingniu/ingest/sink/mq/SinkMqProperties.java | 11 +++++++++-- .../lingniu/ingest/sink/mq/SinkMqPropertiesTest.java | 11 +++++++++++ 3 files changed, 21 insertions(+), 4 deletions(-) diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java index b5a6c61c..b89feaac 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java @@ -23,8 +23,7 @@ import java.util.Properties; *
  • {@code lingniu.ingest.sink.mq.enabled=true}(默认 true)—— 总开关, * 设为 false 时本模块完全不装配,ingest-core 的 DisruptorEventBus 仍然运行但 * 没有 Kafka sink。 - *
  • {@code lingniu.ingest.sink.mq.type=kafka}(默认 kafka)—— 后端选择,预留 - * rocketmq/pulsar 等扩展。 + *
  • {@code lingniu.ingest.sink.mq.type=kafka}(默认 kafka)—— 唯一生产 MQ 后端。 * * *

    Producer 和 Consumer 是两个独立开关:{@code sink.mq.enabled=true} 只表示可以创建 diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java index b7b16299..d24cf2f7 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java @@ -5,6 +5,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; +import java.util.Locale; import java.util.Map; @ConfigurationProperties(prefix = "lingniu.ingest.sink.mq") @@ -18,7 +19,7 @@ public class SinkMqProperties { */ private boolean enabled = true; - /** MQ 后端类型。目前仅支持 kafka,预留 rocketmq/pulsar 等。 */ + /** MQ 后端类型。生产链路只允许 kafka。 */ private String type = "kafka"; /** Kafka bootstrap servers;生产环境应通过环境变量覆盖,不建议使用默认开发地址。 */ private String bootstrapServers = "114.55.58.251:9092"; @@ -35,7 +36,13 @@ public class SinkMqProperties { public boolean isEnabled() { return enabled; } public void setEnabled(boolean enabled) { this.enabled = enabled; } public String getType() { return type; } - public void setType(String type) { this.type = type; } + public void setType(String type) { + String value = type == null || type.isBlank() ? "kafka" : type.trim().toLowerCase(Locale.ROOT); + if (!"kafka".equals(value)) { + throw new IllegalStateException("sink.mq.type supports only kafka; configured value: " + type); + } + this.type = value; + } public String getBootstrapServers() { return bootstrapServers; } public void setBootstrapServers(String bootstrapServers) { this.bootstrapServers = bootstrapServers; } public String getCompressionType() { return compressionType; } diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java index 39bbb3c1..4f3c54dd 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqPropertiesTest.java @@ -3,6 +3,7 @@ package com.lingniu.ingest.sink.mq; import org.junit.jupiter.api.Test; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; class SinkMqPropertiesTest { @@ -15,4 +16,14 @@ class SinkMqPropertiesTest { props.setBootstrapServers("kafka.internal:9092"); assertThat(props.getBootstrapServers()).isEqualTo("kafka.internal:9092"); } + + @Test + void rejectsUnsupportedMqTypeInsteadOfSilentlyDisablingKafkaBeans() { + SinkMqProperties props = new SinkMqProperties(); + + assertThatThrownBy(() -> props.setType("rocketmq")) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("sink.mq.type supports only kafka") + .hasMessageContaining("rocketmq"); + } }