From 99ad73e4c9dfb068c737a405e8823e1eb180bcb8 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 03:50:05 +0800 Subject: [PATCH] refactor: gate mqtt inbound startup --- .../src/main/resources/application.yml | 1 + .../YutongMqttAppCompositionTest.java | 18 +++++++++++++++++- .../YutongMqttAppDefaultsTest.java | 1 + .../config/MqttInboundAutoConfiguration.java | 13 ++++++++++++- .../mqtt/config/MqttInboundProperties.java | 4 ++++ 5 files changed, 35 insertions(+), 2 deletions(-) diff --git a/modules/apps/yutong-mqtt-app/src/main/resources/application.yml b/modules/apps/yutong-mqtt-app/src/main/resources/application.yml index d8053da5..4db89949 100644 --- a/modules/apps/yutong-mqtt-app/src/main/resources/application.yml +++ b/modules/apps/yutong-mqtt-app/src/main/resources/application.yml @@ -31,6 +31,7 @@ lingniu: ingest: mqtt: enabled: ${YUTONG_MQTT_ENABLED:false} + auto-startup: ${YUTONG_MQTT_AUTO_STARTUP:true} endpoints: - name: ${YUTONG_MQTT_ENDPOINT_NAME:yutong} uri: ${YUTONG_MQTT_URI:} diff --git a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java index df4068f0..d47aa351 100644 --- a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java +++ b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppCompositionTest.java @@ -15,12 +15,16 @@ import com.lingniu.ingest.sink.mq.KafkaEventSink; import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.mock; +@ExtendWith(OutputCaptureExtension.class) class YutongMqttAppCompositionTest { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() @@ -34,6 +38,7 @@ class YutongMqttAppCompositionTest { .withBean("kafkaProducer", KafkaProducer.class, YutongMqttAppCompositionTest::kafkaProducer) .withPropertyValues( "lingniu.ingest.mqtt.enabled=true", + "lingniu.ingest.mqtt.auto-startup=false", "lingniu.ingest.mqtt.endpoints[0].name=yutong-test", "lingniu.ingest.mqtt.endpoints[0].uri=tcp://127.0.0.1:1883", "lingniu.ingest.mqtt.endpoints[0].topic=/yutong/#", @@ -51,7 +56,7 @@ class YutongMqttAppCompositionTest { "lingniu.ingest.vehicle-stat.enabled=false"); @Test - void createsMqttIngressAndKafkaSinksWithoutHistoryStorage() { + void createsMqttIngressAndKafkaSinksWithoutHistoryStorage(CapturedOutput output) { contextRunner.run(context -> { assertThat(context).hasSingleBean(MqttEndpointManager.class); assertThat(context).hasSingleBean(MqttProfileRegistry.class); @@ -61,6 +66,17 @@ class YutongMqttAppCompositionTest { assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class); assertThat(context).hasSingleBean(ArchiveStore.class); assertThat(context).hasSingleBean(RawArchiveEventSink.class); + assertThat(context).doesNotHaveBean("mqttEndpointStartupRunner"); + }); + assertThat(output).doesNotContain("mqtt endpoint [yutong-test] initializing"); + assertThat(output).doesNotContain("mqtt endpoint [yutong-test] init failed"); + } + + @Test + void createsStartupRunnerWhenAutoStartupIsEnabled() { + contextRunner.withPropertyValues("lingniu.ingest.mqtt.auto-startup=true").run(context -> { + assertThat(context).hasSingleBean(MqttEndpointManager.class); + assertThat(context).hasBean("mqttEndpointStartupRunner"); }); } diff --git a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java index 968cf326..275b9360 100644 --- a/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java +++ b/modules/apps/yutong-mqtt-app/src/test/java/com/lingniu/ingest/yutongmqttapp/YutongMqttAppDefaultsTest.java @@ -24,6 +24,7 @@ class YutongMqttAppDefaultsTest { .containsEntry("spring.cloud.nacos.config.server-addr", "${NACOS_SERVER_ADDR:127.0.0.1:8848}") .containsEntry("server.port", "${HTTP_PORT:20500}") .containsEntry("lingniu.ingest.mqtt.enabled", "${YUTONG_MQTT_ENABLED:false}") + .containsEntry("lingniu.ingest.mqtt.auto-startup", "${YUTONG_MQTT_AUTO_STARTUP:true}") .containsEntry("lingniu.ingest.mqtt.endpoints[0].uri", "${YUTONG_MQTT_URI:}") .containsEntry("lingniu.ingest.mqtt.endpoints[0].topic", "${YUTONG_MQTT_TOPIC:#}") .containsEntry("lingniu.ingest.mqtt.endpoints[0].profile", "yutong") diff --git a/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundAutoConfiguration.java b/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundAutoConfiguration.java index 1d4d2839..1b173c38 100644 --- a/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundAutoConfiguration.java +++ b/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundAutoConfiguration.java @@ -10,8 +10,10 @@ import com.lingniu.ingest.inbound.mqtt.parser.MqttPayloadParser; import com.lingniu.ingest.inbound.mqtt.profile.MqttProfile; import com.lingniu.ingest.inbound.mqtt.profile.MqttProfileRegistry; import com.lingniu.ingest.inbound.mqtt.profile.YutongMqttProfile; +import org.springframework.boot.ApplicationRunner; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Conditional; @@ -62,7 +64,7 @@ public class MqttInboundAutoConfiguration { return new MqttRealtimeHandler(profileRegistry); } - @Bean(initMethod = "start", destroyMethod = "close") + @Bean(destroyMethod = "close") @Lazy(false) @ConditionalOnMissingBean public MqttEndpointManager mqttEndpointManager(MqttInboundProperties props, @@ -74,6 +76,15 @@ public class MqttInboundAutoConfiguration { return new MqttEndpointManager(props, profileRegistry, identityResolver, dispatcher); } + @Bean + @ConditionalOnProperty(prefix = "lingniu.ingest.mqtt", + name = "auto-startup", + havingValue = "true", + matchIfMissing = true) + public ApplicationRunner mqttEndpointStartupRunner(MqttEndpointManager manager) { + return args -> manager.start(); + } + private static void applyLegacySingleEndpointFallback(MqttInboundProperties props, Environment env) { if (!props.getEndpoints().isEmpty()) { return; diff --git a/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundProperties.java b/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundProperties.java index ca1535fc..5aeced23 100644 --- a/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundProperties.java +++ b/modules/inbound/inbound-mqtt/src/main/java/com/lingniu/ingest/inbound/mqtt/config/MqttInboundProperties.java @@ -16,11 +16,15 @@ public class MqttInboundProperties { /** 总开关。关闭时不会创建任何 MQTT 客户端连接。 */ private boolean enabled = false; + /** 是否随 Spring Boot 应用启动立即连接 broker。测试装配时可关闭,生产默认开启。 */ + private boolean autoStartup = true; /** 多 broker/多 topic 配置;每个 endpoint 独立重连、订阅并转成内部 VehicleEvent。 */ private List endpoints = new ArrayList<>(); public boolean isEnabled() { return enabled; } public void setEnabled(boolean enabled) { this.enabled = enabled; } + public boolean isAutoStartup() { return autoStartup; } + public void setAutoStartup(boolean autoStartup) { this.autoStartup = autoStartup; } public List getEndpoints() { return endpoints; } public void setEndpoints(List endpoints) { this.endpoints = endpoints; }