diff --git a/DECISIONS.md b/DECISIONS.md index d80b4e3c..30e4a9e4 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -99,3 +99,12 @@ 2. raw-archive-store 仅保留在 `optional-raw-archive-store` profile,需要验证本地读写原型时显式构建。 3. 默认生产链路继续以 `sink-archive`、Kafka raw topic、TDengine raw_frames 为准。 - **Consequences**: 默认构建少一个未部署模块;本地 raw archive 原型仍可通过 `-Poptional-raw-archive-store` 保留和验证。 + +## ADR-015 event-file-store 兼容索引:Optional only +- **Status**: Accepted +- **Context**: 默认历史查询已收敛到 TDengine `raw_frames`、`vehicle_locations` 和按需解码;旧 Parquet/DuckDB `event-file-store` 只保留为兼容索引实现,不应把 DuckDB/Parquet 路径拖进默认生产 reactor。 +- **Decision**: + 1. `event-file-store` 不进入默认 Maven reactor。 + 2. event-file-store 仅保留在 `optional-event-file-store` profile,需要兼容排查旧索引时显式构建。 + 3. `EventFileStore`、`EventFileRecord`、`EventFileQuery` 轻量契约归属 `ingest-api`,history 服务只依赖契约,不依赖兼容实现。 +- **Consequences**: 默认构建面少一个 DuckDB/Parquet 兼容模块;旧索引实现仍可通过 `-Poptional-event-file-store` 单独验证。 diff --git a/README.md b/README.md index ff67a362..7f7571b3 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,7 @@ lingniu-vehicle-ingest/ │ ├── sinks/ │ │ ├── sink-mq/ Kafka producer + Protobuf Envelope │ │ ├── sink-archive/ 原始报文冷存 -│ │ └── event-file-store/ 可选兼容索引(旧 Parquet/DuckDB 路径) +│ │ └── event-file-store/ 可选兼容索引(旧 Parquet/DuckDB 路径,仅 -Poptional-event-file-store 显式构建) │ ├── services/ │ │ ├── event-history-service/ Kafka 全字段事件消费 + 历史查询/导出 │ │ ├── vehicle-state-service/ Kafka 全字段事件消费 + Redis 热状态查询(optional-latest-state profile) @@ -82,6 +82,8 @@ JSATL12 附件上传不在默认 Maven reactor 中;需要时使用 `-Poptional Redis 最新状态查询不在默认 Maven reactor 中;需要时使用 `-Poptional-latest-state` 构建 `vehicle-state-service`。 +旧 Parquet/DuckDB 事件索引不在默认 Maven reactor 中;需要兼容排查时使用 `-Poptional-event-file-store` 构建 `event-file-store`。 + ## 迁移说明 本项目是 `lingniu-vehicle-data-reception` 的 v2 重构,采用 strangler fig 渐进式迁移,旧项目保留只读参考。迁移路径与决策参见 `../REFRACTOR_PLAN.md`。 @@ -98,4 +100,4 @@ Redis 最新状态查询不在默认 Maven reactor 中;需要时使用 `-Popti 2. **协议即插拔**:每个 `protocol-*` / `inbound-*` 都可独立开关(`lingniu.ingest..enabled`) 3. **顺序保证**:同一 VIN 严格有序(Disruptor hash + Kafka 分区 key) 4. **幂等消费**:Envelope 带 `eventId`,下游去重 -5. **原始可回放**:协议接入会把 `rawArchiveKey/rawArchiveUri` 追加到事件 metadata;Kafka Envelope、TDengine `raw_frames` 和导出都使用 `archive://...` 逻辑 URI 追溯原始 bytes,实际文件由 `sink-archive` 管理;`event-file-store` 仅作为旧 Parquet/DuckDB 查询索引兼容路径,不是默认生产历史库 +5. **原始可回放**:协议接入会把 `rawArchiveKey/rawArchiveUri` 追加到事件 metadata;Kafka Envelope、TDengine `raw_frames` 和导出都使用 `archive://...` 逻辑 URI 追溯原始 bytes,实际文件由 `sink-archive` 管理;`event-file-store` 仅作为旧 Parquet/DuckDB 查询索引兼容路径,通过 `optional-event-file-store` profile 显式构建,不是默认生产历史库 diff --git a/docs/target-architecture.md b/docs/target-architecture.md index cda92c87..7f455b8f 100644 --- a/docs/target-architecture.md +++ b/docs/target-architecture.md @@ -190,6 +190,8 @@ flowchart LR through `optional-latest-state`. - raw-archive-store is absent from the default Maven reactor and available through `optional-raw-archive-store`. +- event-file-store is absent from the default Maven reactor and available + through `optional-event-file-store`. - History APIs read from TDengine by default; Parquet/DuckDB is optional compatibility only. - Raw frame queries can return complete parsed JSON through `payloadJson.parsed`. diff --git a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java index 46554d33..821f30c6 100644 --- a/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java +++ b/modules/apps/gb32960-ingest-app/src/test/java/com/lingniu/ingest/gb32960app/Gb32960IngestAppCompositionTest.java @@ -63,7 +63,8 @@ class Gb32960IngestAppCompositionTest { assertThat(context).hasSingleBean(ArchiveStore.class); assertThat(context).hasSingleBean(RawArchiveEventSink.class); - assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.EventFileStore"); + assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.DuckDbParquetEventFileStore"); + assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.config.EventFileStoreAutoConfiguration"); }); } diff --git a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java index 206be0c9..c3fada79 100644 --- a/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java +++ b/modules/apps/jt808-ingest-app/src/test/java/com/lingniu/ingest/jt808app/Jt808IngestAppCompositionTest.java @@ -62,7 +62,8 @@ class Jt808IngestAppCompositionTest { assertThat(context).hasSingleBean(ArchiveStore.class); assertThat(context).hasSingleBean(RawArchiveEventSink.class); - assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.EventFileStore"); + assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.DuckDbParquetEventFileStore"); + assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.config.EventFileStoreAutoConfiguration"); }); } diff --git a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java index ad6425a1..a31b5c42 100644 --- a/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java +++ b/modules/apps/vehicle-analytics-app/src/test/java/com/lingniu/ingest/analyticsapp/VehicleAnalyticsAppCompositionTest.java @@ -58,7 +58,9 @@ class VehicleAnalyticsAppCompositionTest { assertTypeNotPresent(context, "com.lingniu.ingest.vehiclestate.VehicleStateEnvelopeIngestor"); assertTypeNotPresent(context, "com.lingniu.ingest.vehiclestate.config.VehicleStateAutoConfiguration"); assertTypeNotPresent(context, "com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer"); - assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.EventFileStore"); + assertTypeNotPresent(context, "com.lingniu.ingest.eventfilestore.DuckDbParquetEventFileStore"); + assertTypeNotPresent(context, + "com.lingniu.ingest.eventfilestore.config.EventFileStoreAutoConfiguration"); }); } diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java index c1098584..0ceca55c 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/MavenModuleProfileTest.java @@ -63,6 +63,22 @@ class MavenModuleProfileTest { .containsExactly("modules/sinks/raw-archive-store"); } + @Test + void eventFileStoreCompatibilityPathIsOptionalAndOutsideDefaultProductionReactor() throws Exception { + Document pom = rootPom(); + Document eventHistoryServicePom = modulePom("modules/services/event-history-service/pom.xml"); + + assertThat(defaultModules(pom)) + .doesNotContain("modules/sinks/event-file-store"); + + assertThat(profileModules(pom, "optional-event-file-store")) + .containsExactly("modules/sinks/event-file-store"); + + assertThat(hasDependency(eventHistoryServicePom, "com.lingniu.ingest", "event-file-store")) + .as("event-history-service should depend only on the lightweight history contract") + .isFalse(); + } + @Test void xindaIsLegacyOnlyAndOutsideDefaultProductionReactor() throws Exception { Document pom = rootPom(); @@ -97,7 +113,7 @@ class MavenModuleProfileTest { .contains("Netty " + projectProperty(pom, "netty.version")) .contains("Eclipse Paho " + projectProperty(pom, "paho-mqtt.version")) .contains("protocol-jsatl12/ 苏标主动安全报警附件(optional-attachments profile)") - .contains("event-file-store/ 可选兼容索引(旧 Parquet/DuckDB 路径)") + .contains("event-file-store/ 可选兼容索引(旧 Parquet/DuckDB 路径,仅 -Poptional-event-file-store 显式构建)") .contains("vehicle-state-service/ Kafka 全字段事件消费 + Redis 热状态查询(optional-latest-state profile)") .contains("vehicle-history-app/ TDengine 历史查询 + RAW JSON") .contains("vehicle-analytics-app/ JT808 每日里程指标消费") @@ -133,6 +149,7 @@ class MavenModuleProfileTest { .contains("JSATL12 仅保留在 `optional-attachments` profile") .contains("vehicle-state-service 仅保留在 `optional-latest-state` profile") .contains("raw-archive-store 仅保留在 `optional-raw-archive-store` profile") + .contains("event-file-store 仅保留在 `optional-event-file-store` profile") .doesNotContain("## ADR-004 信达 Push:保留模块但彻底重写\n- **Status**: Accepted") .doesNotContain("Spring Boot 3.4.x"); } @@ -225,6 +242,28 @@ class MavenModuleProfileTest { return false; } + private static boolean hasDependency(Document pom, String groupId, String artifactId) { + Element dependencies = firstDirectChild(pom.getDocumentElement(), "dependencies"); + if (dependencies == null) { + return false; + } + NodeList children = dependencies.getChildNodes(); + for (int i = 0; i < children.getLength(); i++) { + if (!(children.item(i) instanceof Element dependency) || !"dependency".equals(dependency.getTagName())) { + continue; + } + Element currentGroupId = firstDirectChild(dependency, "groupId"); + Element currentArtifactId = firstDirectChild(dependency, "artifactId"); + if (currentGroupId != null + && currentArtifactId != null + && groupId.equals(currentGroupId.getTextContent().trim()) + && artifactId.equals(currentArtifactId.getTextContent().trim())) { + return true; + } + } + return false; + } + private static Element firstDirectChild(Element parent, String tagName) { NodeList children = parent.getChildNodes(); for (int i = 0; i < children.getLength(); i++) { diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java index 17dcd880..2324b2c6 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java @@ -3,7 +3,7 @@ package com.lingniu.ingest.historyapp; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.eventhistory.Gb32960DecodedFrameService; import com.lingniu.ingest.eventhistory.Gb32960FrameController; diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileQuery.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileQuery.java similarity index 89% rename from modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileQuery.java rename to modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileQuery.java index 488ee7bb..71233eb2 100644 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileQuery.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileQuery.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.eventfilestore; +package com.lingniu.ingest.api.history; import com.lingniu.ingest.api.ProtocolId; @@ -6,10 +6,7 @@ import java.time.Instant; import java.time.LocalDate; /** - * 文件历史库查询条件。 - * - *

