feat: write history envelopes to tdengine
This commit is contained in:
@@ -16,6 +16,10 @@
|
|||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>event-file-store</artifactId>
|
<artifactId>event-file-store</artifactId>
|
||||||
</dependency>
|
</dependency>
|
||||||
|
<dependency>
|
||||||
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
|
<artifactId>tdengine-history-store</artifactId>
|
||||||
|
</dependency>
|
||||||
<dependency>
|
<dependency>
|
||||||
<groupId>com.lingniu.ingest</groupId>
|
<groupId>com.lingniu.ingest</groupId>
|
||||||
<artifactId>sink-mq</artifactId>
|
<artifactId>sink-mq</artifactId>
|
||||||
|
|||||||
@@ -6,8 +6,11 @@ import com.lingniu.ingest.api.consumer.EnvelopeIngestor;
|
|||||||
import com.lingniu.ingest.eventfilestore.EventFileRecord;
|
import com.lingniu.ingest.eventfilestore.EventFileRecord;
|
||||||
import com.lingniu.ingest.eventfilestore.EventFileStore;
|
import com.lingniu.ingest.eventfilestore.EventFileStore;
|
||||||
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
|
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.io.IOException;
|
||||||
|
import java.util.List;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Consumer-side entry point that stores one Kafka envelope value into the
|
* 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 EventFileStore store;
|
||||||
private final TelemetryEnvelopeRecordMapper mapper;
|
private final TelemetryEnvelopeRecordMapper mapper;
|
||||||
|
private final TdengineHistoryWriter tdengineWriter;
|
||||||
|
|
||||||
public EventHistoryEnvelopeIngestor(EventFileStore store, TelemetryEnvelopeRecordMapper mapper) {
|
public EventHistoryEnvelopeIngestor(EventFileStore store, TelemetryEnvelopeRecordMapper mapper) {
|
||||||
|
this(store, mapper, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
public EventHistoryEnvelopeIngestor(EventFileStore store,
|
||||||
|
TelemetryEnvelopeRecordMapper mapper,
|
||||||
|
TdengineHistoryWriter tdengineWriter) {
|
||||||
if (store == null) {
|
if (store == null) {
|
||||||
throw new IllegalArgumentException("store must not be null");
|
throw new IllegalArgumentException("store must not be null");
|
||||||
}
|
}
|
||||||
@@ -30,12 +40,14 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
}
|
}
|
||||||
this.store = store;
|
this.store = store;
|
||||||
this.mapper = mapper;
|
this.mapper = mapper;
|
||||||
|
this.tdengineWriter = tdengineWriter;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void ingest(byte[] kafkaValue) throws IOException {
|
public void ingest(byte[] kafkaValue) throws IOException {
|
||||||
VehicleEnvelope envelope = parse(kafkaValue);
|
VehicleEnvelope envelope = parse(kafkaValue);
|
||||||
EventFileRecord record = mapper.toRecord(envelope);
|
EventFileRecord record = mapper.toRecord(envelope);
|
||||||
store.append(record);
|
store.append(record);
|
||||||
|
writeTdengineFacts(envelope);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -45,6 +57,7 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
envelope = parse(kafkaValue);
|
envelope = parse(kafkaValue);
|
||||||
EventFileRecord record = mapper.toRecord(envelope);
|
EventFileRecord record = mapper.toRecord(envelope);
|
||||||
store.append(record);
|
store.append(record);
|
||||||
|
writeTdengineFacts(envelope);
|
||||||
return EnvelopeIngestResult.stored(record.eventId(), record.vin());
|
return EnvelopeIngestResult.stored(record.eventId(), record.vin());
|
||||||
} catch (IllegalArgumentException ex) {
|
} catch (IllegalArgumentException ex) {
|
||||||
// 协议不可解析或缺少必要字段属于不可重试问题,避免 Kafka consumer 无限重放。
|
// 协议不可解析或缺少必要字段属于不可重试问题,避免 Kafka consumer 无限重放。
|
||||||
@@ -70,4 +83,18 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor {
|
|||||||
throw new IllegalArgumentException("failed to parse VehicleEnvelope", e);
|
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()));
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,12 +12,14 @@ 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.SinkArchiveAutoConfiguration;
|
||||||
import com.lingniu.ingest.sink.archive.config.SinkArchiveProperties;
|
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.AutoConfiguration;
|
||||||
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
|
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
|
||||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
|
||||||
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.autoconfigure.condition.ConditionalOnProperty;
|
||||||
import org.springframework.context.annotation.Bean;
|
import org.springframework.context.annotation.Bean;
|
||||||
|
import org.springframework.beans.factory.ObjectProvider;
|
||||||
|
|
||||||
import java.net.URI;
|
import java.net.URI;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
@@ -49,8 +51,9 @@ public class EventHistoryAutoConfiguration {
|
|||||||
@ConditionalOnBean(EventFileStore.class)
|
@ConditionalOnBean(EventFileStore.class)
|
||||||
@ConditionalOnMissingBean
|
@ConditionalOnMissingBean
|
||||||
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(EventFileStore store,
|
public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(EventFileStore store,
|
||||||
TelemetryEnvelopeRecordMapper mapper) {
|
TelemetryEnvelopeRecordMapper mapper,
|
||||||
return new EventHistoryEnvelopeIngestor(store, mapper);
|
ObjectProvider<TdengineHistoryWriter> writer) {
|
||||||
|
return new EventHistoryEnvelopeIngestor(store, mapper, writer.getIfAvailable());
|
||||||
}
|
}
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
|
|||||||
@@ -10,9 +10,15 @@ import com.lingniu.ingest.eventfilestore.EventFileQuery;
|
|||||||
import com.lingniu.ingest.eventfilestore.EventFileRecord;
|
import com.lingniu.ingest.eventfilestore.EventFileRecord;
|
||||||
import com.lingniu.ingest.eventfilestore.EventFileStore;
|
import com.lingniu.ingest.eventfilestore.EventFileStore;
|
||||||
import com.lingniu.ingest.sink.mq.proto.RawArchiveRef;
|
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.TelemetryField;
|
||||||
import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot;
|
import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot;
|
||||||
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
|
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 org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
@@ -118,6 +124,53 @@ class EventHistoryEnvelopeIngestorTest {
|
|||||||
assertThat(payload.toString()).doesNotContain("rawBytes");
|
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) {
|
private static VehicleEnvelope envelope(String eventId) {
|
||||||
return VehicleEnvelope.newBuilder()
|
return VehicleEnvelope.newBuilder()
|
||||||
.setEventId(eventId)
|
.setEventId(eventId)
|
||||||
@@ -149,4 +202,19 @@ class EventHistoryEnvelopeIngestorTest {
|
|||||||
return List.copyOf(records);
|
return List.copyOf(records);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static final class CapturingTdengineWriter implements TdengineHistoryWriter {
|
||||||
|
private final List<TdengineRawFrameRow> rawFrames = new ArrayList<>();
|
||||||
|
private final List<TdengineLocationRow> locations = new ArrayList<>();
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void appendRawFrames(List<TdengineRawFrameRow> rows) {
|
||||||
|
rawFrames.addAll(rows);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void appendLocations(List<TdengineLocationRow> rows) {
|
||||||
|
locations.addAll(rows);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user