refactor: keep history consumer tdengine only

This commit is contained in:
lingniu
2026-07-01 04:15:43 +08:00
parent f276b6f6f0
commit d5ef061cca
2 changed files with 76 additions and 8 deletions

View File

@@ -2,7 +2,6 @@ package com.lingniu.ingest.historyapp;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink;
import com.lingniu.ingest.eventfilestore.EventFileStore;
import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor;
import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper;
import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerFactory;
@@ -33,17 +32,12 @@ public class VehicleHistoryKafkaConsumerConfiguration {
@Bean
@ConditionalOnMissingBean
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(ObjectProvider<EventFileStore> store,
TelemetryEnvelopeRecordMapper mapper,
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(TelemetryEnvelopeRecordMapper mapper,
ObjectProvider<TdengineHistoryWriter> writer,
@Value("${lingniu.ingest.tdengine-history.telemetry-fields-enabled:false}")
boolean telemetryFieldsEnabled) {
TdengineHistoryWriter tdengineWriter = writer.getIfAvailable();
EventFileStore eventFileStore = store.getIfAvailable();
if (eventFileStore == null) {
return new EventHistoryEnvelopeIngestor(mapper, tdengineWriter, telemetryFieldsEnabled);
}
return new EventHistoryEnvelopeIngestor(eventFileStore, mapper, tdengineWriter, telemetryFieldsEnabled);
return new EventHistoryEnvelopeIngestor(mapper, tdengineWriter, telemetryFieldsEnabled);
}
@Bean

View File

@@ -1,6 +1,7 @@
package com.lingniu.ingest.historyapp;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink;
import com.lingniu.ingest.eventfilestore.EventFileStore;
import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor;
import com.lingniu.ingest.eventhistory.Gb32960DecodedFrameService;
import com.lingniu.ingest.eventhistory.Gb32960FrameController;
@@ -17,9 +18,14 @@ 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.KafkaEventSink;
import com.lingniu.ingest.sink.mq.proto.ParseStatusProto;
import com.lingniu.ingest.sink.mq.proto.RawArchiveRef;
import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration;
import com.lingniu.ingest.tdenginehistory.TdengineHistorySchema;
import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter;
import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow;
import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.junit.jupiter.api.Test;
@@ -28,9 +34,14 @@ import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
class VehicleHistoryAppCompositionTest {
@@ -99,8 +110,71 @@ class VehicleHistoryAppCompositionTest {
});
}
@Test
void historyIngestorIgnoresEventFileStoreBeanAndWritesTdengineOnly() {
EventFileStore eventFileStore = mock(EventFileStore.class);
CapturingTdengineWriter writer = new CapturingTdengineWriter();
new ApplicationContextRunner()
.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")
.run(context -> {
EventHistoryEnvelopeIngestor ingestor = context.getBean(EventHistoryEnvelopeIngestor.class);
ingestor.tryIngest(rawFrameEnvelope().toByteArray());
verify(eventFileStore, never()).appendAll(any());
assertThat(writer.rawFrames)
.extracting(TdengineRawFrameRow::frameId)
.containsExactly("raw-frame-1");
});
}
@SuppressWarnings("unchecked")
private static KafkaProducer<String, byte[]> kafkaProducer() {
return mock(KafkaProducer.class);
}
private static VehicleEnvelope rawFrameEnvelope() {
return VehicleEnvelope.newBuilder()
.setSchemaVersion("1.0")
.setEventId("raw-event-1")
.setVin("VINRAW001")
.setSource("JT808")
.setEventTimeMs(1_782_112_400_000L)
.setIngestTimeMs(1_782_112_401_000L)
.setRawArchive(RawArchiveRef.newBuilder()
.setUri("archive://jt808/2026/06/29/raw-frame-1.bin")
.setSizeBytes(68))
.setRawFrameFact(RawFrameFactPayload.newBuilder()
.setFrameId("raw-frame-1")
.setVehicleKey("jt808:013800000000")
.setVin("VINRAW001")
.setPhone("013800000000")
.setMessageId(0x0200)
.setRawUri("archive://jt808/2026/06/29/raw-frame-1.bin")
.setRawSizeBytes(68)
.setParseStatus(ParseStatusProto.PARSE_STATUS_SUCCEEDED))
.build();
}
private static final class CapturingTdengineWriter implements TdengineHistoryWriter {
private final List<TdengineRawFrameRow> rawFrames = new ArrayList<>();
@Override
public void appendRawFrames(List<TdengineRawFrameRow> rows) {
rawFrames.addAll(rows);
}
@Override
public void appendLocations(List<com.lingniu.ingest.tdenginehistory.TdengineLocationRow> rows) {
}
@Override
public void appendTelemetryFields(List<com.lingniu.ingest.tdenginehistory.TdengineTelemetryFieldRow> rows) {
}
}
}