feat: split history and analytics consumers

This commit is contained in:
lingniu
2026-06-23 17:43:28 +08:00
parent f9aa214dc2
commit 6a8b53bab6
6 changed files with 143 additions and 13 deletions

View File

@@ -17,7 +17,11 @@ public final class VehicleStateEnvelopeIngestor implements EnvelopeIngestor {
}
public void ingest(byte[] kafkaValue) {
updater.update(parse(kafkaValue));
VehicleEnvelope envelope = parse(kafkaValue);
if (!envelope.hasTelemetrySnapshot()) {
return;
}
updater.update(envelope);
}
@Override
@@ -25,6 +29,10 @@ public final class VehicleStateEnvelopeIngestor implements EnvelopeIngestor {
VehicleEnvelope envelope = null;
try {
envelope = parse(kafkaValue);
if (!envelope.hasTelemetrySnapshot()) {
return EnvelopeIngestResult.skipped(
envelope.getEventId(), envelope.getVin(), "envelope telemetry_snapshot is required");
}
// Kafka 消费路径只接受标准 VehicleEnvelope坏消息返回明确结果给处理器写 DLQ。
updater.update(envelope);
return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin());

View File

@@ -2,6 +2,7 @@ package com.lingniu.ingest.vehiclestate;
import com.lingniu.ingest.api.consumer.EnvelopeIngestResult;
import com.lingniu.ingest.api.consumer.EnvelopeIngestor;
import com.lingniu.ingest.sink.mq.proto.RawArchiveRef;
import com.lingniu.ingest.sink.mq.proto.TelemetryField;
import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
@@ -47,6 +48,25 @@ class VehicleStateEnvelopeIngestorTest {
assertThat(repository.getState("VIN001")).isEmpty();
}
@Test
void rawArchiveEnvelopeIsIgnoredWithoutUpdatingRepository() {
InMemoryVehicleStateRepository repository = new InMemoryVehicleStateRepository();
VehicleStateEnvelopeIngestor ingestor = new VehicleStateEnvelopeIngestor(new VehicleStateUpdater(repository));
VehicleEnvelope envelope = rawArchiveEnvelope();
ingestor.ingest(envelope.toByteArray());
var result = ingestor.tryIngest(envelope.toByteArray());
assertThat(result.status()).isEqualTo(EnvelopeIngestResult.Status.SKIPPED);
assertThat(result.eventId()).isEqualTo("raw-event-1");
assertThat(result.vin()).isEqualTo("VIN001");
assertThat(result.message()).contains("telemetry_snapshot");
assertThat(repository.getState("VIN001")).isEmpty();
assertThat(repository.getLocation("VIN001")).isEmpty();
assertThat(repository.getSafety("VIN001")).isEmpty();
assertThat(repository.getLastEvent("VIN001")).isEmpty();
}
private static VehicleEnvelope envelope() {
return VehicleEnvelope.newBuilder()
.setEventId("event-1")
@@ -64,6 +84,20 @@ class VehicleStateEnvelopeIngestorTest {
.build();
}
private static VehicleEnvelope rawArchiveEnvelope() {
return VehicleEnvelope.newBuilder()
.setEventId("raw-event-1")
.setVin("VIN001")
.setSource("GB32960")
.setEventTimeMs(1_782_112_400_000L)
.setIngestTimeMs(1_782_112_401_000L)
.setRawArchive(RawArchiveRef.newBuilder()
.setUri("archive://2026/06/23/GB32960/VIN001/raw-event-1.bin")
.setSizeBytes(128)
.build())
.build();
}
private static final class EmptyRepository implements VehicleStateRepository {
@Override public void putState(String vin, String json) {}
@Override public void putLocation(String vin, String json) {}