docs: align module flow with tdengine history
This commit is contained in:
@@ -259,12 +259,12 @@
|
|||||||
<p>
|
<p>
|
||||||
当前项目不是业务库写入服务。它的核心职责是把 GB/T 32960、JT/T 808、JT/T 1078、JSATL12、
|
当前项目不是业务库写入服务。它的核心职责是把 GB/T 32960、JT/T 808、JT/T 1078、JSATL12、
|
||||||
MQTT 等不同来源的车辆数据统一转换为 <code>VehicleEvent</code> 和全字段 <code>TelemetrySnapshot</code>,
|
MQTT 等不同来源的车辆数据统一转换为 <code>VehicleEvent</code> 和全字段 <code>TelemetrySnapshot</code>,
|
||||||
再通过 Kafka、原始报文归档、Parquet 明细文件库、Redis 热状态和独立统计模块交给下游业务。
|
再通过 Kafka、原始报文归档、TDengine 历史索引、Redis 热状态和独立统计模块交给下游业务。
|
||||||
</p>
|
</p>
|
||||||
<div class="summary-grid">
|
<div class="summary-grid">
|
||||||
<div class="summary-card">
|
<div class="summary-card">
|
||||||
<strong>接入方式</strong>
|
<strong>接入方式</strong>
|
||||||
<span>Netty TCP、MQTT、HTTP Push、后续可扩展更多 Inbound Adapter。</span>
|
<span>Netty TCP、MQTT、后续可扩展更多 Inbound Adapter;Xinda Push 不在默认生产面。</span>
|
||||||
</div>
|
</div>
|
||||||
<div class="summary-card">
|
<div class="summary-card">
|
||||||
<strong>统一模型</strong>
|
<strong>统一模型</strong>
|
||||||
@@ -272,7 +272,7 @@
|
|||||||
</div>
|
</div>
|
||||||
<div class="summary-card">
|
<div class="summary-card">
|
||||||
<strong>输出方式</strong>
|
<strong>输出方式</strong>
|
||||||
<span>Kafka Protobuf Envelope + 原始 bytes 冷存 + Parquet/DuckDB + Redis 热状态。</span>
|
<span>Kafka Protobuf Envelope + 原始 bytes 冷存 + TDengine raw_frames/locations + Redis 热状态。</span>
|
||||||
</div>
|
</div>
|
||||||
<div class="summary-card">
|
<div class="summary-card">
|
||||||
<strong>业务边界</strong>
|
<strong>业务边界</strong>
|
||||||
@@ -307,7 +307,7 @@ flowchart TB
|
|||||||
subgraph protocolBase["核心与公共能力 modules/core"]
|
subgraph protocolBase["核心与公共能力 modules/core"]
|
||||||
codec["ingest-codec-common<br/>BCD / CRC / BCC / bit 工具"]
|
codec["ingest-codec-common<br/>BCD / CRC / BCC / bit 工具"]
|
||||||
session["session-core<br/>设备会话、命令分发接口、会话存储"]
|
session["session-core<br/>设备会话、命令分发接口、会话存储"]
|
||||||
identity["vehicle-identity<br/>跨协议身份解析、外部标识绑定、文件持久化"]
|
identity["vehicle-identity<br/>跨协议身份解析、MySQL 绑定表"]
|
||||||
obs["observability<br/>metrics / tracing / health"]
|
obs["observability<br/>metrics / tracing / health"]
|
||||||
api["ingest-api<br/>ProtocolId / RawFrame / VehicleEvent / 注解 SPI"]
|
api["ingest-api<br/>ProtocolId / RawFrame / VehicleEvent / 注解 SPI"]
|
||||||
registry["HandlerRegistry<br/>按协议、命令、infoType 路由"]
|
registry["HandlerRegistry<br/>按协议、命令、infoType 路由"]
|
||||||
@@ -318,13 +318,14 @@ flowchart TB
|
|||||||
subgraph sink["输出与明细存储 modules/sinks"]
|
subgraph sink["输出与明细存储 modules/sinks"]
|
||||||
mq["sink-mq<br/>VehicleEvent 到 Protobuf Envelope 到 Kafka"]
|
mq["sink-mq<br/>VehicleEvent 到 Protobuf Envelope 到 Kafka"]
|
||||||
archive["sink-archive<br/>RawArchive 到本地/S3/OSS 冷存"]
|
archive["sink-archive<br/>RawArchive 到本地/S3/OSS 冷存"]
|
||||||
fileStore["event-file-store<br/>VehicleEvent/RawArchive 索引到 Parquet,DuckDB 按日查询"]
|
tdengineStore["tdengine-history-store<br/>TDengine raw_frames + locations"]
|
||||||
|
fileStore["event-file-store<br/>Parquet/DuckDB 兼容可选"]
|
||||||
end
|
end
|
||||||
|
|
||||||
subgraph services["消费服务 modules/services"]
|
subgraph services["消费服务 modules/services"]
|
||||||
history["event-history-service<br/>Kafka 全字段事件到历史查询/导出"]
|
history["event-history-service<br/>Kafka event/raw 到 TDengine 历史查询"]
|
||||||
state["vehicle-state-service<br/>Kafka 全字段事件到 Redis 热状态"]
|
state["vehicle-state-service<br/>Kafka 全字段事件到 Redis 热状态"]
|
||||||
stat["vehicle-stat-service<br/>Kafka 全字段事件到可配置日统计"]
|
stat["vehicle-stat-service<br/>Kafka 全字段事件到指标表"]
|
||||||
end
|
end
|
||||||
|
|
||||||
subgraph apps["应用入口 modules/apps"]
|
subgraph apps["应用入口 modules/apps"]
|
||||||
@@ -356,10 +357,11 @@ flowchart TB
|
|||||||
dispatcher --> bus
|
dispatcher --> bus
|
||||||
bus --> mq
|
bus --> mq
|
||||||
bus --> archive
|
bus --> archive
|
||||||
bus --> fileStore
|
|
||||||
mq --> history
|
mq --> history
|
||||||
mq --> state
|
mq --> state
|
||||||
mq --> stat
|
mq --> stat
|
||||||
|
history --> tdengineStore
|
||||||
|
history -.兼容可选.-> fileStore
|
||||||
|
|
||||||
gateway --> session
|
gateway --> session
|
||||||
session --> jt808
|
session --> jt808
|
||||||
@@ -375,7 +377,6 @@ flowchart TB
|
|||||||
bootall -.装配.-> mqtt
|
bootall -.装配.-> mqtt
|
||||||
bootall -.装配.-> mq
|
bootall -.装配.-> mq
|
||||||
bootall -.装配.-> archive
|
bootall -.装配.-> archive
|
||||||
bootall -.装配.-> fileStore
|
|
||||||
bootall -.装配.-> history
|
bootall -.装配.-> history
|
||||||
bootall -.装配.-> state
|
bootall -.装配.-> state
|
||||||
bootall -.装配.-> stat
|
bootall -.装配.-> stat
|
||||||
@@ -388,7 +389,7 @@ flowchart TB
|
|||||||
|
|
||||||
class gb,jt808,jt1078,jsatl12,mqtt,mqttProfile entry;
|
class gb,jt808,jt1078,jsatl12,mqtt,mqttProfile entry;
|
||||||
class api,registry,dispatcher,bus,identity core;
|
class api,registry,dispatcher,bus,identity core;
|
||||||
class mq,archive,fileStore sink;
|
class mq,archive,tdengineStore,fileStore sink;
|
||||||
class gateway,session command;
|
class gateway,session command;
|
||||||
class codec,obs,bootall support;
|
class codec,obs,bootall support;
|
||||||
</pre>
|
</pre>
|
||||||
@@ -408,17 +409,18 @@ flowchart TB
|
|||||||
sequenceDiagram
|
sequenceDiagram
|
||||||
autonumber
|
autonumber
|
||||||
participant Device as 车辆/平台
|
participant Device as 车辆/平台
|
||||||
participant Inbound as Inbound Adapter<br/>Netty/MQTT/Push
|
participant Inbound as Inbound Adapter<br/>Netty/MQTT
|
||||||
participant Decoder as FrameDecoder<br/>MessageDecoder
|
participant Decoder as FrameDecoder<br/>MessageDecoder
|
||||||
participant Dispatcher as Dispatcher
|
participant Dispatcher as Dispatcher
|
||||||
participant Handler as Protocol Handler<br/>Mapper
|
participant Handler as Protocol Handler<br/>Mapper
|
||||||
participant EventBus as DisruptorEventBus
|
participant EventBus as DisruptorEventBus
|
||||||
participant MQ as sink-mq<br/>Kafka Envelope
|
participant MQ as sink-mq<br/>Kafka Envelope
|
||||||
participant Archive as sink-archive<br/>ArchiveStore
|
participant Archive as sink-archive<br/>ArchiveStore
|
||||||
participant FileStore as event-file-store<br/>Parquet + DuckDB
|
participant History as vehicle-history-app
|
||||||
|
participant TDengine as tdengine-history-store<br/>raw_frames + locations
|
||||||
participant Consumer as 下游业务消费者
|
participant Consumer as 下游业务消费者
|
||||||
|
|
||||||
Device->>Inbound: 原始报文 bytes / MQTT message / HTTP payload
|
Device->>Inbound: 原始报文 bytes / MQTT message
|
||||||
Inbound->>Decoder: 粘包拆帧、校验、协议体解析
|
Inbound->>Decoder: 粘包拆帧、校验、协议体解析
|
||||||
Inbound->>Inbound: MQTT 按 endpoint profile 选择厂商解析器
|
Inbound->>Inbound: MQTT 按 endpoint profile 选择厂商解析器
|
||||||
Decoder->>Dispatcher: RawFrame(protocol, command, infoType, payload, rawBytes)
|
Decoder->>Dispatcher: RawFrame(protocol, command, infoType, payload, rawBytes)
|
||||||
@@ -429,11 +431,11 @@ sequenceDiagram
|
|||||||
Dispatcher->>Handler: 绑定 rawArchiveKey / rawArchiveUri
|
Dispatcher->>Handler: 绑定 rawArchiveKey / rawArchiveUri
|
||||||
Handler->>EventBus: VehicleEvent(Realtime / Location / Alarm / Login ...)
|
Handler->>EventBus: VehicleEvent(Realtime / Location / Alarm / Login ...)
|
||||||
EventBus->>MQ: 投递规整后的业务事件
|
EventBus->>MQ: 投递规整后的业务事件
|
||||||
EventBus->>FileStore: 写入解析后明细事件
|
|
||||||
MQ->>MQ: EnvelopeMapper 转 Protobuf
|
MQ->>MQ: EnvelopeMapper 转 Protobuf
|
||||||
MQ->>Consumer: Kafka topic,key=vin,单车有序
|
MQ->>History: Kafka event/raw topic,key=vin,单车有序
|
||||||
FileStore->>Consumer: 按日期范围正序/倒序查询和导出
|
History->>TDengine: 写 raw_frames、locations 和 raw parsed JSON
|
||||||
Consumer->>Consumer: 历史明细查询、Redis 热状态、可配置日统计、告警工单
|
TDengine->>Consumer: 分页历史查询、原始帧追溯、位置查询
|
||||||
|
MQ->>Consumer: Redis 热状态、指标统计、告警工单等下游消费
|
||||||
</pre>
|
</pre>
|
||||||
</div>
|
</div>
|
||||||
<div class="callout">
|
<div class="callout">
|
||||||
@@ -526,9 +528,9 @@ flowchart LR
|
|||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
<td>modules/core/vehicle-identity</td>
|
<td>modules/core/vehicle-identity</td>
|
||||||
<td>跨协议车辆身份解析,维护 phone、deviceId、plate 到 VIN 的绑定关系;支持 memory 和 append-only JSONL file 后端,避免重启丢失绑定。</td>
|
<td>跨协议车辆身份解析,维护 phone、deviceId、plate 到 VIN 的绑定关系;生产运行使用 MySQL 绑定表,避免重启丢失绑定并保证多实例一致。</td>
|
||||||
<td>协议模块传入的外部标识。</td>
|
<td>协议模块传入的外部标识。</td>
|
||||||
<td>稳定 VIN、解析来源、是否已命中绑定。</td>
|
<td>稳定 VIN、解析来源、是否已命中绑定、MySQL 绑定记录。</td>
|
||||||
<td><span class="tag red">统计主键</span><span class="tag">多协议归一</span><span class="tag green">持久化绑定</span></td>
|
<td><span class="tag red">统计主键</span><span class="tag">多协议归一</span><span class="tag green">持久化绑定</span></td>
|
||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
@@ -589,17 +591,24 @@ flowchart LR
|
|||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
<td>modules/sinks/event-file-store</td>
|
<td>modules/sinks/event-file-store</td>
|
||||||
<td>文件型明细库,使用 Parquet 存储解析后的事件明细和 RawArchive 可查询索引,使用 DuckDB 按日期范围正序/倒序查询;RawArchive 只存 archive:// 逻辑 URI、命令码、长度和 metadata,不嵌入原始 bytes,原始 bytes 本体仍由 sink-archive 管理;读取历史文件时未知 protocol 会降级为 UNKNOWN,避免查询中断。</td>
|
<td>兼容文件型明细库,使用 Parquet 存储解析后的事件明细和 RawArchive 可查询索引,使用 DuckDB 按日期范围正序/倒序查询;默认生产部署不启用,仅用于旧低层记录接口或本地调查场景。</td>
|
||||||
<td>VehicleEvent 或 EventFileRecord。</td>
|
<td>VehicleEvent 或 EventFileRecord。</td>
|
||||||
<td>protocol/date 分区的 Parquet 文件和查询结果。</td>
|
<td>protocol/date 分区的 Parquet 文件和兼容查询结果。</td>
|
||||||
<td><span class="tag green">按日查询</span><span class="tag green">导出</span><span class="tag">字段配置</span></td>
|
<td><span class="tag">兼容查询</span><span class="tag amber">默认关闭</span></td>
|
||||||
|
</tr>
|
||||||
|
<tr>
|
||||||
|
<td>modules/sinks/tdengine-history-store</td>
|
||||||
|
<td>生产历史库边界,负责 TDengine schema、raw_frames/location/telemetry field 行映射、批量写入和分页查询语句;raw_frames 保留完整 payloadJson.parsed,位置表保持轻量字段并通过 rawUri 关联原始帧。</td>
|
||||||
|
<td>Kafka Protobuf Envelope、Raw Envelope。</td>
|
||||||
|
<td>TDengine <code>raw_frames</code>、位置表、可选 telemetry field 表。</td>
|
||||||
|
<td><span class="tag green">生产历史</span><span class="tag green">分页查询</span><span class="tag amber">raw 追溯</span></td>
|
||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
<td>modules/services/event-history-service</td>
|
<td>modules/services/event-history-service</td>
|
||||||
<td>消费 Kafka 全字段事件,写入文件型明细库,并提供历史查询和 CSV 导出边界;导出支持通过 <code>fields</code> 参数选择内部字段并平铺为 CSV 列,也支持 <code>metadata.xxx</code> 导出协议侧诊断、媒体归档、报警类型等 metadata 字段;未知 Envelope source 会进入 UNKNOWN 协议分区,并在 metadata 保留 <code>originalSource</code> 便于回溯;生产消费可使用 <code>tryIngest</code> 和自动装配的 <code>EnvelopeConsumerProcessor</code> 将坏 protobuf、缺快照和存储异常收敛为结构化结果,并通过 <code>EnvelopeDeadLetterSink</code> 发布死信,避免阻塞 Kafka 分区。</td>
|
<td>消费 32960、808、宇通 MQTT 的 Kafka event/raw topic,生产默认写入 TDengine raw_frames 和位置表,并提供 raw 帧、位置历史等分页查询边界;Raw 查询返回完整 parsed JSON,位置历史只保留核心字段并通过 rawUri 关联原始帧;生产消费可使用 <code>tryIngest</code> 和自动装配的 <code>EnvelopeConsumerProcessor</code> 将坏 protobuf、缺快照和存储异常收敛为结构化结果,并通过 <code>EnvelopeDeadLetterSink</code> 发布死信,避免阻塞 Kafka 分区。</td>
|
||||||
<td>Kafka Protobuf Envelope。</td>
|
<td>Kafka Protobuf Envelope。</td>
|
||||||
<td>EventFileRecord、查询结果、可配置字段 CSV。</td>
|
<td>TDengine raw_frames、位置分页结果、raw archive 逻辑引用。</td>
|
||||||
<td><span class="tag green">历史明细</span><span class="tag green">多日导出</span><span class="tag">字段配置</span></td>
|
<td><span class="tag green">历史明细</span><span class="tag green">分页查询</span><span class="tag amber">raw JSON</span></td>
|
||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
<td>modules/services/vehicle-state-service</td>
|
<td>modules/services/vehicle-state-service</td>
|
||||||
@@ -681,8 +690,8 @@ flowchart LR
|
|||||||
</tr>
|
</tr>
|
||||||
<tr>
|
<tr>
|
||||||
<td>明细展示和导出</td>
|
<td>明细展示和导出</td>
|
||||||
<td><code>event-file-store</code></td>
|
<td><code>vehicle-history-app</code> + <code>tdengine-history-store</code></td>
|
||||||
<td>解析后的事件明细和 RawArchive 索引按 protocol/date 写 Parquet。DuckDB 查询时按日期范围读取,并按 event_time 正序或倒序返回。</td>
|
<td>生产查询通过 TDengine raw_frames、位置表和 rawUri 关联完成;Parquet/DuckDB 仅作为兼容文件型明细库,不进入默认热路径。</td>
|
||||||
</tr>
|
</tr>
|
||||||
</tbody>
|
</tbody>
|
||||||
</table>
|
</table>
|
||||||
|
|||||||
Reference in New Issue
Block a user