{@code dateFrom/dateTo} 用于定位分区目录,{@code eventTimeFrom/eventTimeTo} - * 用于在分区内做精确时间过滤。这样既支持按天快速裁剪,也支持接口精确到秒的查询。 + * Query criteria for the optional file-backed historical event index. */ public record EventFileQuery( ProtocolId protocol, @@ -43,7 +40,6 @@ public record EventFileQuery( if (eventTimeFrom != null && eventTimeTo != null && eventTimeTo.isBefore(eventTimeFrom)) { throw new IllegalArgumentException("eventTimeTo must not be before eventTimeFrom"); } - // 归一化可选字段,避免空字符串进入 SQL 谓词或文件分区路径。 order = order == null ? Order.ASC : order; limit = limit <= 0 ? 100 : limit; vin = vin == null || vin.isBlank() ? null : vin.trim(); diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileRecord.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileRecord.java similarity index 68% rename from modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileRecord.java rename to modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileRecord.java index 98c0ad4c..1f93cf64 100644 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileRecord.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileRecord.java @@ -1,4 +1,4 @@ -package com.lingniu.ingest.eventfilestore; +package com.lingniu.ingest.api.history; import com.lingniu.ingest.api.ProtocolId; @@ -6,13 +6,7 @@ import java.time.Instant; import java.util.Map; /** - * 文件型明细库的一条可查询记录。 - * - *

{@code payloadJson} 保存统一 telemetry snapshot JSON;公共列只保留排序、分区、 - * 追溯和通用展示需要的最小字段。 - * - *

对 GB32960 来说,{@code rawArchiveUri} 是最关键的追溯列:通用查询可以直接展示 - * payloadJson,专用全字段查询则通过 rawArchiveUri 回读原始包重新解码。 + * Minimal historical event record shared by optional compatibility stores and history queries. */ public record EventFileRecord( String eventId, @@ -38,7 +32,6 @@ public record EventFileRecord( if (ingestTime == null) { throw new IllegalArgumentException("ingestTime must not be null"); } - // 允许 eventType/vin/rawArchiveUri 为空字符串,便于保存平台登录、RAW 索引等非车辆事件。 eventType = eventType == null ? "" : eventType; vin = vin == null ? "" : vin; rawArchiveUri = rawArchiveUri == null ? "" : rawArchiveUri; diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileStore.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileStore.java new file mode 100644 index 00000000..4864d1c8 --- /dev/null +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/history/EventFileStore.java @@ -0,0 +1,26 @@ +package com.lingniu.ingest.api.history; + +import java.io.IOException; +import java.util.List; + +/** + * Historical event index contract. + * + *

The default production history path uses TDengine directly. This contract remains in core so + * optional compatibility implementations can still plug into query code without pulling their + * storage engines into the default reactor. + */ +public interface EventFileStore { + + default void append(EventFileRecord record) throws IOException { + appendAll(List.of(record)); + } + + void appendAll(List records) throws IOException; + + List query(EventFileQuery query) throws IOException; + + default EventFileRecord findByRawArchiveUri(String rawArchiveUri) throws IOException { + return null; + } +} diff --git a/modules/services/event-history-service/pom.xml b/modules/services/event-history-service/pom.xml index 7932ff45..19cfc4e4 100644 --- a/modules/services/event-history-service/pom.xml +++ b/modules/services/event-history-service/pom.xml @@ -14,7 +14,7 @@ com.lingniu.ingest - event-file-store + ingest-api com.lingniu.ingest diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryController.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryController.java index 25bd56a1..6ff8f016 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryController.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryController.java @@ -1,9 +1,9 @@ package com.lingniu.ingest.eventhistory; import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.eventfilestore.EventFileQuery; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.web.bind.annotation.GetMapping; 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 7d53379e..07bc47b9 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 @@ -3,8 +3,8 @@ package com.lingniu.ingest.eventhistory; import com.google.protobuf.InvalidProtocolBufferException; import com.lingniu.ingest.api.consumer.EnvelopeBatchIngestor; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; import com.lingniu.ingest.tdenginehistory.TdengineEnvelopeRows; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameService.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameService.java index 227ec8a7..77370116 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameService.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameService.java @@ -3,9 +3,9 @@ package com.lingniu.ingest.eventhistory; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.eventfilestore.EventFileQuery; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import com.lingniu.ingest.protocol.gb32960.model.Gb32960Message; import com.lingniu.ingest.protocol.gb32960.model.InfoBlock; diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960FrameController.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960FrameController.java index c0cbdd5b..951a4f67 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960FrameController.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/Gb32960FrameController.java @@ -1,6 +1,6 @@ package com.lingniu.ingest.eventhistory; -import com.lingniu.ingest.eventfilestore.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileQuery; import com.lingniu.ingest.tdenginehistory.TdenginePageCursor; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java index d8f13f8d..ffef1246 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapper.java @@ -4,7 +4,7 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.event.RawArchiveKeys; -import com.lingniu.ingest.eventfilestore.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileRecord; import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; import com.google.protobuf.InvalidProtocolBufferException; 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 a72ac3b9..f1dba5dd 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 @@ -2,7 +2,7 @@ package com.lingniu.ingest.eventhistory.config; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.eventhistory.EventHistoryController; import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.eventhistory.Gb32960DecodedFrameService; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/ApiExceptionHandlerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/ApiExceptionHandlerTest.java index b34710fe..597957ee 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/ApiExceptionHandlerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/ApiExceptionHandlerTest.java @@ -1,7 +1,7 @@ package com.lingniu.ingest.eventhistory; -import com.lingniu.ingest.eventfilestore.EventFileStore; -import com.lingniu.ingest.eventfilestore.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; +import com.lingniu.ingest.api.history.EventFileRecord; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import jakarta.servlet.http.HttpServletRequest; import org.junit.jupiter.api.Test; @@ -93,7 +93,7 @@ class ApiExceptionHandlerTest { } @Override - public List query(com.lingniu.ingest.eventfilestore.EventFileQuery query) { + public List query(com.lingniu.ingest.api.history.EventFileQuery query) { return List.of(); } }; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryControllerTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryControllerTest.java index 7e11b7ae..b4a4794c 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryControllerTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryControllerTest.java @@ -1,9 +1,9 @@ package com.lingniu.ingest.eventhistory; import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.eventfilestore.EventFileQuery; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import org.junit.jupiter.api.Test; import java.io.IOException; 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 a04016cb..31f2d2be 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 @@ -6,9 +6,9 @@ import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.api.consumer.EnvelopeIngestResult; import com.lingniu.ingest.api.consumer.EnvelopeIngestor; import com.lingniu.ingest.api.event.RawArchiveKeys; -import com.lingniu.ingest.eventfilestore.EventFileQuery; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.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; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameServiceTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameServiceTest.java index d5db9f1f..f0590cb8 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameServiceTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/Gb32960DecodedFrameServiceTest.java @@ -2,9 +2,9 @@ package com.lingniu.ingest.eventhistory; import com.lingniu.ingest.api.ProtocolId; import com.lingniu.ingest.codec.BccChecksum; -import com.lingniu.ingest.eventfilestore.EventFileQuery; -import com.lingniu.ingest.eventfilestore.EventFileRecord; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960BodyParser; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import com.lingniu.ingest.protocol.gb32960.codec.InfoBlockParser; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java index 875a7415..14f6d38e 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/TelemetryEnvelopeRecordMapperTest.java @@ -1,7 +1,7 @@ package com.lingniu.ingest.eventhistory; import com.lingniu.ingest.api.ProtocolId; -import com.lingniu.ingest.eventfilestore.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileRecord; import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; import com.lingniu.ingest.sink.mq.proto.TelemetryField; import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot; diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java index 9f6f5b9a..754bb8d4 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java @@ -2,7 +2,7 @@ package com.lingniu.ingest.eventhistory.config; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.eventhistory.EventHistoryController; import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.eventhistory.Gb32960DecodedFrameService; diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStore.java b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStore.java index 029b0d46..de35a442 100644 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStore.java +++ b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStore.java @@ -3,6 +3,9 @@ package com.lingniu.ingest.eventfilestore; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import com.lingniu.ingest.api.ProtocolId; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import java.io.IOException; import java.nio.file.Files; diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStore.java b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStore.java deleted file mode 100644 index cfe9300b..00000000 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStore.java +++ /dev/null @@ -1,32 +0,0 @@ -package com.lingniu.ingest.eventfilestore; - -import java.io.IOException; -import java.util.List; - -/** - * 历史明细文件库抽象。 - * - *

写入侧只接受已经标准化的 {@link EventFileRecord};具体实现可以落 Parquet、维护 DuckDB - * sidecar 索引,或在测试中用内存实现。32960 专用 snapshot 查询会先用这里按 VIN/日期找到 - * rawArchiveUri,再回读原始 .bin 解码完整字段。 - */ -public interface EventFileStore { - - /** 单条追加的便捷方法,最终仍走批量写入路径,保证实现只维护一种落盘语义。 */ - default void append(EventFileRecord record) throws IOException { - appendAll(List.of(record)); - } - - void appendAll(List records) throws IOException; - - /** 按协议、日期、VIN、事件类型等条件查询标准化历史记录。 */ - List query(EventFileQuery query) throws IOException; - - /** - * 按 rawArchiveUri 回查索引记录。默认实现返回 null,允许轻量测试实现不维护该索引。 - * 生产 DuckDB/Parquet 实现必须覆盖,用于 snapshots 的 sourceFrames 反查。 - */ - default EventFileRecord findByRawArchiveUri(String rawArchiveUri) throws IOException { - return null; - } -} diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStoreSink.java b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStoreSink.java index 0b8153c9..10233c7f 100644 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStoreSink.java +++ b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/EventFileStoreSink.java @@ -5,6 +5,8 @@ import com.lingniu.ingest.api.event.RawArchiveKeys; import com.lingniu.ingest.api.event.TelemetrySnapshot; import com.lingniu.ingest.api.event.VehicleEvent; import com.lingniu.ingest.api.event.VehicleEventTelemetrySnapshotMapper; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.api.sink.EventSink; import java.io.IOException; diff --git a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfiguration.java b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfiguration.java index 305d7ba1..c92e5a56 100644 --- a/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfiguration.java +++ b/modules/sinks/event-file-store/src/main/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfiguration.java @@ -2,8 +2,8 @@ package com.lingniu.ingest.eventfilestore.config; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.eventfilestore.DuckDbParquetEventFileStore; -import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.eventfilestore.EventFileStoreSink; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; diff --git a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStoreTest.java b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStoreTest.java index bdb2688e..567d169c 100644 --- a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStoreTest.java +++ b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/DuckDbParquetEventFileStoreTest.java @@ -1,7 +1,10 @@ package com.lingniu.ingest.eventfilestore; -import com.lingniu.ingest.api.ProtocolId; import com.fasterxml.jackson.databind.ObjectMapper; +import com.lingniu.ingest.api.ProtocolId; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/EventFileStoreSinkTest.java b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/EventFileStoreSinkTest.java index 36e37c4a..cc52f074 100644 --- a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/EventFileStoreSinkTest.java +++ b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/EventFileStoreSinkTest.java @@ -7,6 +7,9 @@ import com.lingniu.ingest.api.event.LocationPayload; import com.lingniu.ingest.api.event.RawArchiveKeys; import com.lingniu.ingest.api.event.RealtimePayload; import com.lingniu.ingest.api.event.VehicleEvent; +import com.lingniu.ingest.api.history.EventFileQuery; +import com.lingniu.ingest.api.history.EventFileRecord; +import com.lingniu.ingest.api.history.EventFileStore; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfigurationTest.java b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfigurationTest.java index 55c5b645..66f58f2e 100644 --- a/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfigurationTest.java +++ b/modules/sinks/event-file-store/src/test/java/com/lingniu/ingest/eventfilestore/config/EventFileStoreAutoConfigurationTest.java @@ -1,6 +1,6 @@ package com.lingniu.ingest.eventfilestore.config; -import com.lingniu.ingest.eventfilestore.EventFileStore; +import com.lingniu.ingest.api.history.EventFileStore; import com.lingniu.ingest.eventfilestore.EventFileStoreSink; import org.junit.jupiter.api.Test; import org.springframework.boot.autoconfigure.AutoConfigurations; diff --git a/pom.xml b/pom.xml index 3e549617..1d029be7 100644 --- a/pom.xml +++ b/pom.xml @@ -26,7 +26,6 @@ modules/core/observability modules/sinks/sink-mq modules/sinks/sink-archive - modules/sinks/event-file-store modules/sinks/tdengine-history-store modules/services/event-history-service modules/services/vehicle-stat-service @@ -66,6 +65,12 @@ modules/sinks/raw-archive-store + + optional-event-file-store + + modules/sinks/event-file-store + + legacy-xinda