refactor: slim history app consumer config

This commit is contained in:
lingniu
2026-07-01 04:46:47 +08:00
parent f1f29ddd4e
commit 84a098d84c
3 changed files with 23 additions and 29 deletions

View File

@@ -3,15 +3,11 @@ package com.lingniu.ingest.historyapp;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink;
import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor;
import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper;
import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerFactory;
import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner;
import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerWorker;
import com.lingniu.ingest.sink.mq.SinkMqProperties;
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
@@ -24,29 +20,6 @@ import java.util.Map;
@Configuration(proxyBeanMethods = false)
public class VehicleHistoryKafkaConsumerConfiguration {
@Bean
@ConditionalOnMissingBean
public TelemetryEnvelopeRecordMapper telemetryEnvelopeRecordMapper() {
return new TelemetryEnvelopeRecordMapper();
}
@Bean
@ConditionalOnMissingBean
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(TelemetryEnvelopeRecordMapper mapper,
ObjectProvider<TdengineHistoryWriter> writer,
@Value("${lingniu.ingest.tdengine-history.telemetry-fields-enabled:false}")
boolean telemetryFieldsEnabled) {
TdengineHistoryWriter tdengineWriter = writer.getIfAvailable();
return new EventHistoryEnvelopeIngestor(mapper, tdengineWriter, telemetryFieldsEnabled);
}
@Bean
@ConditionalOnMissingBean(name = "eventHistoryEnvelopeConsumerProcessor")
public EnvelopeConsumerProcessor eventHistoryEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor,
EnvelopeDeadLetterSink deadLetterSink) {
return new EnvelopeConsumerProcessor("event-history", ingestor, deadLetterSink);
}
@Bean
@ConditionalOnMissingBean(name = "eventHistoryRawEnvelopeConsumerProcessor")
public EnvelopeConsumerProcessor eventHistoryRawEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor,

View File

@@ -32,9 +32,12 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Bean;
import java.lang.reflect.Method;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
@@ -48,6 +51,19 @@ class VehicleHistoryAppCompositionTest {
@TempDir
Path tempDir;
@Test
void kafkaConsumerConfigurationDoesNotRedeclareServiceOwnedHistoryBeans() {
List<String> beanMethods = Arrays.stream(VehicleHistoryKafkaConsumerConfiguration.class.getDeclaredMethods())
.filter(method -> method.isAnnotationPresent(Bean.class))
.map(Method::getName)
.toList();
assertThat(beanMethods)
.doesNotContain("telemetryEnvelopeRecordMapper")
.doesNotContain("eventHistoryEnvelopeIngestor")
.doesNotContain("eventHistoryEnvelopeConsumerProcessor");
}
@Test
void createsTdengineHistoryStorageAndQueryBeansWithoutGb32960TcpServerOrEventFileStore() {
new ApplicationContextRunner()
@@ -99,10 +115,12 @@ class VehicleHistoryAppCompositionTest {
@Test
void createsHistoryIngestorWithoutEventFileStoreForTdengineOnlyRuntime() {
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(EventHistoryAutoConfiguration.class))
.withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class)
.withBean(TdengineHistoryWriter.class, () -> mock(TdengineHistoryWriter.class))
.withBean(EnvelopeDeadLetterSink.class, () -> mock(EnvelopeDeadLetterSink.class))
.withPropertyValues(
"lingniu.ingest.event-history.enabled=true",
"lingniu.ingest.sink.mq.consumer.enabled=false")
.run(context -> {
assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class);
@@ -116,11 +134,14 @@ class VehicleHistoryAppCompositionTest {
CapturingTdengineWriter writer = new CapturingTdengineWriter();
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(EventHistoryAutoConfiguration.class))
.withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class)
.withBean(EventFileStore.class, () -> eventFileStore)
.withBean(TdengineHistoryWriter.class, () -> writer)
.withBean(EnvelopeDeadLetterSink.class, () -> mock(EnvelopeDeadLetterSink.class))
.withPropertyValues("lingniu.ingest.sink.mq.consumer.enabled=false")
.withPropertyValues(
"lingniu.ingest.event-history.enabled=true",
"lingniu.ingest.sink.mq.consumer.enabled=false")
.run(context -> {
EventHistoryEnvelopeIngestor ingestor = context.getBean(EventHistoryEnvelopeIngestor.class);

View File

@@ -61,7 +61,7 @@ public class EventHistoryAutoConfiguration {
@Bean
@ConditionalOnBean(EventFileStore.class)
@ConditionalOnMissingBean
@ConditionalOnMissingBean(value = {TdengineHistoryWriter.class, EventHistoryEnvelopeIngestor.class})
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(EventFileStore store,
TelemetryEnvelopeRecordMapper mapper,
ObjectProvider<TdengineHistoryWriter> writer) {