From b1190cd7c75aedf8838bc975cbf47b8d411d1e57 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 01:42:24 +0800 Subject: [PATCH] refactor: route jt808 mileage through dedicated stream --- .../VehicleStatEnvelopeIngestor.java | 14 ++++++--- .../VehicleStatEnvelopeIngestorTest.java | 29 +++++++++++++++++-- 2 files changed, 36 insertions(+), 7 deletions(-) diff --git a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java index a7bef3f6..3ea8c2cc 100644 --- a/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java +++ b/modules/services/vehicle-stat-service/src/main/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestor.java @@ -29,8 +29,10 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor { if (!envelope.hasTelemetrySnapshot()) { return; } + if (processJt808Mileage(envelope)) { + return; + } processor.process(envelope); - processJt808Mileage(envelope); } @Override @@ -42,8 +44,10 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor { return EnvelopeIngestResult.skipped( envelope.getEventId(), envelope.getVin(), "envelope telemetry_snapshot is required"); } + if (processJt808Mileage(envelope)) { + return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin()); + } processor.process(envelope); - processJt808Mileage(envelope); return EnvelopeIngestResult.processed(envelope.getEventId(), envelope.getVin()); } catch (IllegalArgumentException ex) { return envelope == null @@ -52,10 +56,12 @@ public final class VehicleStatEnvelopeIngestor implements EnvelopeIngestor { } } - private void processJt808Mileage(VehicleEnvelope envelope) { - if (jt808MileageProcessor != null) { + private boolean processJt808Mileage(VehicleEnvelope envelope) { + if (jt808MileageProcessor != null && "JT808".equalsIgnoreCase(envelope.getSource())) { jt808MileageProcessor.process(envelope); + return true; } + return false; } private static VehicleEnvelope parse(byte[] kafkaValue) { diff --git a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java index ac8a12a8..d3b8dec9 100644 --- a/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java +++ b/modules/services/vehicle-stat-service/src/test/java/com/lingniu/ingest/vehiclestat/VehicleStatEnvelopeIngestorTest.java @@ -17,6 +17,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; class VehicleStatEnvelopeIngestorTest { @@ -35,10 +36,27 @@ class VehicleStatEnvelopeIngestorTest { } @Test - void fansOutTelemetryEnvelopeToOptionalJt808MileageProcessor() { + void doesNotFanOutNonJt808TelemetryToJt808MileageProcessor() { InMemoryVehicleStatRepository repository = new InMemoryVehicleStatRepository(); Jt808MileageStreamProcessor jt808Processor = mock(Jt808MileageStreamProcessor.class); - VehicleEnvelope envelope = envelope(); + VehicleEnvelope envelope = envelope("GB32960"); + VehicleStatEnvelopeIngestor ingestor = new VehicleStatEnvelopeIngestor( + new VehicleStatEventProcessor(repository), + jt808Processor); + + ingestor.ingest(envelope.toByteArray()); + + verifyNoInteractions(jt808Processor); + assertThat(repository.mileagePoints("VIN001", LocalDate.of(2026, 6, 22))) + .extracting(MileagePoint::totalMileageKm) + .containsExactly(123.45); + } + + @Test + void routesJt808TelemetryOnlyToJt808MileageProcessorWhenAvailable() { + InMemoryVehicleStatRepository repository = new InMemoryVehicleStatRepository(); + Jt808MileageStreamProcessor jt808Processor = mock(Jt808MileageStreamProcessor.class); + VehicleEnvelope envelope = envelope("JT808"); VehicleStatEnvelopeIngestor ingestor = new VehicleStatEnvelopeIngestor( new VehicleStatEventProcessor(repository), jt808Processor); @@ -46,6 +64,7 @@ class VehicleStatEnvelopeIngestorTest { ingestor.ingest(envelope.toByteArray()); verify(jt808Processor).process(envelope); + assertThat(repository.mileagePoints("VIN001", LocalDate.of(2026, 6, 22))).isEmpty(); } @Test @@ -110,10 +129,14 @@ class VehicleStatEnvelopeIngestorTest { } private static VehicleEnvelope envelope() { + return envelope("GB32960"); + } + + private static VehicleEnvelope envelope(String source) { return VehicleEnvelope.newBuilder() .setEventId("event-1") .setVin("VIN001") - .setSource("GB32960") + .setSource(source) .setEventTimeMs(1_782_112_400_000L) .setIngestTimeMs(1_782_112_401_000L) .setTelemetrySnapshot(TelemetrySnapshot.newBuilder()