diff --git a/modules/services/event-history-service/pom.xml b/modules/services/event-history-service/pom.xml index 959d5483..10389f5c 100644 --- a/modules/services/event-history-service/pom.xml +++ b/modules/services/event-history-service/pom.xml @@ -16,6 +16,10 @@ com.lingniu.ingest event-file-store + + com.lingniu.ingest + tdengine-history-store + com.lingniu.ingest sink-mq diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java index 3b79d8bf..2e8ecd97 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java @@ -6,8 +6,11 @@ import com.lingniu.ingest.api.consumer.EnvelopeIngestor; import com.lingniu.ingest.eventfilestore.EventFileRecord; import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.tdenginehistory.TdengineEnvelopeRows; +import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import java.io.IOException; +import java.util.List; /** * Consumer-side entry point that stores one Kafka envelope value into the @@ -20,8 +23,15 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { private final EventFileStore store; private final TelemetryEnvelopeRecordMapper mapper; + private final TdengineHistoryWriter tdengineWriter; public EventHistoryEnvelopeIngestor(EventFileStore store, TelemetryEnvelopeRecordMapper mapper) { + this(store, mapper, null); + } + + public EventHistoryEnvelopeIngestor(EventFileStore store, + TelemetryEnvelopeRecordMapper mapper, + TdengineHistoryWriter tdengineWriter) { if (store == null) { throw new IllegalArgumentException("store must not be null"); } @@ -30,12 +40,14 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { } this.store = store; this.mapper = mapper; + this.tdengineWriter = tdengineWriter; } public void ingest(byte[] kafkaValue) throws IOException { VehicleEnvelope envelope = parse(kafkaValue); EventFileRecord record = mapper.toRecord(envelope); store.append(record); + writeTdengineFacts(envelope); } @Override @@ -45,6 +57,7 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { envelope = parse(kafkaValue); EventFileRecord record = mapper.toRecord(envelope); store.append(record); + writeTdengineFacts(envelope); return EnvelopeIngestResult.stored(record.eventId(), record.vin()); } catch (IllegalArgumentException ex) { // 协议不可解析或缺少必要字段属于不可重试问题,避免 Kafka consumer 无限重放。 @@ -70,4 +83,18 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { throw new IllegalArgumentException("failed to parse VehicleEnvelope", e); } } + + private void writeTdengineFacts(VehicleEnvelope envelope) throws IOException { + if (tdengineWriter == null) { + return; + } + var rawFrame = TdengineEnvelopeRows.rawFrame(envelope); + if (rawFrame.isPresent()) { + tdengineWriter.appendRawFrames(List.of(rawFrame.get())); + } + var location = TdengineEnvelopeRows.location(envelope); + if (location.isPresent()) { + tdengineWriter.appendLocations(List.of(location.get())); + } + } } diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java index 530d5b51..2bc1f6c8 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java @@ -12,12 +12,14 @@ import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; 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.TdengineHistoryWriter; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; +import org.springframework.beans.factory.ObjectProvider; import java.net.URI; import java.nio.file.Path; @@ -49,8 +51,9 @@ public class EventHistoryAutoConfiguration { @ConditionalOnBean(EventFileStore.class) @ConditionalOnMissingBean public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(EventFileStore store, - TelemetryEnvelopeRecordMapper mapper) { - return new EventHistoryEnvelopeIngestor(store, mapper); + TelemetryEnvelopeRecordMapper mapper, + ObjectProvider writer) { + return new EventHistoryEnvelopeIngestor(store, mapper, writer.getIfAvailable()); } @Bean diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java index 6a142bda..cbcf2d93 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java @@ -10,9 +10,15 @@ import com.lingniu.ingest.eventfilestore.EventFileQuery; import com.lingniu.ingest.eventfilestore.EventFileRecord; import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; +import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; import com.lingniu.ingest.sink.mq.proto.TelemetryField; import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; +import com.lingniu.ingest.sink.mq.proto.LocationPayload; +import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; +import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; +import com.lingniu.ingest.tdenginehistory.TdengineLocationRow; +import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow; import org.junit.jupiter.api.Test; import java.io.IOException; @@ -118,6 +124,53 @@ class EventHistoryEnvelopeIngestorTest { assertThat(payload.toString()).doesNotContain("rawBytes"); } + @Test + void tryIngestWritesRawAndLocationFactsToTdengineWhenWriterExists() throws Exception { + CapturingStore store = new CapturingStore(); + CapturingTdengineWriter tdengineWriter = new CapturingTdengineWriter(); + EventHistoryEnvelopeIngestor ingestor = new EventHistoryEnvelopeIngestor( + store, new TelemetryEnvelopeRecordMapper(), tdengineWriter); + + VehicleEnvelope envelope = VehicleEnvelope.newBuilder() + .setSchemaVersion("1.0") + .setEventId("jt808-location-1") + .setVin("VINJT808001") + .setSource("JT808") + .setEventTimeMs(1_782_112_400_000L) + .setIngestTimeMs(1_782_112_401_000L) + .putMetadata("vehicle_key", "jt808:g7gps") + .putMetadata("frame_id", "frame-jt808-1") + .setRawArchive(RawArchiveRef.newBuilder() + .setUri("archive://jt808/2026/06/29/frame-jt808-1.bin") + .setSizeBytes(68)) + .setRawFrameFact(RawFrameFactPayload.newBuilder() + .setFrameId("frame-jt808-1") + .setVehicleKey("jt808:g7gps") + .setVin("VINJT808001") + .setPhone("013800000000") + .setMessageId(0x0200) + .setRawUri("archive://jt808/2026/06/29/frame-jt808-1.bin") + .setRawSizeBytes(68) + .setParseStatus(ParseStatusProto.PARSE_STATUS_SUCCEEDED)) + .setLocation(LocationPayload.newBuilder() + .setLongitude(113.12) + .setLatitude(23.45) + .setSpeedKmh(42.5)) + .build(); + + EnvelopeIngestResult result = ingestor.tryIngest(envelope.toByteArray()); + + assertThat(result.status()).isEqualTo(EnvelopeIngestResult.Status.STORED); + assertThat(store.records).hasSize(1); + assertThat(tdengineWriter.rawFrames) + .extracting(TdengineRawFrameRow::frameId) + .containsExactly("frame-jt808-1"); + assertThat(tdengineWriter.locations) + .extracting(TdengineLocationRow::factId) + .containsExactly("jt808-location-1"); + assertThat(tdengineWriter.locations.getFirst().longitude()).isEqualTo(113.12); + } + private static VehicleEnvelope envelope(String eventId) { return VehicleEnvelope.newBuilder() .setEventId(eventId) @@ -149,4 +202,19 @@ class EventHistoryEnvelopeIngestorTest { return List.copyOf(records); } } + + private static final class CapturingTdengineWriter implements TdengineHistoryWriter { + private final List rawFrames = new ArrayList<>(); + private final List locations = new ArrayList<>(); + + @Override + public void appendRawFrames(List rows) { + rawFrames.addAll(rows); + } + + @Override + public void appendLocations(List rows) { + locations.addAll(rows); + } + } }