refactor: remove archive sink from history runtime
This commit is contained in:
@@ -25,10 +25,6 @@
|
|||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>sink-mq</artifactId>
|
<artifactId>sink-mq</artifactId>
|
||||||
</dependency>
|
</dependency>
|
||||||
<dependency>
|
|
||||||
<groupId>com.lingniu.ingest</groupId>
|
|
||||||
<artifactId>sink-archive</artifactId>
|
|
||||||
</dependency>
|
|
||||||
<dependency>
|
<dependency>
|
||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>tdengine-history-store</artifactId>
|
<artifactId>tdengine-history-store</artifactId>
|
||||||
|
|||||||
@@ -81,9 +81,7 @@ lingniu:
|
|||||||
topics:
|
topics:
|
||||||
- ${KAFKA_TOPIC_YUTONG_MQTT_RAW:vehicle.raw.mqtt-yutong.v1}
|
- ${KAFKA_TOPIC_YUTONG_MQTT_RAW:vehicle.raw.mqtt-yutong.v1}
|
||||||
archive:
|
archive:
|
||||||
enabled: ${SINK_ARCHIVE_ENABLED:true}
|
enabled: false
|
||||||
type: local
|
|
||||||
path: ${SINK_ARCHIVE_PATH:./archive/}
|
|
||||||
tdengine-history:
|
tdengine-history:
|
||||||
enabled: ${TDENGINE_HISTORY_ENABLED:false}
|
enabled: ${TDENGINE_HISTORY_ENABLED:false}
|
||||||
database: ${TDENGINE_HISTORY_DATABASE:vehicle_history}
|
database: ${TDENGINE_HISTORY_DATABASE:vehicle_history}
|
||||||
|
|||||||
@@ -15,9 +15,6 @@ import com.lingniu.ingest.eventhistory.TelemetryFieldHistoryController;
|
|||||||
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
||||||
import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration;
|
import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration;
|
||||||
import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer;
|
import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer;
|
||||||
import com.lingniu.ingest.sink.archive.ArchiveStore;
|
|
||||||
import com.lingniu.ingest.sink.archive.RawArchiveEventSink;
|
|
||||||
import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration;
|
|
||||||
import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink;
|
import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink;
|
||||||
import com.lingniu.ingest.sink.mq.KafkaEventSink;
|
import com.lingniu.ingest.sink.mq.KafkaEventSink;
|
||||||
import com.lingniu.ingest.sink.mq.SinkMqProperties;
|
import com.lingniu.ingest.sink.mq.SinkMqProperties;
|
||||||
@@ -87,10 +84,9 @@ class VehicleHistoryAppCompositionTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void createsTdengineHistoryStorageAndQueryBeansWithoutGb32960TcpServerOrEventFileStore() {
|
void createsTdengineHistoryStorageAndQueryBeansWithoutGb32960TcpServerEventFileStoreOrLocalArchiveSink() {
|
||||||
new ApplicationContextRunner()
|
new ApplicationContextRunner()
|
||||||
.withConfiguration(AutoConfigurations.of(
|
.withConfiguration(AutoConfigurations.of(
|
||||||
SinkArchiveAutoConfiguration.class,
|
|
||||||
TdengineHistoryAutoConfiguration.class,
|
TdengineHistoryAutoConfiguration.class,
|
||||||
SinkMqAutoConfiguration.class,
|
SinkMqAutoConfiguration.class,
|
||||||
Gb32960AutoConfiguration.class,
|
Gb32960AutoConfiguration.class,
|
||||||
@@ -99,9 +95,6 @@ class VehicleHistoryAppCompositionTest {
|
|||||||
.withAllowBeanDefinitionOverriding(true)
|
.withAllowBeanDefinitionOverriding(true)
|
||||||
.withBean("kafkaProducer", KafkaProducer.class, VehicleHistoryAppCompositionTest::kafkaProducer)
|
.withBean("kafkaProducer", KafkaProducer.class, VehicleHistoryAppCompositionTest::kafkaProducer)
|
||||||
.withPropertyValues(
|
.withPropertyValues(
|
||||||
"lingniu.ingest.sink.archive.enabled=true",
|
|
||||||
"lingniu.ingest.sink.archive.type=local",
|
|
||||||
"lingniu.ingest.sink.archive.path=" + tempDir.resolve("archive"),
|
|
||||||
"lingniu.ingest.tdengine-history.enabled=true",
|
"lingniu.ingest.tdengine-history.enabled=true",
|
||||||
"lingniu.ingest.tdengine-history.database=vehicle_history_test",
|
"lingniu.ingest.tdengine-history.database=vehicle_history_test",
|
||||||
"lingniu.ingest.event-history.enabled=true",
|
"lingniu.ingest.event-history.enabled=true",
|
||||||
@@ -114,8 +107,8 @@ class VehicleHistoryAppCompositionTest {
|
|||||||
"lingniu.ingest.vehicle-state.enabled=false",
|
"lingniu.ingest.vehicle-state.enabled=false",
|
||||||
"lingniu.ingest.vehicle-stat.enabled=false")
|
"lingniu.ingest.vehicle-stat.enabled=false")
|
||||||
.run(context -> {
|
.run(context -> {
|
||||||
assertThat(context).hasSingleBean(ArchiveStore.class);
|
assertThat(context).doesNotHaveBean("archiveStore");
|
||||||
assertThat(context).hasSingleBean(RawArchiveEventSink.class);
|
assertThat(context).doesNotHaveBean("rawArchiveEventSink");
|
||||||
assertThat(context).doesNotHaveBean("eventFileStore");
|
assertThat(context).doesNotHaveBean("eventFileStore");
|
||||||
assertThat(context).doesNotHaveBean("eventFileStoreSink");
|
assertThat(context).doesNotHaveBean("eventFileStoreSink");
|
||||||
assertThat(context).hasSingleBean(TdengineHistorySchema.class);
|
assertThat(context).hasSingleBean(TdengineHistorySchema.class);
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ class VehicleHistoryAppDefaultsTest {
|
|||||||
.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("lingniu.ingest.gb32960.enabled", true)
|
.containsEntry("lingniu.ingest.gb32960.enabled", true)
|
||||||
.containsEntry("lingniu.ingest.gb32960.server.enabled", false)
|
.containsEntry("lingniu.ingest.gb32960.server.enabled", false)
|
||||||
.containsEntry("lingniu.ingest.sink.archive.enabled", "${SINK_ARCHIVE_ENABLED:true}")
|
.containsEntry("lingniu.ingest.sink.archive.enabled", false)
|
||||||
.containsEntry("lingniu.ingest.tdengine-history.enabled", "${TDENGINE_HISTORY_ENABLED:false}")
|
.containsEntry("lingniu.ingest.tdengine-history.enabled", "${TDENGINE_HISTORY_ENABLED:false}")
|
||||||
.containsEntry("lingniu.ingest.tdengine-history.database", "${TDENGINE_HISTORY_DATABASE:vehicle_history}")
|
.containsEntry("lingniu.ingest.tdengine-history.database", "${TDENGINE_HISTORY_DATABASE:vehicle_history}")
|
||||||
.containsEntry(
|
.containsEntry(
|
||||||
@@ -104,6 +104,9 @@ class VehicleHistoryAppDefaultsTest {
|
|||||||
|
|
||||||
assertThat(properties.stringPropertyNames())
|
assertThat(properties.stringPropertyNames())
|
||||||
.noneMatch(name -> name.startsWith("lingniu.ingest.event-file-store."));
|
.noneMatch(name -> name.startsWith("lingniu.ingest.event-file-store."));
|
||||||
|
assertThat(properties.stringPropertyNames())
|
||||||
|
.noneMatch(name -> name.equals("lingniu.ingest.sink.archive.type")
|
||||||
|
|| name.equals("lingniu.ingest.sink.archive.path"));
|
||||||
}
|
}
|
||||||
|
|
||||||
private static Properties applicationProperties() {
|
private static Properties applicationProperties() {
|
||||||
|
|||||||
@@ -24,10 +24,6 @@
|
|||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>sink-mq</artifactId>
|
<artifactId>sink-mq</artifactId>
|
||||||
</dependency>
|
</dependency>
|
||||||
<dependency>
|
|
||||||
<groupId>com.lingniu.ingest</groupId>
|
|
||||||
<artifactId>sink-archive</artifactId>
|
|
||||||
</dependency>
|
|
||||||
<dependency>
|
<dependency>
|
||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>protocol-gb32960</artifactId>
|
<artifactId>protocol-gb32960</artifactId>
|
||||||
|
|||||||
@@ -15,8 +15,6 @@ import com.lingniu.ingest.eventhistory.TelemetryFieldHistoryController;
|
|||||||
import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper;
|
import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper;
|
||||||
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
||||||
import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration;
|
import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration;
|
||||||
import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration;
|
|
||||||
import com.lingniu.ingest.sink.archive.config.SinkArchiveProperties;
|
|
||||||
import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration;
|
import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration;
|
||||||
import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader;
|
import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader;
|
||||||
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
|
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
|
||||||
@@ -47,7 +45,6 @@ import java.nio.file.Path;
|
|||||||
@AutoConfiguration
|
@AutoConfiguration
|
||||||
@AutoConfigureAfter({
|
@AutoConfigureAfter({
|
||||||
Gb32960AutoConfiguration.class,
|
Gb32960AutoConfiguration.class,
|
||||||
SinkArchiveAutoConfiguration.class,
|
|
||||||
TdengineHistoryAutoConfiguration.class
|
TdengineHistoryAutoConfiguration.class
|
||||||
})
|
})
|
||||||
@ConditionalOnProperty(prefix = "lingniu.ingest.event-history", name = "enabled", havingValue = "true")
|
@ConditionalOnProperty(prefix = "lingniu.ingest.event-history", name = "enabled", havingValue = "true")
|
||||||
@@ -108,10 +105,11 @@ public class EventHistoryAutoConfiguration {
|
|||||||
public Gb32960DecodedFrameService gb32960DecodedFrameService(ObjectProvider<EventFileStore> store,
|
public Gb32960DecodedFrameService gb32960DecodedFrameService(ObjectProvider<EventFileStore> store,
|
||||||
ObjectProvider<TdengineHistoryReader> reader,
|
ObjectProvider<TdengineHistoryReader> reader,
|
||||||
Gb32960MessageDecoder decoder,
|
Gb32960MessageDecoder decoder,
|
||||||
SinkArchiveProperties archiveProperties) {
|
@Value("${lingniu.ingest.event-history.archive-path:${SINK_ARCHIVE_PATH:./archive/}}")
|
||||||
|
String archivePath) {
|
||||||
// 优先使用 EventFileStore 索引;TDengine-only 高吞吐运行时可直接从 raw_frames 找 rawUri。
|
// 优先使用 EventFileStore 索引;TDengine-only 高吞吐运行时可直接从 raw_frames 找 rawUri。
|
||||||
return new Gb32960DecodedFrameService(store.getIfAvailable(), reader.getIfAvailable(),
|
return new Gb32960DecodedFrameService(store.getIfAvailable(), reader.getIfAvailable(),
|
||||||
decoder, archiveRoot(archiveProperties.getPath()), null);
|
decoder, archiveRoot(archivePath), null);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
|
|||||||
@@ -16,7 +16,6 @@ import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper;
|
|||||||
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder;
|
||||||
import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader;
|
import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader;
|
||||||
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
|
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
|
||||||
import com.lingniu.ingest.sink.archive.config.SinkArchiveProperties;
|
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
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;
|
||||||
@@ -62,14 +61,13 @@ class EventHistoryAutoConfigurationTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void createsGb32960FrameBeansWhenDecoderAndArchivePropertiesExist() {
|
void createsGb32960FrameBeansWhenDecoderExists() {
|
||||||
contextRunner
|
contextRunner
|
||||||
.withPropertyValues(
|
.withPropertyValues(
|
||||||
"lingniu.ingest.event-history.enabled=true",
|
"lingniu.ingest.event-history.enabled=true",
|
||||||
"lingniu.ingest.event-history.api.specialized-enabled=true",
|
"lingniu.ingest.event-history.api.specialized-enabled=true",
|
||||||
"lingniu.ingest.sink.archive.path=/tmp/lingniu-test-archive")
|
"lingniu.ingest.event-history.archive-path=/tmp/lingniu-test-archive")
|
||||||
.withBean(Gb32960MessageDecoder.class, () -> mock(Gb32960MessageDecoder.class))
|
.withBean(Gb32960MessageDecoder.class, () -> mock(Gb32960MessageDecoder.class))
|
||||||
.withBean(SinkArchiveProperties.class, SinkArchiveProperties::new)
|
|
||||||
.run(context -> {
|
.run(context -> {
|
||||||
assertThat(context).hasSingleBean(Gb32960DecodedFrameService.class);
|
assertThat(context).hasSingleBean(Gb32960DecodedFrameService.class);
|
||||||
assertThat(context).hasSingleBean(Gb32960FrameController.class);
|
assertThat(context).hasSingleBean(Gb32960FrameController.class);
|
||||||
@@ -83,10 +81,9 @@ class EventHistoryAutoConfigurationTest {
|
|||||||
.withPropertyValues(
|
.withPropertyValues(
|
||||||
"lingniu.ingest.event-history.enabled=true",
|
"lingniu.ingest.event-history.enabled=true",
|
||||||
"lingniu.ingest.event-history.api.specialized-enabled=true",
|
"lingniu.ingest.event-history.api.specialized-enabled=true",
|
||||||
"lingniu.ingest.sink.archive.path=/tmp/lingniu-test-archive")
|
"lingniu.ingest.event-history.archive-path=/tmp/lingniu-test-archive")
|
||||||
.withBean(Gb32960MessageDecoder.class, () -> mock(Gb32960MessageDecoder.class))
|
.withBean(Gb32960MessageDecoder.class, () -> mock(Gb32960MessageDecoder.class))
|
||||||
.withBean(TdengineHistoryReader.class, () -> mock(TdengineHistoryReader.class))
|
.withBean(TdengineHistoryReader.class, () -> mock(TdengineHistoryReader.class))
|
||||||
.withBean(SinkArchiveProperties.class, SinkArchiveProperties::new)
|
|
||||||
.run(context -> {
|
.run(context -> {
|
||||||
assertThat(context).doesNotHaveBean(EventFileStore.class);
|
assertThat(context).doesNotHaveBean(EventFileStore.class);
|
||||||
assertThat(context).hasSingleBean(Gb32960DecodedFrameService.class);
|
assertThat(context).hasSingleBean(Gb32960DecodedFrameService.class);
|
||||||
|
|||||||
Reference in New Issue
Block a user