From b36f993059e445d9d921bf6132210d33d2bbf904 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 04:06:41 +0800 Subject: [PATCH] docs: align comments with tdengine history path --- .../java/com/lingniu/ingest/api/event/RawArchiveKeys.java | 4 ++-- .../com/lingniu/ingest/api/event/TelemetryFieldValue.java | 4 ++-- .../java/com/lingniu/ingest/api/event/TelemetrySnapshot.java | 4 ++-- .../lingniu/ingest/core/concurrency/DisruptorEventBus.java | 2 +- .../ingest/core/config/IngestCoreAutoConfiguration.java | 2 +- .../ingest/core/pipeline/builtin/DedupInterceptor.java | 2 +- .../protocol/gb32960/handler/Gb32960RealtimeHandler.java | 2 +- .../protocol/gb32960/inbound/Gb32960ChannelHandler.java | 2 +- .../ingest/protocol/jt808/inbound/Jt808ChannelHandler.java | 2 +- .../lingniu/ingest/vehiclestate/VehicleStateController.java | 2 +- 10 files changed, 13 insertions(+), 13 deletions(-) diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/RawArchiveKeys.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/RawArchiveKeys.java index f2f14a7b..014098d9 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/RawArchiveKeys.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/RawArchiveKeys.java @@ -11,8 +11,8 @@ import java.util.Map; * Shared raw-archive key conventions used by Dispatcher, archive sink, and query services. * *

