refactor: gate mqtt inbound startup
This commit is contained in:
@@ -31,6 +31,7 @@ lingniu:
|
|||||||
ingest:
|
ingest:
|
||||||
mqtt:
|
mqtt:
|
||||||
enabled: ${YUTONG_MQTT_ENABLED:false}
|
enabled: ${YUTONG_MQTT_ENABLED:false}
|
||||||
|
auto-startup: ${YUTONG_MQTT_AUTO_STARTUP:true}
|
||||||
endpoints:
|
endpoints:
|
||||||
- name: ${YUTONG_MQTT_ENDPOINT_NAME:yutong}
|
- name: ${YUTONG_MQTT_ENDPOINT_NAME:yutong}
|
||||||
uri: ${YUTONG_MQTT_URI:}
|
uri: ${YUTONG_MQTT_URI:}
|
||||||
|
|||||||
@@ -15,12 +15,16 @@ import com.lingniu.ingest.sink.mq.KafkaEventSink;
|
|||||||
import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration;
|
import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration;
|
||||||
import org.apache.kafka.clients.producer.KafkaProducer;
|
import org.apache.kafka.clients.producer.KafkaProducer;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.junit.jupiter.api.extension.ExtendWith;
|
||||||
import org.springframework.boot.autoconfigure.AutoConfigurations;
|
import org.springframework.boot.autoconfigure.AutoConfigurations;
|
||||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
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.assertj.core.api.Assertions.assertThat;
|
||||||
import static org.mockito.Mockito.mock;
|
import static org.mockito.Mockito.mock;
|
||||||
|
|
||||||
|
@ExtendWith(OutputCaptureExtension.class)
|
||||||
class YutongMqttAppCompositionTest {
|
class YutongMqttAppCompositionTest {
|
||||||
|
|
||||||
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
|
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
|
||||||
@@ -34,6 +38,7 @@ class YutongMqttAppCompositionTest {
|
|||||||
.withBean("kafkaProducer", KafkaProducer.class, YutongMqttAppCompositionTest::kafkaProducer)
|
.withBean("kafkaProducer", KafkaProducer.class, YutongMqttAppCompositionTest::kafkaProducer)
|
||||||
.withPropertyValues(
|
.withPropertyValues(
|
||||||
"lingniu.ingest.mqtt.enabled=true",
|
"lingniu.ingest.mqtt.enabled=true",
|
||||||
|
"lingniu.ingest.mqtt.auto-startup=false",
|
||||||
"lingniu.ingest.mqtt.endpoints[0].name=yutong-test",
|
"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].uri=tcp://127.0.0.1:1883",
|
||||||
"lingniu.ingest.mqtt.endpoints[0].topic=/yutong/#",
|
"lingniu.ingest.mqtt.endpoints[0].topic=/yutong/#",
|
||||||
@@ -51,7 +56,7 @@ class YutongMqttAppCompositionTest {
|
|||||||
"lingniu.ingest.vehicle-stat.enabled=false");
|
"lingniu.ingest.vehicle-stat.enabled=false");
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void createsMqttIngressAndKafkaSinksWithoutHistoryStorage() {
|
void createsMqttIngressAndKafkaSinksWithoutHistoryStorage(CapturedOutput output) {
|
||||||
contextRunner.run(context -> {
|
contextRunner.run(context -> {
|
||||||
assertThat(context).hasSingleBean(MqttEndpointManager.class);
|
assertThat(context).hasSingleBean(MqttEndpointManager.class);
|
||||||
assertThat(context).hasSingleBean(MqttProfileRegistry.class);
|
assertThat(context).hasSingleBean(MqttProfileRegistry.class);
|
||||||
@@ -61,6 +66,17 @@ class YutongMqttAppCompositionTest {
|
|||||||
assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class);
|
assertThat(context).hasSingleBean(KafkaEnvelopeDeadLetterSink.class);
|
||||||
assertThat(context).hasSingleBean(ArchiveStore.class);
|
assertThat(context).hasSingleBean(ArchiveStore.class);
|
||||||
assertThat(context).hasSingleBean(RawArchiveEventSink.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");
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ class YutongMqttAppDefaultsTest {
|
|||||||
.containsEntry("spring.cloud.nacos.config.server-addr", "${NACOS_SERVER_ADDR:127.0.0.1:8848}")
|
.containsEntry("spring.cloud.nacos.config.server-addr", "${NACOS_SERVER_ADDR:127.0.0.1:8848}")
|
||||||
.containsEntry("server.port", "${HTTP_PORT:20500}")
|
.containsEntry("server.port", "${HTTP_PORT:20500}")
|
||||||
.containsEntry("lingniu.ingest.mqtt.enabled", "${YUTONG_MQTT_ENABLED:false}")
|
.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].uri", "${YUTONG_MQTT_URI:}")
|
||||||
.containsEntry("lingniu.ingest.mqtt.endpoints[0].topic", "${YUTONG_MQTT_TOPIC:#}")
|
.containsEntry("lingniu.ingest.mqtt.endpoints[0].topic", "${YUTONG_MQTT_TOPIC:#}")
|
||||||
.containsEntry("lingniu.ingest.mqtt.endpoints[0].profile", "yutong")
|
.containsEntry("lingniu.ingest.mqtt.endpoints[0].profile", "yutong")
|
||||||
|
|||||||
@@ -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.MqttProfile;
|
||||||
import com.lingniu.ingest.inbound.mqtt.profile.MqttProfileRegistry;
|
import com.lingniu.ingest.inbound.mqtt.profile.MqttProfileRegistry;
|
||||||
import com.lingniu.ingest.inbound.mqtt.profile.YutongMqttProfile;
|
import com.lingniu.ingest.inbound.mqtt.profile.YutongMqttProfile;
|
||||||
|
import org.springframework.boot.ApplicationRunner;
|
||||||
import org.springframework.boot.autoconfigure.AutoConfiguration;
|
import org.springframework.boot.autoconfigure.AutoConfiguration;
|
||||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
|
||||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||||
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
||||||
import org.springframework.context.annotation.Bean;
|
import org.springframework.context.annotation.Bean;
|
||||||
import org.springframework.context.annotation.Conditional;
|
import org.springframework.context.annotation.Conditional;
|
||||||
@@ -62,7 +64,7 @@ public class MqttInboundAutoConfiguration {
|
|||||||
return new MqttRealtimeHandler(profileRegistry);
|
return new MqttRealtimeHandler(profileRegistry);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean(initMethod = "start", destroyMethod = "close")
|
@Bean(destroyMethod = "close")
|
||||||
@Lazy(false)
|
@Lazy(false)
|
||||||
@ConditionalOnMissingBean
|
@ConditionalOnMissingBean
|
||||||
public MqttEndpointManager mqttEndpointManager(MqttInboundProperties props,
|
public MqttEndpointManager mqttEndpointManager(MqttInboundProperties props,
|
||||||
@@ -74,6 +76,15 @@ public class MqttInboundAutoConfiguration {
|
|||||||
return new MqttEndpointManager(props, profileRegistry, identityResolver, dispatcher);
|
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) {
|
private static void applyLegacySingleEndpointFallback(MqttInboundProperties props, Environment env) {
|
||||||
if (!props.getEndpoints().isEmpty()) {
|
if (!props.getEndpoints().isEmpty()) {
|
||||||
return;
|
return;
|
||||||
|
|||||||
@@ -16,11 +16,15 @@ public class MqttInboundProperties {
|
|||||||
|
|
||||||
/** 总开关。关闭时不会创建任何 MQTT 客户端连接。 */
|
/** 总开关。关闭时不会创建任何 MQTT 客户端连接。 */
|
||||||
private boolean enabled = false;
|
private boolean enabled = false;
|
||||||
|
/** 是否随 Spring Boot 应用启动立即连接 broker。测试装配时可关闭,生产默认开启。 */
|
||||||
|
private boolean autoStartup = true;
|
||||||
/** 多 broker/多 topic 配置;每个 endpoint 独立重连、订阅并转成内部 VehicleEvent。 */
|
/** 多 broker/多 topic 配置;每个 endpoint 独立重连、订阅并转成内部 VehicleEvent。 */
|
||||||
private List<Endpoint> endpoints = new ArrayList<>();
|
private List<Endpoint> endpoints = new ArrayList<>();
|
||||||
|
|
||||||
public boolean isEnabled() { return enabled; }
|
public boolean isEnabled() { return enabled; }
|
||||||
public void setEnabled(boolean enabled) { this.enabled = enabled; }
|
public void setEnabled(boolean enabled) { this.enabled = enabled; }
|
||||||
|
public boolean isAutoStartup() { return autoStartup; }
|
||||||
|
public void setAutoStartup(boolean autoStartup) { this.autoStartup = autoStartup; }
|
||||||
public List<Endpoint> getEndpoints() { return endpoints; }
|
public List<Endpoint> getEndpoints() { return endpoints; }
|
||||||
public void setEndpoints(List<Endpoint> endpoints) { this.endpoints = endpoints; }
|
public void setEndpoints(List<Endpoint> endpoints) { this.endpoints = endpoints; }
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user