32960 生产链路里,RAW .bin 文件本体按 {@code 日期/协议/VIN/eventId.bin} 落在 archive 根目录, - * DuckDB 只保存 {@code archive://...} 引用和查询索引。这个类集中维护 key 规则,避免接收端、落盘端、 - * snapshot 查询端对目录分区的理解不一致。

+ * TDengine raw_frames 或兼容索引只保存 {@code archive://...} 引用。这个类集中维护 key 规则, + * 避免接收端、落盘端、snapshot 查询端对目录分区的理解不一致。

*/ public final class RawArchiveKeys { diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetryFieldValue.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetryFieldValue.java index b6fec6b5..d101ad25 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetryFieldValue.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetryFieldValue.java @@ -3,11 +3,11 @@ package com.lingniu.ingest.api.event; /** * One normalized telemetry field value. * - *

Protocol mappers use stable internal field keys so Kafka, Parquet, Redis, + *

Protocol mappers use stable internal field keys so Kafka, TDengine, Redis, * and statistics do not depend on protocol-specific names. * *

{@code value} 统一用字符串承载,真实类型由 {@code valueType} 表达。这样 Kafka - * protobuf、CSV 导出、DuckDB JSON 字段和前端展示可以共用同一份字段结构;数值精度/格式化 + * protobuf、CSV 导出、JSON 字段和前端展示可以共用同一份字段结构;数值精度/格式化 * 由查询层或前端按 {@code valueType + unit} 决定。 */ public record TelemetryFieldValue( diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetrySnapshot.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetrySnapshot.java index 555b5227..2035eb55 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetrySnapshot.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/event/TelemetrySnapshot.java @@ -12,9 +12,9 @@ import java.util.Optional; /** * Normalized full-field telemetry snapshot for one parsed vehicle event. * - *

This is the shared contract for Kafka, Parquet history, Redis hot state, + *

This is the shared contract for Kafka, TDengine history, Redis hot state, * and statistics. It intentionally lives in {@code ingest-api} and has no - * Spring, Kafka, Redis, or DuckDB dependency. + * Spring, Kafka, Redis, or TDengine dependency. */ public record TelemetrySnapshot( String eventId, diff --git a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java index 92760e25..8e084cd6 100644 --- a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java +++ b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/concurrency/DisruptorEventBus.java @@ -52,7 +52,7 @@ public final class DisruptorEventBus implements AutoCloseable { EventHandler[] handlers = this.sinks.stream() .map(this::toHandler) .toArray(EventHandler[]::new); - // 每个 sink 一个独立 handler,事件按扇出模式同时写 Kafka、event-file-store 等目标。 + // 每个 sink 一个独立 handler,事件按扇出模式同时写 Kafka、archive、历史索引等目标。 disruptor.handleEventsWith(handlers); disruptor.start(); log.info("DisruptorEventBus started ringBuffer={} wait={} sinks={}", diff --git a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/config/IngestCoreAutoConfiguration.java b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/config/IngestCoreAutoConfiguration.java index 3249e85f..d3de0447 100644 --- a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/config/IngestCoreAutoConfiguration.java +++ b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/config/IngestCoreAutoConfiguration.java @@ -24,7 +24,7 @@ import java.util.List; * *

生产链路顺序:入口收到 {@code RawFrame} → {@link InterceptorChain} 做准入 → * {@link Dispatcher} 定位协议 Handler → Handler 产出 {@code VehicleEvent} → - * {@link DisruptorEventBus} 扇出到 Archive、DuckDB、Kafka 等 {@code EventSink}。 + * {@link DisruptorEventBus} 扇出到 Archive、Kafka、历史索引等 {@code EventSink}。 */ @AutoConfiguration @EnableConfigurationProperties(IngestCoreProperties.class) diff --git a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/pipeline/builtin/DedupInterceptor.java b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/pipeline/builtin/DedupInterceptor.java index ecffb669..7c4f72f4 100644 --- a/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/pipeline/builtin/DedupInterceptor.java +++ b/modules/core/ingest-core/src/main/java/com/lingniu/ingest/core/pipeline/builtin/DedupInterceptor.java @@ -18,7 +18,7 @@ import java.util.Map; * 优先使用已解析 VIN;VIN 未解析时使用 phone/deviceId/terminalId/plate,避免 JT808 多终端在 * {@code vin=unknown} 时同流水号互相去重。如果上游已经在 {@link RawFrame#sourceMeta()} 里带了 * {@code seq},优先用真实流水号;否则退化为 raw bytes fingerprint,避免设备重连或 TCP 重发导致 - * 同一帧重复进入 DuckDB/Archive。 + * 同一帧重复进入 Archive/Kafka/历史索引。 * *

当前缓存是进程内的,只能保证单实例去重。32960 若做多副本水平扩容,需要把这个位置替换成 * Redis/集中式幂等键,否则不同实例仍可能各自接收一次相同原始包。 diff --git a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/handler/Gb32960RealtimeHandler.java b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/handler/Gb32960RealtimeHandler.java index 87d25919..b0480d15 100644 --- a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/handler/Gb32960RealtimeHandler.java +++ b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/handler/Gb32960RealtimeHandler.java @@ -18,7 +18,7 @@ import java.util.List; *

这是 GB32960 协议在通用 Dispatcher 管线里的**唯一入口** Bean——移除将导致所有 * GB32960 帧产出的事件丢失。构造器注入 {@link Gb32960EventMapper},全部业务逻辑限制 * 在纯函数里,不触碰数据库 / 不做 IO。产出的事件由 Dispatcher 统一投递到 Disruptor; - * 下游 sink 再决定写 raw archive、DuckDB/Parquet 历史库或转发 Kafka。 + * 下游 sink 再决定写 raw archive、TDengine 历史库或转发 Kafka。 * *

当前 32960 专用 snapshot 查询不依赖 Realtime/Location 事件表,而是通过 raw archive URI * 回读原始 .bin 并即时解码合并,避免不同车辆或不同子包之间做大范围 join。 diff --git a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java index 437a5539..8df345c4 100644 --- a/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java +++ b/modules/protocols/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java @@ -29,7 +29,7 @@ import java.util.Map; * Netty 入站处理器:字节 → Gb32960Message → 认证 → RawFrame → Dispatcher。 * *

这是 32960 生产链路的入口边界。它只做协议层职责:解码、鉴权、应答、组装 - * {@link RawFrame} 并投递 Dispatcher;原始包归档、DuckDB/Parquet 索引、Kafka 转发都在 + * {@link RawFrame} 并投递 Dispatcher;原始包归档、TDengine 历史索引、Kafka 转发都在 * 下游 sink 中完成,避免 Netty EventLoop 被存储 IO 阻塞。 * *

认证逻辑: diff --git a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java index efc26b63..f849f288 100644 --- a/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java +++ b/modules/protocols/protocol-jt808/src/main/java/com/lingniu/ingest/protocol/jt808/inbound/Jt808ChannelHandler.java @@ -111,7 +111,7 @@ public class Jt808ChannelHandler extends SimpleChannelInboundHandler { processDecodedFrame(ctx, frame, msg); } catch (RuntimeException e) { // 解码成功但业务处理失败时仍派发一条 malformed RawFrame, - // 保证原始帧能被 archive/event-file-store 留痕,方便事后排查。 + // 保证原始帧能被 archive、Kafka 和历史索引留痕,方便事后排查。 log.warn("[jt808] processing failed peer={} phone={} msgId=0x{} error={}", addr(ctx), msg.header().phone(), Integer.toHexString(msg.header().messageId()), e.getMessage()); log.debug("[jt808] processing failed stack peer={}", addr(ctx), e); diff --git a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateController.java b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateController.java index ecb9f973..507114c3 100644 --- a/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateController.java +++ b/modules/services/vehicle-state-service/src/main/java/com/lingniu/ingest/vehiclestate/VehicleStateController.java @@ -28,7 +28,7 @@ public final class VehicleStateController { @GetMapping("/{vin}") public ResponseEntity state(@PathVariable String vin) { - // 查询的是 Redis 最新状态快照,不是 event-file-store 的历史 snapshot。 + // 查询的是 Redis 最新状态快照,不是历史明细库。 return json(repository.getState(vin)); }