mirror of
https://gitee.com/san-bing/JChargePointProtocol
synced 2026-07-20 13:47:51 +08:00
wiki生成
This commit is contained in:
638
docs/核心模块详解/协议实现模块/云快充协议实现/上行消息处理.md
Normal file
638
docs/核心模块详解/协议实现模块/云快充协议实现/上行消息处理.md
Normal file
@@ -0,0 +1,638 @@
|
||||
# 上行消息处理
|
||||
|
||||
<cite>
|
||||
**本文档引用的文件**
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java)
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java)
|
||||
- [YunKuaiChongUplinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongUplinkCmdExe.java)
|
||||
- [YunKuaiChongUplinkMessage.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongUplinkMessage.java)
|
||||
- [KafkaForwarder.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/forwarder/KafkaForwarder.java)
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java)
|
||||
- [YunKuaiChongV150LoginULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150LoginULCmd.java)
|
||||
- [YunKuaiChongV150StartChargeULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150StartChargeULCmd.java)
|
||||
- [YunKuaiChongV150RealTimeDataULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RealTimeDataULCmd.java)
|
||||
- [YunKuaiChongV150TransactionRecordULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150TransactionRecordULCmd.java)
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java)
|
||||
</cite>
|
||||
|
||||
## 目录
|
||||
|
||||
1. [概述](#概述)
|
||||
2. [系统架构](#系统架构)
|
||||
3. [消息处理流程](#消息处理流程)
|
||||
4. [协议命令注解机制](#协议命令注解机制)
|
||||
5. [关键上行命令详解](#关键上行命令详解)
|
||||
6. [消息转发机制](#消息转发机制)
|
||||
7. [数据持久化与状态管理](#数据持久化与状态管理)
|
||||
8. [性能优化策略](#性能优化策略)
|
||||
9. [总结](#总结)
|
||||
|
||||
## 概述
|
||||
|
||||
云快充协议上行消息处理系统是一个高度模块化的分布式消息处理框架,负责接收、解析、分发和处理来自充电桩的各种上行消息。该系统采用事件驱动架构,通过反射机制和命令路由模式实现灵活的消息处理。
|
||||
|
||||
核心特性包括:
|
||||
|
||||
- **高性能TCP监听**:基于Netty的异步I/O处理
|
||||
- **智能协议解析**:支持多种协议版本的统一处理
|
||||
- **反射命令路由**:基于注解的自动化命令分发
|
||||
- **消息转发队列**:通过Kafka实现异步消息传递
|
||||
- **状态管理**:实时更新充电桩状态和属性
|
||||
|
||||
## 系统架构
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "网络层"
|
||||
TCP[TCP监听器]
|
||||
Session[TCP会话管理]
|
||||
end
|
||||
subgraph "协议处理层"
|
||||
Processor[云快充协议处理器]
|
||||
Router[命令路由器]
|
||||
Message[消息解析器]
|
||||
end
|
||||
subgraph "命令执行层"
|
||||
Heartbeat[心跳命令]
|
||||
Login[登录命令]
|
||||
Charge[充电命令]
|
||||
RealTime[实时数据]
|
||||
Transaction[交易记录]
|
||||
end
|
||||
subgraph "消息转发层"
|
||||
Forwarder[Kafka转发器]
|
||||
Queue[消息队列]
|
||||
end
|
||||
subgraph "应用服务层"
|
||||
App[jcpp-app服务]
|
||||
DB[数据库存储]
|
||||
end
|
||||
TCP --> Session
|
||||
Session --> Processor
|
||||
Processor --> Router
|
||||
Router --> Message
|
||||
Message --> Heartbeat
|
||||
Message --> Login
|
||||
Message --> Charge
|
||||
Message --> RealTime
|
||||
Message --> Transaction
|
||||
Heartbeat --> Forwarder
|
||||
Login --> Forwarder
|
||||
Charge --> Forwarder
|
||||
RealTime --> Forwarder
|
||||
Transaction --> Forwarder
|
||||
Forwarder --> Queue
|
||||
Queue --> App
|
||||
App --> DB
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L27-L61)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L25-L40)
|
||||
|
||||
## 消息处理流程
|
||||
|
||||
### TCP监听器接收数据
|
||||
|
||||
系统从TcpListener接收到原始数据包开始,首先进行基本的完整性检查:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([接收数据包]) --> CheckLength{数据包长度检查}
|
||||
CheckLength --> |长度不足| Reject[拒绝处理]
|
||||
CheckLength --> |长度足够| CheckHeader{头部标识检查}
|
||||
CheckHeader --> |头部无效| Reject
|
||||
CheckHeader --> |头部有效| ParseHeader[解析协议头]
|
||||
ParseHeader --> CheckBoundary{边界检查}
|
||||
CheckBoundary --> |边界无效| Reject
|
||||
CheckBoundary --> |边界有效| ParseFields[解析字段]
|
||||
ParseFields --> CheckCRC{校验和验证}
|
||||
CheckCRC --> |校验失败| Reject
|
||||
CheckCRC --> |校验成功| BuildMessage[构建消息对象]
|
||||
BuildMessage --> RouteCommand[路由命令]
|
||||
RouteCommand --> ExecuteCommand[执行命令]
|
||||
ExecuteCommand --> End([处理完成])
|
||||
Reject --> End
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L63-L154)
|
||||
|
||||
### 消息解析与验证
|
||||
|
||||
协议处理器执行严格的验证流程:
|
||||
|
||||
1. **长度验证**:确保数据包包含最小长度(8字节)
|
||||
2. **头部验证**:检查起始标识符(0x68)
|
||||
3. **边界检查**:验证数据长度字段的合理性
|
||||
4. **CRC校验**:双重校验和验证(小端和大端)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L63-L154)
|
||||
|
||||
### 命令路由与分发
|
||||
|
||||
解析完成后,消息通过命令路由器进行分发:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Processor as 协议处理器
|
||||
participant Router as 命令路由器
|
||||
participant Executor as 命令执行器
|
||||
participant Forwarder as 消息转发器
|
||||
Processor->>Router : 查找命令执行器
|
||||
Router->>Router : 构建路由键(protocol : cmd)
|
||||
Router->>Router : 查询执行器映射表
|
||||
Router-->>Processor : 返回执行器实例
|
||||
Processor->>Executor : 调用execute方法
|
||||
Executor->>Executor : 处理业务逻辑
|
||||
Executor->>Forwarder : 发送消息到队列
|
||||
Forwarder-->>Executor : 确认发送
|
||||
Executor-->>Processor : 返回处理结果
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L170-L184)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L85-L95)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L170-L184)
|
||||
|
||||
## 协议命令注解机制
|
||||
|
||||
### @ProtocolCmd注解设计
|
||||
|
||||
系统采用基于注解的命令映射机制,通过`@ProtocolCmd`注解实现声明式命令注册:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ProtocolCmd {
|
||||
+int value()
|
||||
+String[] protocolNames()
|
||||
}
|
||||
class ProtocolCommandRouter {
|
||||
-Map~String,T~ executorMap
|
||||
+ProtocolCommandRouter(Class, Predicate)
|
||||
+getExecutor(String, int) T
|
||||
-initializeRoutes(Class, Predicate)
|
||||
-registerExecutor(Class)
|
||||
}
|
||||
class YunKuaiChongUplinkCmdExe {
|
||||
+execute(TcpSession, YunKuaiChongUplinkMessage, ProtocolContext)
|
||||
#uplinkMessageBuilder(String, TcpSession, YunKuaiChongUplinkMessage)
|
||||
}
|
||||
ProtocolCmd --> YunKuaiChongUplinkCmdExe : 注解标记
|
||||
ProtocolCommandRouter --> YunKuaiChongUplinkCmdExe : 管理实例
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java#L20-L32)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L25-L40)
|
||||
|
||||
### 反射机制工作原理
|
||||
|
||||
命令路由器通过以下步骤实现反射机制:
|
||||
|
||||
1. **类扫描**:扫描指定包路径下带有`@ProtocolCmd`注解的类
|
||||
2. **注解解析**:提取命令字和协议名称数组
|
||||
3. **实例化**:动态创建命令执行器实例
|
||||
4. **路由注册**:构建`protocol:cmd`键值对存储到映射表
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L42-L85)
|
||||
|
||||
### 多版本协议支持
|
||||
|
||||
系统通过协议名称数组支持多个协议版本的统一处理:
|
||||
|
||||
| 命令类型 | 协议版本 | 命令字 | 功能描述 |
|
||||
|------|----------------|------|---------|
|
||||
| 心跳命令 | V150/V160/V170 | 0x03 | 充电桩状态报告 |
|
||||
| 登录命令 | V150/V160/V170 | 0x01 | 充电桩认证注册 |
|
||||
| 充电启动 | V150/V160/V170 | 0x31 | 主动充电请求 |
|
||||
| 实时数据 | V150/V160/V170 | 0x13 | 充电过程监控 |
|
||||
| 交易记录 | V150/V160 | 0x3B | 充电完成记录 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java#L25-L40)
|
||||
|
||||
## 关键上行命令详解
|
||||
|
||||
### 心跳包处理
|
||||
|
||||
心跳命令(0x03)负责维持连接状态和报告充电桩状态:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Pile as 充电桩
|
||||
participant Processor as 协议处理器
|
||||
participant HeartbeatCmd as 心跳命令
|
||||
participant Session as 会话管理
|
||||
participant Forwarder as 消息转发器
|
||||
Pile->>Processor : 发送心跳数据
|
||||
Processor->>HeartbeatCmd : 路由到心跳处理器
|
||||
HeartbeatCmd->>HeartbeatCmd : 解析桩编号和枪状态
|
||||
HeartbeatCmd->>Session : 刷新会话激活状态
|
||||
HeartbeatCmd->>HeartbeatCmd : 构建心跳请求消息
|
||||
HeartbeatCmd->>Forwarder : 发送到Kafka队列
|
||||
HeartbeatCmd->>Pile : 发送心跳响应
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java#L30-L84)
|
||||
|
||||
**业务逻辑特点**:
|
||||
|
||||
- **状态更新**:刷新会话激活时间
|
||||
- **响应机制**:自动回复心跳确认
|
||||
- **信息提取**:解析桩编号、枪号、状态等关键信息
|
||||
- **转发处理**:将心跳信息发送到后端服务
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java#L30-L84)
|
||||
|
||||
### 登录请求处理
|
||||
|
||||
登录命令(0x01)负责充电桩的身份认证和会话建立:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
LoginReq[登录请求] --> ParseInfo[解析设备信息]
|
||||
ParseInfo --> ExtractPileCode[提取桩编号]
|
||||
ExtractPileCode --> ExtractType[提取桩类型]
|
||||
ExtractType --> ExtractGuns[提取枪数量]
|
||||
ExtractGuns --> ExtractVersion[提取软件版本]
|
||||
ExtractVersion --> ExtractSIM[提取SIM卡信息]
|
||||
ExtractSIM --> RegisterSession[注册会话]
|
||||
RegisterSession --> BuildLoginMsg[构建登录消息]
|
||||
BuildLoginMsg --> ForwardToBackend[转发到后端]
|
||||
ForwardToBackend --> Complete[处理完成]
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150LoginULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150LoginULCmd.java#L30-L84)
|
||||
|
||||
**认证流程**:
|
||||
|
||||
1. **身份验证**:通过桩编号验证设备合法性
|
||||
2. **会话注册**:建立和维护设备会话状态
|
||||
3. **信息收集**:收集设备硬件和软件信息
|
||||
4. **后端通知**:通知应用服务新设备上线
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150LoginULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150LoginULCmd.java#L30-L84)
|
||||
|
||||
### 启动充电处理
|
||||
|
||||
充电启动命令(0x31)处理充电桩主动发起的充电请求:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
StartCharge[充电启动请求] --> ParseBasic[解析基础信息]
|
||||
ParseBasic --> ParseCard[解析卡片信息]
|
||||
ParseCard --> ParseVIN[解析VIN码]
|
||||
ParseVIN --> ValidatePassword{需要密码?}
|
||||
ValidatePassword --> |是| HashPassword[密码哈希处理]
|
||||
ValidatePassword --> |否| SkipPassword[跳过密码处理]
|
||||
HashPassword --> ReverseVIN[反转VIN码]
|
||||
SkipPassword --> ReverseVIN
|
||||
ReverseVIN --> BuildRequest[构建启动请求]
|
||||
BuildRequest --> ForwardMessage[转发消息]
|
||||
ForwardMessage --> Complete[处理完成]
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150StartChargeULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150StartChargeULCmd.java#L35-L124)
|
||||
|
||||
**处理要点**:
|
||||
|
||||
- **卡片验证**:处理物理卡片号和密码验证
|
||||
- **VIN码处理**:特殊格式的VIN码反转处理
|
||||
- **启动方式**:支持多种启动方式(APP、卡片、离线卡等)
|
||||
- **安全机制**:密码MD5哈希处理
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150StartChargeULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150StartChargeULCmd.java#L35-L124)
|
||||
|
||||
### 实时数据处理
|
||||
|
||||
实时数据命令(0x13)负责传输充电过程中的监控数据:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class RealTimeDataProcessor {
|
||||
+execute(TcpSession, YunKuaiChongUplinkMessage, ProtocolContext)
|
||||
-parseGunStatus(int, int, String) GunRunStatus
|
||||
-parseFaults(byte[]) boolean[]
|
||||
-getFaultDescriptions(boolean[]) String[]
|
||||
-reduceMagnification(long, int) BigDecimal
|
||||
}
|
||||
class FaultDescription {
|
||||
+急停按钮动作故障
|
||||
+无可用整流模块
|
||||
+出风口温度过高
|
||||
+交流防雷故障
|
||||
+交直流模块通信中断
|
||||
+绝缘检测模块通信中断
|
||||
+电度表通信中断
|
||||
+读卡器通信中断
|
||||
+RC10通信中断
|
||||
+风扇调速板故障
|
||||
+直流熔断器故障
|
||||
+高压接触器故障
|
||||
+门打开
|
||||
}
|
||||
RealTimeDataProcessor --> FaultDescription : 使用
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150RealTimeDataULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RealTimeDataULCmd.java#L35-L239)
|
||||
|
||||
**数据处理能力**:
|
||||
|
||||
- **状态解析**:解析充电枪的各种运行状态
|
||||
- **故障检测**:14种硬件故障的bit位解析
|
||||
- **数值转换**:精密的数值放大倍率处理
|
||||
- **多消息转发**:同时发送状态和进度消息
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150RealTimeDataULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RealTimeDataULCmd.java#L35-L239)
|
||||
|
||||
### 交易记录处理
|
||||
|
||||
交易记录命令(0x3B)处理充电完成后的结算信息:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
TransactionRecord[交易记录] --> ParseTradeNo[解析交易流水号]
|
||||
ParseTradeNo --> ParsePileGun[解析桩号和枪号]
|
||||
ParsePileGun --> ParseTimes[解析时间信息]
|
||||
ParseTimes --> ParseEnergy[解析电量数据]
|
||||
ParseEnergy --> ParseAmount[解析金额信息]
|
||||
ParseAmount --> ParseMeter[解析电表数据]
|
||||
ParseMeter --> ParseVIN[解析VIN码]
|
||||
ParseVIN --> ParseStartFlag[解析启动方式]
|
||||
ParseStartFlag --> ParseStopReason[解析停止原因]
|
||||
ParseStopReason --> ParseCard[解析卡片信息]
|
||||
ParseCard --> BuildTransaction[构建交易记录]
|
||||
BuildTransaction --> ForwardToBackend[转发到后端]
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150TransactionRecordULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150TransactionRecordULCmd.java#L35-L305)
|
||||
|
||||
**结算信息完整性**:
|
||||
|
||||
- **分时电价**:支持尖峰平谷四种电价模式
|
||||
- **电量统计**:精确到4位小数的电量计量
|
||||
- **费用明细**:详细的电费和服务费计算
|
||||
- **异常处理**:128种充电异常原因的完整映射
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150TransactionRecordULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150TransactionRecordULCmd.java#L35-L305)
|
||||
|
||||
## 消息转发机制
|
||||
|
||||
### Kafka转发器架构
|
||||
|
||||
系统采用Kafka作为消息中间件,实现异步消息传递:
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "消息生产者"
|
||||
Processor[协议处理器]
|
||||
CmdExe[命令执行器]
|
||||
end
|
||||
subgraph "Kafka集群"
|
||||
Broker1[Kafka Broker 1]
|
||||
Broker2[Kafka Broker 2]
|
||||
Broker3[Kafka Broker 3]
|
||||
end
|
||||
subgraph "消息消费者"
|
||||
Consumer1[jcpp-app服务1]
|
||||
Consumer2[jcpp-app服务2]
|
||||
Consumer3[jcpp-app服务3]
|
||||
end
|
||||
Processor --> CmdExe
|
||||
CmdExe --> KafkaForwarder[Kafka转发器]
|
||||
KafkaForwarder --> Broker1
|
||||
KafkaForwarder --> Broker2
|
||||
KafkaForwarder --> Broker3
|
||||
Broker1 --> Consumer1
|
||||
Broker2 --> Consumer2
|
||||
Broker3 --> Consumer3
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [KafkaForwarder.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/forwarder/KafkaForwarder.java#L40-L80)
|
||||
|
||||
### 消息序列化与传输
|
||||
|
||||
转发器支持多种消息格式和传输模式:
|
||||
|
||||
| 传输模式 | 格式类型 | 性能特点 | 使用场景 |
|
||||
|-----------|----------|------|-------|
|
||||
| Monolith | Protobuf | 高性能 | 单体部署 |
|
||||
| Partition | JSON | 易调试 | 分布式部署 |
|
||||
| Custom | 自定义 | 灵活扩展 | 特殊需求 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [KafkaForwarder.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/forwarder/KafkaForwarder.java#L120-L180)
|
||||
|
||||
### 消息路由策略
|
||||
|
||||
系统实现了智能的消息路由机制:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Message[上行消息] --> CheckMode{检查传输模式}
|
||||
CheckMode --> |Monolith| DirectSend[直接发送]
|
||||
CheckMode --> |Partition| PartitionSend[分区发送]
|
||||
CheckMode --> |Custom| CustomSend[自定义发送]
|
||||
DirectSend --> MonolithQueue[单体队列]
|
||||
PartitionSend --> PartitionKey[分区键]
|
||||
PartitionKey --> KafkaTopic[Kafka主题]
|
||||
CustomSend --> CustomLogic[自定义逻辑]
|
||||
MonolithQueue --> Consumer[jcpp-app消费者]
|
||||
KafkaTopic --> Consumer
|
||||
CustomLogic --> Consumer
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [KafkaForwarder.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/forwarder/KafkaForwarder.java#L120-L180)
|
||||
|
||||
## 数据持久化与状态管理
|
||||
|
||||
### 属性系统架构
|
||||
|
||||
系统采用属性(Attribute)系统进行数据持久化:
|
||||
|
||||
```mermaid
|
||||
erDiagram
|
||||
ENTITY {
|
||||
uuid id PK
|
||||
string entity_type
|
||||
timestamp created_at
|
||||
timestamp last_update_ts
|
||||
}
|
||||
ATTRIBUTE {
|
||||
uuid entity_id FK
|
||||
string attr_key
|
||||
string str_value
|
||||
double dbl_value
|
||||
bigint long_value
|
||||
boolean bool_value
|
||||
json json_value
|
||||
timestamp last_update_ts
|
||||
int version
|
||||
}
|
||||
ENTITY ||--o{ ATTRIBUTE : contains
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DefaultAttributeRepository.java](file://jcpp-app/src/main/java/sanbing/jcpp/app/dal/repository/attribute/DefaultAttributeRepository.java#L112-L142)
|
||||
|
||||
### 状态更新流程
|
||||
|
||||
系统实现了高效的状态更新机制:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Cmd as 命令处理器
|
||||
participant AttrSvc as 属性服务
|
||||
participant Repo as 属性仓库
|
||||
participant Queue as 批量队列
|
||||
participant DB as 数据库
|
||||
Cmd->>AttrSvc : 请求状态更新
|
||||
AttrSvc->>AttrSvc : 构建属性列表
|
||||
AttrSvc->>Repo : 添加到批量队列
|
||||
Repo->>Queue : 排队等待
|
||||
Queue->>DB : 批量插入/更新
|
||||
DB-->>Queue : 确认写入
|
||||
Queue-->>Repo : 返回版本号
|
||||
Repo-->>AttrSvc : 返回操作结果
|
||||
AttrSvc-->>Cmd : 返回更新状态
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DefaultAttributeRepository.java](file://jcpp-app/src/main/java/sanbing/jcpp/app/dal/repository/attribute/DefaultAttributeRepository.java#L71-L110)
|
||||
|
||||
### 批量处理优化
|
||||
|
||||
系统采用批量处理策略提升性能:
|
||||
|
||||
| 优化策略 | 实现方式 | 性能提升 |
|
||||
|------|------------------|-----------|
|
||||
| 批量队列 | SqlBlockingQueue | 减少数据库连接开销 |
|
||||
| 异步处理 | Future模式 | 提升并发处理能力 |
|
||||
| 版本控制 | 冲突检测 | 确保数据一致性 |
|
||||
| 缓存机制 | 多级缓存 | 减少重复查询 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [DefaultAttributeRepository.java](file://jcpp-app/src/main/java/sanbing/jcpp/app/dal/repository/attribute/DefaultAttributeRepository.java#L71-L110)
|
||||
|
||||
## 性能优化策略
|
||||
|
||||
### 异步处理架构
|
||||
|
||||
系统采用完全异步的处理架构:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "网络层"
|
||||
Netty[Netty I/O线程]
|
||||
end
|
||||
subgraph "处理层"
|
||||
Processor[协议处理器]
|
||||
Router[命令路由器]
|
||||
Executor[命令执行器]
|
||||
end
|
||||
subgraph "异步层"
|
||||
AsyncPool[异步线程池]
|
||||
Future[Future模式]
|
||||
end
|
||||
subgraph "存储层"
|
||||
DB[数据库]
|
||||
Cache[缓存]
|
||||
end
|
||||
Netty --> Processor
|
||||
Processor --> Router
|
||||
Router --> Executor
|
||||
Executor --> AsyncPool
|
||||
AsyncPool --> Future
|
||||
Future --> DB
|
||||
Future --> Cache
|
||||
```
|
||||
|
||||
### 内存优化策略
|
||||
|
||||
系统实现了多层次的内存优化:
|
||||
|
||||
| 优化层级 | 策略 | 效果 |
|
||||
|------|---------------|---------|
|
||||
| 对象池 | ByteBuf复用 | 减少GC压力 |
|
||||
| 缓存 | Caffeine本地缓存 | 提升访问速度 |
|
||||
| 序列化 | Protobuf二进制格式 | 减少序列化开销 |
|
||||
| 批量 | 批量操作 | 减少网络往返 |
|
||||
|
||||
### 并发控制机制
|
||||
|
||||
系统采用多种并发控制策略:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Request[请求到达] --> Lock{获取锁}
|
||||
Lock --> |成功| Process[处理请求]
|
||||
Lock --> |失败| Queue[加入队列]
|
||||
Process --> Release[释放锁]
|
||||
Queue --> Wait[等待执行]
|
||||
Wait --> Process
|
||||
Release --> Complete[完成处理]
|
||||
```
|
||||
|
||||
## 总结
|
||||
|
||||
云快充协议上行消息处理系统是一个高度优化的分布式消息处理框架,具有以下核心优势:
|
||||
|
||||
### 技术亮点
|
||||
|
||||
1. **模块化设计**:通过注解驱动的命令路由实现高度可扩展的架构
|
||||
2. **高性能处理**:基于Netty的异步I/O和批量处理策略
|
||||
3. **强一致性**:通过事务性和批量操作保证数据一致性
|
||||
4. **可观测性**:完善的日志记录和监控指标体系
|
||||
|
||||
### 业务价值
|
||||
|
||||
1. **实时响应**:毫秒级的消息处理延迟
|
||||
2. **高可靠性**:多重校验和容错机制
|
||||
3. **可扩展性**:支持多种协议版本和未来扩展
|
||||
4. **运维友好**:清晰的日志和监控界面
|
||||
|
||||
### 应用场景
|
||||
|
||||
该系统适用于大规模充电桩管理平台,能够处理数万甚至数十万充电桩的并发消息处理需求,为电动车充电服务提供稳定可靠的技术支撑。
|
||||
475
docs/核心模块详解/协议实现模块/云快充协议实现/下行消息处理.md
Normal file
475
docs/核心模块详解/协议实现模块/云快充协议实现/下行消息处理.md
Normal file
@@ -0,0 +1,475 @@
|
||||
# 云快充协议下行消息处理链路
|
||||
|
||||
<cite>
|
||||
**本文档引用的文件**
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java)
|
||||
- [DownlinkGrpcService.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkGrpcService.java)
|
||||
- [YunKuaiChongDownlinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongDownlinkCmdExe.java)
|
||||
- [YunKuaiChongDownlinkCmdConverter.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/mapping/YunKuaiChongDownlinkCmdConverter.java)
|
||||
- [YunKuaiChongV150RemoteStartDLCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RemoteStartDLCmd.java)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java)
|
||||
- [AbstractYunKuaiChongCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/AbstractYunKuaiChongCmdExe.java)
|
||||
- [ProtocolSession.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/ProtocolSession.java)
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java)
|
||||
- [DownlinkCmdEnum.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/DownlinkCmdEnum.java)
|
||||
</cite>
|
||||
|
||||
## 目录
|
||||
|
||||
1. [概述](#概述)
|
||||
2. [系统架构](#系统架构)
|
||||
3. [外部接口层](#外部接口层)
|
||||
4. [协议处理层](#协议处理层)
|
||||
5. [命令转换层](#命令转换层)
|
||||
6. [消息构建层](#消息构建层)
|
||||
7. [传输层](#传输层)
|
||||
8. [错误处理与超时机制](#错误处理与超时机制)
|
||||
9. [完整时序图](#完整时序图)
|
||||
10. [总结](#总结)
|
||||
|
||||
## 概述
|
||||
|
||||
云快充协议下行消息处理链路是一个完整的从外部系统发起控制请求到最终发送到充电桩设备的处理流程。该系统支持两种主要的外部接口:REST
|
||||
API 和 gRPC 服务,通过统一的协议处理框架实现对云快充协议的下行指令处理。
|
||||
|
||||
## 系统架构
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "外部接口层"
|
||||
REST[REST API<br/>DownlinkController]
|
||||
GRPC[gRPC服务<br/>DownlinkGrpcService]
|
||||
end
|
||||
subgraph "协议处理层"
|
||||
PMM[协议消息处理器<br/>YunKuaiChongProtocolMessageProcessor]
|
||||
DCC[命令转换器<br/>YunKuaiChongDownlinkCmdConverter]
|
||||
DCE[命令执行器基类<br/>AbstractYunKuaiChongCmdExe]
|
||||
end
|
||||
subgraph "命令执行层"
|
||||
RS[远程启动命令<br/>YunKuaiChongV150RemoteStartDLCmd]
|
||||
SS[停止充电命令<br/>YunKuaiChongV150RemoteStopDLCmd]
|
||||
TS[时间同步命令<br/>YunKuaiChongV150TimeSyncDLCmd]
|
||||
end
|
||||
subgraph "传输层"
|
||||
PS[协议会话<br/>ProtocolSession]
|
||||
TC[TCP通道处理器<br/>TcpChannelHandler]
|
||||
CH[通道<br/>Channel]
|
||||
end
|
||||
REST --> PMM
|
||||
GRPC --> PMM
|
||||
PMM --> DCC
|
||||
PMM --> DCE
|
||||
DCE --> RS
|
||||
DCE --> SS
|
||||
DCE --> TS
|
||||
RS --> PS
|
||||
PS --> TC
|
||||
TC --> CH
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L1-L76)
|
||||
- [DownlinkGrpcService.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkGrpcService.java#L1-L185)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L1-L204)
|
||||
|
||||
## 外部接口层
|
||||
|
||||
### REST API 接口
|
||||
|
||||
REST API 提供了一个标准的 HTTP 接口,通过 `/api/onDownlink` 端点接收下行控制请求。
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as 外部客户端
|
||||
participant Controller as DownlinkController
|
||||
participant Registry as SessionRegistry
|
||||
participant Session as ProtocolSession
|
||||
Client->>Controller : POST /api/onDownlink<br/>protobuf消息
|
||||
Controller->>Controller : 解析消息获取sessionId
|
||||
Controller->>Registry : 获取ProtocolSession
|
||||
Registry-->>Controller : 返回session或null
|
||||
alt session存在
|
||||
Controller->>Session : onDownlink(downlinkMsg)
|
||||
Session-->>Controller : 处理完成
|
||||
Controller-->>Client : 200 OK
|
||||
else session不存在
|
||||
Controller-->>Client : 404 Not Found
|
||||
end
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L37-L75)
|
||||
|
||||
### gRPC 服务接口
|
||||
|
||||
gRPC 服务提供了高性能的双向流式通信接口,支持实时的下行控制请求处理。
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as gRPC客户端
|
||||
participant Service as DownlinkGrpcService
|
||||
participant Registry as SessionRegistry
|
||||
participant Session as ProtocolSession
|
||||
Client->>Service : onDownlink(stream)<br/>连接请求
|
||||
Service->>Service : 建立连接状态
|
||||
Client->>Service : onDownlink(stream)<br/>下行请求
|
||||
Service->>Service : 解析消息获取sessionId
|
||||
Service->>Registry : 获取ProtocolSession
|
||||
Registry-->>Service : 返回session或null
|
||||
alt session存在
|
||||
Service->>Session : onDownlink(downlinkMsg)
|
||||
Session-->>Service : 处理完成
|
||||
Service-->>Client : 响应确认
|
||||
else session不存在
|
||||
Service-->>Client : 记录日志
|
||||
end
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkGrpcService.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkGrpcService.java#L123-L184)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L1-L76)
|
||||
- [DownlinkGrpcService.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkGrpcService.java#L1-L185)
|
||||
|
||||
## 协议处理层
|
||||
|
||||
### 协议消息处理器
|
||||
|
||||
协议消息处理器负责接收来自外部接口的下行请求,并将其转换为内部的消息格式进行处理。
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([接收下行请求]) --> ParseMsg["解析DownlinkRequestMessage"]
|
||||
ParseMsg --> ConvertCmd["转换命令类型"]
|
||||
ConvertCmd --> CheckSupport{"命令是否支持?"}
|
||||
CheckSupport --> |否| LogWarn["记录警告日志"]
|
||||
CheckSupport --> |是| CreateMsg["创建YunKuaiChongDwonlinkMessage"]
|
||||
CreateMsg --> GetExecutor["获取命令执行器"]
|
||||
GetExecutor --> ExecutorExists{"执行器是否存在?"}
|
||||
ExecutorExists --> |否| LogInfo["记录未知指令日志"]
|
||||
ExecutorExists --> |是| ExecuteCmd["执行命令"]
|
||||
ExecuteCmd --> End([处理完成])
|
||||
LogWarn --> End
|
||||
LogInfo --> End
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L140-L202)
|
||||
|
||||
### 命令枚举定义
|
||||
|
||||
系统定义了完整的下行命令枚举,涵盖了所有支持的云快充协议指令。
|
||||
|
||||
| 命令类型 | 描述 | 协议支持 |
|
||||
|----------------------------|---------|-------|
|
||||
| LOGIN_ACK | 登录应答 | 全版本 |
|
||||
| HEARTBEAT_ACK | 心跳应答 | 全版本 |
|
||||
| SET_PRICING | 设置定价模型 | V1.5+ |
|
||||
| REMOTE_START_CHARGING | 远程启动充电 | V1.5+ |
|
||||
| REMOTE_STOP_CHARGING | 远程停止充电 | V1.5+ |
|
||||
| TRANSACTION_RECORD_ACK | 交易记录确认 | V1.5+ |
|
||||
| SYNC_TIME_REQUEST | 同步时间请求 | V1.5+ |
|
||||
| OFFLINE_CARD_QUERY_REQUEST | 离线卡查询请求 | V1.5+ |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L1-L204)
|
||||
- [DownlinkCmdEnum.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/DownlinkCmdEnum.java#L1-L55)
|
||||
|
||||
## 命令转换层
|
||||
|
||||
### 命令转换器
|
||||
|
||||
命令转换器负责将通用的 `DownlinkCmdEnum` 转换为云快充协议特定的命令码。
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class DownlinkCmdConverter {
|
||||
<<interface>>
|
||||
+convertToCmd(DownlinkCmdEnum) Integer
|
||||
+supports(DownlinkCmdEnum) boolean
|
||||
+getProtocolName() String
|
||||
}
|
||||
class YunKuaiChongDownlinkCmdConverter {
|
||||
-COMMAND_MAP Map~DownlinkCmdEnum,Integer~
|
||||
-INSTANCE YunKuaiChongDownlinkCmdConverter
|
||||
+getInstance() YunKuaiChongDownlinkCmdConverter
|
||||
+convertToCmd(DownlinkCmdEnum) Integer
|
||||
+getProtocolName() String
|
||||
}
|
||||
DownlinkCmdConverter <|-- YunKuaiChongDownlinkCmdConverter
|
||||
note for YunKuaiChongDownlinkCmdConverter "单例模式,使用ConcurrentHashMap存储转换关系\n提供O(1)查找性能"
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongDownlinkCmdConverter.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/mapping/YunKuaiChongDownlinkCmdConverter.java#L1-L89)
|
||||
|
||||
### 转换映射表
|
||||
|
||||
以下是部分关键命令的转换映射:
|
||||
|
||||
| 通用命令 | 云快充协议命令码 | 功能描述 |
|
||||
|------------------------|----------|--------|
|
||||
| REMOTE_START_CHARGING | 0x34 | 远程启动充电 |
|
||||
| REMOTE_STOP_CHARGING | 0x36 | 远程停止充电 |
|
||||
| SET_PRICING | 0x58 | 设置定价模型 |
|
||||
| SYNC_TIME_REQUEST | 0x56 | 同步时间请求 |
|
||||
| TRANSACTION_RECORD_ACK | 0x40 | 交易记录确认 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongDownlinkCmdConverter.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/mapping/YunKuaiChongDownlinkCmdConverter.java#L1-L89)
|
||||
|
||||
## 消息构建层
|
||||
|
||||
### 命令执行器基类
|
||||
|
||||
抽象命令执行器基类提供了消息编码和发送的核心功能。
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class AbstractYunKuaiChongCmdExe {
|
||||
-DOWNLINK_CMD_CONVERTER DownlinkCmdConverter
|
||||
+encode(int, int, int, ByteBuf) byte[]
|
||||
+encodeAndWriteFlush(int, int, int, ByteBuf, TcpSession) void
|
||||
+encodeAndWriteFlush(int, ByteBuf, TcpSession) void
|
||||
+encodeAndWriteFlush(DownlinkCmdEnum, int, int, ByteBuf, TcpSession) void
|
||||
+encodeAndWriteFlush(DownlinkCmdEnum, ByteBuf, TcpSession) void
|
||||
+encodePileCode(String) byte[]
|
||||
+encodeGunCode(String) byte[]
|
||||
+encodeTradeNo(String) byte[]
|
||||
+encodeLogicalCardNo(String) byte[]
|
||||
+encodePhysicalCardNo(String) long
|
||||
}
|
||||
class YunKuaiChongV150RemoteStartDLCmd {
|
||||
+execute(TcpSession, YunKuaiChongDwonlinkMessage, ProtocolContext) void
|
||||
}
|
||||
AbstractYunKuaiChongCmdExe <|-- YunKuaiChongV150RemoteStartDLCmd
|
||||
note for AbstractYunKuaiChongCmdExe "提供协议特定的消息编码功能\n包括BCD编码、CRC校验等"
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [AbstractYunKuaiChongCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/AbstractYunKuaiChongCmdExe.java#L1-L310)
|
||||
|
||||
### 远程启动命令实现
|
||||
|
||||
以远程启动充电命令为例,展示具体的命令处理流程。
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Cmd as YunKuaiChongV150RemoteStartDLCmd
|
||||
participant Base as AbstractYunKuaiChongCmdExe
|
||||
participant Session as TcpSession
|
||||
participant Channel as Netty Channel
|
||||
Cmd->>Cmd : 解析请求参数
|
||||
Note over Cmd : pileCode, gunCode, tradeNo<br/>logicalCardNo, physicalCardNo<br/>limitYuan
|
||||
Cmd->>Cmd : 创建ByteBuf
|
||||
Cmd->>Cmd : 写入交易流水号(32字节BCD)
|
||||
Cmd->>Cmd : 写入桩编码(14字节BCD)
|
||||
Cmd->>Cmd : 写入枪号(2字节BCD)
|
||||
Cmd->>Cmd : 写入逻辑卡号(16字节BCD)
|
||||
Cmd->>Cmd : 写入物理卡号(8字节LE Long)
|
||||
Cmd->>Cmd : 写入账户余额(4字节LE Int)
|
||||
Cmd->>Base : encodeAndWriteFlush(cmd, msgBody, session)
|
||||
Base->>Base : encode(cmd, seqNo, flag, msgBody)
|
||||
Note over Base : 添加协议头、长度、序列号<br/>加密标志、命令字、数据域<br/>计算CRC校验和
|
||||
Base->>Session : writeAndFlush(encodedMessage)
|
||||
Session->>Channel : 发送到网络
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150RemoteStartDLCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RemoteStartDLCmd.java#L1-L72)
|
||||
- [AbstractYunKuaiChongCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/AbstractYunKuaiChongCmdExe.java#L180-L220)
|
||||
|
||||
### 协议消息格式
|
||||
|
||||
云快充协议的下行消息格式如下:
|
||||
|
||||
| 字段 | 长度 | 描述 |
|
||||
|------|-----|-----------------|
|
||||
| 帧头 | 1字节 | 固定值 0x68 |
|
||||
| 数据长度 | 1字节 | 包含后续所有字段的长度 |
|
||||
| 序列号 | 2字节 | LE格式,自动递增 |
|
||||
| 加密标志 | 1字节 | 0表示正常加密 |
|
||||
| 命令字 | 1字节 | 下行命令类型 |
|
||||
| 数据域 | 可变 | 具体命令的数据内容 |
|
||||
| 校验和 | 2字节 | CRC校验,多项式0x180D |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [AbstractYunKuaiChongCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/AbstractYunKuaiChongCmdExe.java#L1-L310)
|
||||
- [YunKuaiChongV150RemoteStartDLCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RemoteStartDLCmd.java#L1-L72)
|
||||
|
||||
## 传输层
|
||||
|
||||
### 协议会话管理
|
||||
|
||||
协议会话负责维护与充电桩设备的连接状态和消息路由。
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ProtocolSession {
|
||||
<<abstract>>
|
||||
#protocolName String
|
||||
#id UUID
|
||||
#lastActivityTime LocalDateTime
|
||||
#pileCodeSet Set~String~
|
||||
#requestCache Cache~String,Object~
|
||||
+onDownlink(DownlinkRequestMessage) void
|
||||
+close(SessionCloseReason) void
|
||||
+addPileCode(String) void
|
||||
+addSchedule(String, Function) void
|
||||
}
|
||||
class TcpSession {
|
||||
-sendDownlinkConsumer Consumer~DownlinkRequestMessage~
|
||||
-writeAndFlushConsumer Consumer~ByteBuf~
|
||||
-sequenceNumber AtomicInteger
|
||||
+nextSeqNo(SequenceNumberLength) int
|
||||
+writeAndFlush(ByteBuf) void
|
||||
+onDownlink(DownlinkRequestMessage) void
|
||||
}
|
||||
ProtocolSession <|-- TcpSession
|
||||
note for TcpSession "继承自ProtocolSession<br/>提供TCP连接的下行消息发送功能"
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolSession.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/ProtocolSession.java#L1-L124)
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L30-L128)
|
||||
|
||||
### TCP通道处理器
|
||||
|
||||
TCP通道处理器负责实际的网络数据传输和消息发送。
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([消息发送请求]) --> CheckSession{"会话是否有效?"}
|
||||
CheckSession --> |否| LogError["记录错误日志"]
|
||||
CheckSession --> |是| EncodeMsg["编码协议消息"]
|
||||
EncodeMsg --> WriteToChannel["写入Netty Channel"]
|
||||
WriteToChannel --> CheckResult{"发送是否成功?"}
|
||||
CheckResult --> |否| LogFailure["记录发送失败"]
|
||||
CheckResult --> |是| LogSuccess["记录发送成功"]
|
||||
LogError --> End([结束])
|
||||
LogFailure --> End
|
||||
LogSuccess --> End
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L175-L218)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [ProtocolSession.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/ProtocolSession.java#L1-L124)
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L1-L218)
|
||||
|
||||
## 错误处理与超时机制
|
||||
|
||||
### 超时处理
|
||||
|
||||
系统在多个层面实现了超时处理机制:
|
||||
|
||||
1. **HTTP超时**: REST API 默认超时时间为3秒
|
||||
2. **gRPC超时**: 连接建立和消息处理都有超时保护
|
||||
3. **消息发送超时**: TCP通道处理器记录发送时间并监控结果
|
||||
|
||||
### 错误恢复策略
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Error([发生错误]) --> CheckType{"错误类型"}
|
||||
CheckType --> |网络错误| Retry["重试机制"]
|
||||
CheckType --> |协议错误| LogError["记录错误日志"]
|
||||
CheckType --> |超时错误| TimeoutAction["超时处理"]
|
||||
Retry --> RetryCount{"重试次数检查"}
|
||||
RetryCount --> |未超限| DelayRetry["延迟重试"]
|
||||
RetryCount --> |超限| FailFast["快速失败"]
|
||||
DelayRetry --> SendAgain["重新发送"]
|
||||
SendAgain --> Success{"发送成功?"}
|
||||
Success --> |是| LogSuccess["记录成功"]
|
||||
Success --> |否| LogError
|
||||
TimeoutAction --> LogTimeout["记录超时"]
|
||||
LogError --> Cleanup["清理资源"]
|
||||
FailFast --> Cleanup
|
||||
LogSuccess --> Cleanup
|
||||
LogTimeout --> Cleanup
|
||||
Cleanup --> End([结束])
|
||||
```
|
||||
|
||||
### 异常处理机制
|
||||
|
||||
系统提供了完善的异常处理机制,包括:
|
||||
|
||||
- **DownlinkException**: 下行消息处理异常
|
||||
- **IllegalArgumentException**: 参数验证异常
|
||||
- **RuntimeException**: 通用运行时异常处理
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L37-L75)
|
||||
- [DownlinkGrpcService.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkGrpcService.java#L123-L184)
|
||||
|
||||
## 完整时序图
|
||||
|
||||
以下展示了从外部系统发起远程启动充电请求到设备响应的完整时序:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as 外部客户端
|
||||
participant Controller as DownlinkController
|
||||
participant PMM as 协议消息处理器
|
||||
participant Converter as 命令转换器
|
||||
participant Executor as 命令执行器
|
||||
participant Session as 协议会话
|
||||
participant Channel as TCP通道
|
||||
participant Device as 充电桩设备
|
||||
Note over Client,Device : 远程启动充电流程
|
||||
Client->>Controller : POST /api/onDownlink<br/>RemoteStartChargingRequest
|
||||
Controller->>Controller : 解析protobuf消息
|
||||
Controller->>Session : 获取ProtocolSession
|
||||
Session-->>Controller : 返回会话对象
|
||||
Controller->>PMM : downlinkHandle(sessionToHandlerMsg)
|
||||
PMM->>Converter : convertToCmd(REMOTE_START_CHARGING)
|
||||
Converter-->>PMM : 返回命令码 0x34
|
||||
PMM->>Executor : 获取YunKuaiChongV150RemoteStartDLCmd
|
||||
Executor->>Executor : 解析请求参数
|
||||
Executor->>Executor : 构建消息体
|
||||
Note over Executor : 交易流水号+桩编码+枪号<br/>卡号+余额等信息
|
||||
Executor->>Executor : encodeAndWriteFlush(0x34, msgBody, session)
|
||||
Executor->>Session : writeAndFlush(encodedMessage)
|
||||
Session->>Channel : 发送到网络
|
||||
Channel->>Device : 传输协议消息
|
||||
Device-->>Channel : 设备响应
|
||||
Channel-->>Session : 接收响应
|
||||
Session-->>PMM : 处理上行消息
|
||||
PMM-->>Controller : 处理完成
|
||||
Controller-->>Client : 返回成功响应
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L37-L75)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L140-L202)
|
||||
- [YunKuaiChongV150RemoteStartDLCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150RemoteStartDLCmd.java#L30-L71)
|
||||
|
||||
## 总结
|
||||
|
||||
云快充协议下行消息处理链路展现了现代工业物联网系统的典型架构特点:
|
||||
|
||||
1. **分层架构**: 清晰的分层设计使得系统具有良好的可维护性和扩展性
|
||||
2. **多接口支持**: 同时支持REST API和gRPC两种接口,满足不同场景需求
|
||||
3. **协议抽象**: 通过命令转换器实现协议无关的设计
|
||||
4. **高效传输**: 基于Netty的异步非阻塞IO,确保高并发处理能力
|
||||
5. **健壮性**: 完善的错误处理和超时机制保证系统稳定性
|
||||
|
||||
该系统为云快充平台提供了可靠、高效的下行指令处理能力,支撑着大规模充电桩的远程控制需求。通过模块化的架构设计和标准化的协议处理流程,系统能够快速适配新的协议版本和功能需求。
|
||||
490
docs/核心模块详解/协议实现模块/云快充协议实现/云快充协议实现.md
Normal file
490
docs/核心模块详解/协议实现模块/云快充协议实现/云快充协议实现.md
Normal file
@@ -0,0 +1,490 @@
|
||||
# 云快充协议实现
|
||||
|
||||
<cite>
|
||||
**本文档中引用的文件**
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java)
|
||||
- [YunkuaichongV160ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v160/YunkuaichongV160ProtocolBootstrap.java)
|
||||
- [YunkuaichongV170ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v170/YunkuaichongV170ProtocolBootstrap.java)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java)
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java)
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java)
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java)
|
||||
- [YunKuaiChongDownlinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongDownlinkCmdExe.java)
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java)
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java)
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java)
|
||||
- [DownlinkCmdEnum.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/domain/DownlinkCmdEnum.java)
|
||||
</cite>
|
||||
|
||||
## 目录
|
||||
|
||||
1. [简介](#简介)
|
||||
2. [项目结构](#项目结构)
|
||||
3. [核心组件](#核心组件)
|
||||
4. [架构概览](#架构概览)
|
||||
5. [详细组件分析](#详细组件分析)
|
||||
6. [协议版本管理](#协议版本管理)
|
||||
7. [命令处理机制](#命令处理机制)
|
||||
8. [消息编解码流程](#消息编解码流程)
|
||||
9. [性能考虑](#性能考虑)
|
||||
10. [故障排除指南](#故障排除指南)
|
||||
11. [结论](#结论)
|
||||
|
||||
## 简介
|
||||
|
||||
云快充协议是一个专为电动汽车充电站设计的通信协议实现,采用模块化架构支持多个协议版本(v150、v160、v170),提供了完整的上行和下行消息处理能力。该协议实现了基于Netty的高性能TCP通信,支持心跳检测、远程控制、参数配置等核心功能。
|
||||
|
||||
## 项目结构
|
||||
|
||||
云快充协议的项目结构遵循分层架构设计,主要包含以下模块:
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "协议API层"
|
||||
API[Protocol API]
|
||||
Bootstrap[Protocol Bootstrap]
|
||||
MessageProcessor[Message Processor]
|
||||
end
|
||||
subgraph "云快充实现层"
|
||||
V150[Version 1.5.0]
|
||||
V160[Version 1.6.0]
|
||||
V170[Version 1.7.0]
|
||||
Constants[Protocol Constants]
|
||||
end
|
||||
subgraph "命令处理层"
|
||||
UplinkCmd[Uplink Commands]
|
||||
DownlinkCmd[Downlink Commands]
|
||||
Router[Command Router]
|
||||
end
|
||||
API --> Bootstrap
|
||||
Bootstrap --> MessageProcessor
|
||||
MessageProcessor --> V150
|
||||
MessageProcessor --> V160
|
||||
MessageProcessor --> V170
|
||||
V150 --> UplinkCmd
|
||||
V150 --> DownlinkCmd
|
||||
Router --> UplinkCmd
|
||||
Router --> DownlinkCmd
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L1-L127)
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L1-L48)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L1-L48)
|
||||
- [YunkuaichongV160ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v160/YunkuaichongV160ProtocolBootstrap.java#L1-L48)
|
||||
- [YunkuaichongV170ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v170/YunkuaichongV170ProtocolBootstrap.java#L1-L49)
|
||||
|
||||
## 核心组件
|
||||
|
||||
### YunkuaichongV150ProtocolBootstrap
|
||||
|
||||
`YunkuaichongV150ProtocolBootstrap` 是云快充协议的核心入口点,负责协议的初始化和配置。它继承自 `ProtocolBootstrap`
|
||||
,实现了协议的基本生命周期管理。
|
||||
|
||||
#### 主要特性:
|
||||
|
||||
- **协议标识**:使用 `@ProtocolComponent` 注解标记协议名称
|
||||
- **消息处理器**:创建并返回 `YunKuaiChongProtocolMessageProcessor` 实例
|
||||
- **版本支持**:支持 v150、v160、v170 三个版本
|
||||
|
||||
#### 初始化流程:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Bootstrap as YunkuaichongV150ProtocolBootstrap
|
||||
participant Parent as ProtocolBootstrap
|
||||
participant Processor as MessageProcessor
|
||||
participant Listener as TCP Listener
|
||||
Bootstrap->>Parent : init()
|
||||
Parent->>Parent : loadConfig(protocolName)
|
||||
Parent->>Parent : createForwarder()
|
||||
Parent->>Listener : new TcpListener()
|
||||
Bootstrap->>Processor : new YunKuaiChongProtocolMessageProcessor()
|
||||
Processor->>Processor : initializeRouters()
|
||||
Processor-->>Bootstrap : 返回处理器实例
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L30-L47)
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L40-L80)
|
||||
|
||||
### YunKuaiChongProtocolMessageProcessor
|
||||
|
||||
这是协议消息处理的核心组件,负责解析原始字节流并路由到具体的命令处理器。
|
||||
|
||||
#### 核心功能:
|
||||
|
||||
- **上行消息处理**:解析TCP接收到的原始字节流
|
||||
- **下行消息处理**:处理来自REST API的下行请求
|
||||
- **命令路由**:使用 `ProtocolCommandRouter` 进行命令分发
|
||||
- **校验和验证**:支持多种校验和算法
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L30-L47)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L1-L204)
|
||||
|
||||
## 架构概览
|
||||
|
||||
云快充协议采用分层架构,从底层的网络通信到高层的业务逻辑处理:
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "应用层"
|
||||
REST[REST API Controller]
|
||||
GRPC[gRPC Service]
|
||||
end
|
||||
subgraph "协议层"
|
||||
Bootstrap[Protocol Bootstrap]
|
||||
MessageProcessor[Message Processor]
|
||||
Router[Command Router]
|
||||
end
|
||||
subgraph "传输层"
|
||||
TCP[TCP Listener]
|
||||
ChannelHandler[TCP Channel Handler]
|
||||
end
|
||||
subgraph "网络层"
|
||||
Netty[Netty Framework]
|
||||
Codec[Message Codec]
|
||||
end
|
||||
REST --> Bootstrap
|
||||
GRPC --> Bootstrap
|
||||
Bootstrap --> MessageProcessor
|
||||
MessageProcessor --> Router
|
||||
Router --> ChannelHandler
|
||||
ChannelHandler --> TCP
|
||||
TCP --> Netty
|
||||
Netty --> Codec
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L1-L76)
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L1-L234)
|
||||
|
||||
## 详细组件分析
|
||||
|
||||
### 协议命令注解系统
|
||||
|
||||
`@ProtocolCmd` 注解是云快充协议命令映射的核心机制:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ProtocolCmd {
|
||||
+int value()
|
||||
+String[] protocolNames()
|
||||
}
|
||||
class YunKuaiChongV150HeartbeatULCmd {
|
||||
+execute(TcpSession, YunKuaiChongUplinkMessage, ProtocolContext)
|
||||
}
|
||||
class ProtocolCommandRouter {
|
||||
-Map~String,T~ executorMap
|
||||
+getExecutor(String, int) T
|
||||
+registerExecutor(Class)
|
||||
}
|
||||
ProtocolCmd --> YunKuaiChongV150HeartbeatULCmd : "注解"
|
||||
ProtocolCommandRouter --> YunKuaiChongV150HeartbeatULCmd : "路由"
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java#L1-L33)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L1-L105)
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java#L1-L85)
|
||||
|
||||
#### 注解机制工作原理:
|
||||
|
||||
1. **扫描阶段**:`ProtocolCommandRouter` 使用反射扫描带有 `@ProtocolCmd` 注解的类
|
||||
2. **注册阶段**:将命令字与执行器类建立映射关系
|
||||
3. **路由阶段**:根据协议名称和命令字查找对应的执行器
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java#L1-L33)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L40-L80)
|
||||
|
||||
### 上行命令处理流程
|
||||
|
||||
上行命令处理是协议的核心功能之一,负责解析设备发送的消息:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([接收TCP消息]) --> ValidateHeader["验证协议头<br/>0x68开头"]
|
||||
ValidateHeader --> ParseLength["解析数据长度字段"]
|
||||
ParseLength --> ValidateBounds["边界检查<br/>确保消息完整"]
|
||||
ValidateBounds --> ExtractFields["提取关键字段<br/>序列号、加密标志、帧类型"]
|
||||
ExtractFields --> ChecksumValidation["校验和验证<br/>支持LE和BE两种格式"]
|
||||
ChecksumValidation --> Success{"校验成功?"}
|
||||
Success --> |否| LogError["记录错误日志"]
|
||||
Success --> |是| ParseBody["解析消息体"]
|
||||
ParseBody --> RouteCommand["路由到命令处理器"]
|
||||
RouteCommand --> ExecuteCmd["执行具体命令"]
|
||||
ExecuteCmd --> End([处理完成])
|
||||
LogError --> End
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L50-L150)
|
||||
|
||||
### 下行命令处理流程
|
||||
|
||||
下行命令处理负责响应REST API请求,向设备发送控制指令:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Controller as DownlinkController
|
||||
participant Session as TcpSession
|
||||
participant Processor as MessageProcessor
|
||||
participant Router as CommandRouter
|
||||
participant Executor as CmdExecutor
|
||||
participant Channel as TCP Channel
|
||||
Controller->>Session : 查找会话
|
||||
Session->>Processor : onDownlink(downlinkMsg)
|
||||
Processor->>Processor : convertToCmd(downlinkCmd)
|
||||
Processor->>Router : getExecutor(protocolName, cmd)
|
||||
Router-->>Processor : 返回执行器
|
||||
Processor->>Executor : execute(session, message, context)
|
||||
Executor->>Executor : encodeAndWriteFlush()
|
||||
Executor->>Channel : 写入TCP缓冲区
|
||||
Channel-->>Controller : 返回响应
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L40-L75)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L150-L204)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L40-L75)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L150-L204)
|
||||
|
||||
## 协议版本管理
|
||||
|
||||
云快充协议支持多个版本,每个版本都有自己的 `ProtocolBootstrap` 实现:
|
||||
|
||||
### 版本差异对比
|
||||
|
||||
| 功能特性 | v150 | v160 | v170 |
|
||||
|------|----------|--------|--------|
|
||||
| 基础协议 | ✓ | ✓ | ✓ |
|
||||
| 并行启动 | ✗ | ✓ | ✓ |
|
||||
| 新增命令 | 心跳、登录、对时 | 并行启动命令 | 交易记录命令 |
|
||||
| 向后兼容 | ✓ | ✓ | ✓ |
|
||||
|
||||
### 版本兼容性策略
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "版本兼容性"
|
||||
V150[v1.5.0]
|
||||
V160[v1.6.0]
|
||||
V170[v1.7.0]
|
||||
end
|
||||
subgraph "共享组件"
|
||||
Common[公共命令集]
|
||||
Router[统一路由]
|
||||
Processor[通用处理器]
|
||||
end
|
||||
V150 --> Common
|
||||
V160 --> Common
|
||||
V170 --> Common
|
||||
Common --> Router
|
||||
Router --> Processor
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L20-L25)
|
||||
- [YunkuaichongV160ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v160/YunkuaichongV160ProtocolBootstrap.java#L20-L25)
|
||||
- [YunkuaichongV170ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v170/YunkuaichongV170ProtocolBootstrap.java#L20-L25)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L20-L25)
|
||||
- [YunkuaichongV160ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v160/YunkuaichongV160ProtocolBootstrap.java#L20-L25)
|
||||
- [YunkuaichongV170ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v170/YunkuaichongV170ProtocolBootstrap.java#L20-L25)
|
||||
|
||||
## 命令处理机制
|
||||
|
||||
### 上行命令处理
|
||||
|
||||
以心跳命令为例,展示完整的命令处理流程:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class YunKuaiChongV150HeartbeatULCmd {
|
||||
+execute(TcpSession, YunKuaiChongUplinkMessage, ProtocolContext)
|
||||
-pingAck(TcpSession, YunKuaiChongUplinkMessage, byte[], byte)
|
||||
}
|
||||
class YunKuaiChongUplinkCmdExe {
|
||||
<<abstract>>
|
||||
+execute(TcpSession, YunKuaiChongUplinkMessage, ProtocolContext)*
|
||||
+uplinkMessageBuilder(String, TcpSession, YunKuaiChongUplinkMessage)
|
||||
+encodeAndWriteFlush(DownlinkCmdEnum, int, int, ByteBuf, TcpSession)
|
||||
}
|
||||
class YunKuaiChongUplinkMessage {
|
||||
+int getCmd()
|
||||
+byte[] getMsgBody()
|
||||
+int getSequenceNumber()
|
||||
+int getEncryptionFlag()
|
||||
}
|
||||
YunKuaiChongUplinkCmdExe <|-- YunKuaiChongV150HeartbeatULCmd
|
||||
YunKuaiChongV150HeartbeatULCmd --> YunKuaiChongUplinkMessage : "使用"
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java#L25-L85)
|
||||
|
||||
### 下行命令处理
|
||||
|
||||
下行命令处理涉及REST API调用到设备响应的完整流程:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
RESTReq[REST请求] --> Controller[DownlinkController]
|
||||
Controller --> Session[查找TCP会话]
|
||||
Session --> MsgProc[MessageProcessor]
|
||||
MsgProc --> Converter[命令转换器]
|
||||
Converter --> Router[CommandRouter]
|
||||
Router --> Executor[命令执行器]
|
||||
Executor --> Encoder[消息编码]
|
||||
Encoder --> TCP[发送到设备]
|
||||
Controller --> Success[HTTP 200]
|
||||
Controller --> Timeout[HTTP 408]
|
||||
Controller --> NotFound[HTTP 404]
|
||||
Controller --> Error[HTTP 500]
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L40-L75)
|
||||
- [YunKuaiChongDownlinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongDownlinkCmdExe.java#L1-L19)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongV150HeartbeatULCmd.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/cmd/YunKuaiChongV150HeartbeatULCmd.java#L25-L85)
|
||||
- [DownlinkController.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/adapter/DownlinkController.java#L40-L75)
|
||||
|
||||
## 消息编解码流程
|
||||
|
||||
### TCP消息处理
|
||||
|
||||
`TcpChannelHandler` 负责Netty框架层面的消息处理:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Netty as Netty Channel
|
||||
participant Handler as TcpChannelHandler
|
||||
participant Processor as MessageProcessor
|
||||
participant Stats as Metrics
|
||||
Netty->>Handler : channelRead0(msg)
|
||||
Handler->>Handler : process(msg)
|
||||
Handler->>Processor : uplinkHandleAsync()
|
||||
Processor->>Stats : 记录统计信息
|
||||
Processor-->>Handler : 处理完成
|
||||
Handler->>Handler : writeAndFlush()
|
||||
Handler->>Netty : 写入响应
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L60-L120)
|
||||
|
||||
### 消息格式规范
|
||||
|
||||
云快充协议的消息格式采用固定长度头部加可变长度体的设计:
|
||||
|
||||
| 字段 | 长度 | 描述 |
|
||||
|------|-----|-----------|
|
||||
| 标识符 | 1字节 | 固定值 0x68 |
|
||||
| 数据长度 | 1字节 | 包含头部的总长度 |
|
||||
| 序列号 | 2字节 | 请求-响应配对 |
|
||||
| 加密标志 | 1字节 | 0=明文,1=加密 |
|
||||
| 帧类型 | 1字节 | 命令字 |
|
||||
| 消息体 | 可变 | 具体命令数据 |
|
||||
| 校验和 | 2字节 | CRC校验 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [TcpChannelHandler.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpChannelHandler.java#L60-L120)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L50-L150)
|
||||
|
||||
## 性能考虑
|
||||
|
||||
### 异步处理机制
|
||||
|
||||
协议采用异步处理模式,避免阻塞主线程:
|
||||
|
||||
- **消息队列**:使用Netty的事件循环处理消息
|
||||
- **并发控制**:支持多线程并发处理
|
||||
- **内存优化**:使用ByteBuf进行零拷贝操作
|
||||
|
||||
### 缓存策略
|
||||
|
||||
- **命令路由缓存**:`ProtocolCommandRouter` 使用ConcurrentHashMap缓存命令映射
|
||||
- **会话管理**:TCP会话状态缓存
|
||||
- **统计信息**:实时监控指标缓存
|
||||
|
||||
### 错误处理
|
||||
|
||||
- **快速失败**:无效消息立即丢弃
|
||||
- **优雅降级**:部分功能不可用时保持核心功能
|
||||
- **日志记录**:详细的错误日志便于问题排查
|
||||
|
||||
## 故障排除指南
|
||||
|
||||
### 常见问题及解决方案
|
||||
|
||||
#### 1. 协议版本不匹配
|
||||
|
||||
**症状**:设备无法识别命令
|
||||
**解决方案**:检查协议版本配置,确保客户端和服务端版本一致
|
||||
|
||||
#### 2. 校验和错误
|
||||
|
||||
**症状**:日志显示校验失败
|
||||
**解决方案**:检查消息编码和传输过程中的数据完整性
|
||||
|
||||
#### 3. 命令路由失败
|
||||
|
||||
**症状**:未知命令错误日志
|
||||
**解决方案**:确认命令注解配置正确,检查命令处理器注册
|
||||
|
||||
#### 4. TCP连接问题
|
||||
|
||||
**症状**:连接频繁断开
|
||||
**解决方案**:检查网络稳定性,调整心跳间隔配置
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L100-L150)
|
||||
|
||||
## 结论
|
||||
|
||||
云快充协议实现展现了现代通信协议设计的最佳实践:
|
||||
|
||||
### 主要优势
|
||||
|
||||
1. **模块化设计**:清晰的分层架构便于维护和扩展
|
||||
2. **版本兼容性**:支持多版本并保持向后兼容
|
||||
3. **高性能**:基于Netty的异步处理机制
|
||||
4. **可扩展性**:灵活的命令路由机制支持新功能添加
|
||||
5. **可靠性**:完善的错误处理和监控机制
|
||||
|
||||
### 技术特色
|
||||
|
||||
- **注解驱动开发**:简化命令注册和映射
|
||||
- **统一消息处理**:标准化的消息编解码流程
|
||||
- **灵活的版本管理**:支持协议演进和向后兼容
|
||||
- **完善的监控体系**:实时性能监控和错误追踪
|
||||
|
||||
该协议为电动汽车充电站提供了稳定、高效的通信解决方案,具备良好的可维护性和扩展性,能够满足不同场景下的通信需求。
|
||||
444
docs/核心模块详解/协议实现模块/云快充协议实现/核心架构与初始化.md
Normal file
444
docs/核心模块详解/协议实现模块/云快充协议实现/核心架构与初始化.md
Normal file
@@ -0,0 +1,444 @@
|
||||
# 云快充协议核心架构与初始化
|
||||
|
||||
<cite>
|
||||
**本文档中引用的文件**
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java)
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java)
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java)
|
||||
- [ProtocolMessageProcessor.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolMessageProcessor.java)
|
||||
- [ProtocolContext.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolContext.java)
|
||||
- [TcpListener.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpListener.java)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java)
|
||||
- [YunKuaiChongUplinkMessage.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongUplinkMessage.java)
|
||||
- [YunKuaiChongUplinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongUplinkCmdExe.java)
|
||||
- [YunKuaiChongDownlinkCmdExe.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongDownlinkCmdExe.java)
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java)
|
||||
- [TcpCfg.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/cfg/TcpCfg.java)
|
||||
- [ProtocolCmd.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/annotation/ProtocolCmd.java)
|
||||
</cite>
|
||||
|
||||
## 目录
|
||||
|
||||
1. [引言](#引言)
|
||||
2. [项目结构概览](#项目结构概览)
|
||||
3. [核心架构设计](#核心架构设计)
|
||||
4. [协议栈初始化流程](#协议栈初始化流程)
|
||||
5. [消息处理机制](#消息处理机制)
|
||||
6. [协议版本管理](#协议版本管理)
|
||||
7. [基础设施组件注入](#基础设施组件注入)
|
||||
8. [性能优化策略](#性能优化策略)
|
||||
9. [故障排除指南](#故障排除指南)
|
||||
10. [总结](#总结)
|
||||
|
||||
## 引言
|
||||
|
||||
云快充协议是JChargePointProtocol项目中的核心通信协议,负责电动汽车充电桩与中央管理系统之间的数据交换。本文档深入解析了云快充协议的核心架构设计,重点阐述了YunkuaichongV150ProtocolBootstrap如何继承ProtocolBootstrap并实现抽象方法,完成协议栈的初始化流程,包括TCP监听器的创建、端口配置、编解码器的注册以及会话管理器的设置。
|
||||
|
||||
## 项目结构概览
|
||||
|
||||
云快充协议采用模块化设计,主要包含以下核心模块:
|
||||
|
||||
```mermaid
|
||||
graph TB
|
||||
subgraph "协议层"
|
||||
Bootstrap[协议引导器]
|
||||
Processor[消息处理器]
|
||||
Constants[协议常量]
|
||||
end
|
||||
subgraph "版本管理"
|
||||
V150[V150协议]
|
||||
V160[V160协议]
|
||||
V170[V170协议]
|
||||
end
|
||||
subgraph "基础设施"
|
||||
Context[协议上下文]
|
||||
Router[命令路由器]
|
||||
Listener[TCP监听器]
|
||||
end
|
||||
Bootstrap --> Processor
|
||||
Bootstrap --> V150
|
||||
Bootstrap --> V160
|
||||
Bootstrap --> V170
|
||||
Processor --> Router
|
||||
Processor --> Context
|
||||
Bootstrap --> Listener
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L1-L48)
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L1-L50)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L1-L48)
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java#L1-L44)
|
||||
|
||||
## 核心架构设计
|
||||
|
||||
### 协议引导器层次结构
|
||||
|
||||
云快充协议采用了清晰的分层架构设计,通过继承和组合的方式实现了高度的可扩展性和可维护性:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ProtocolBootstrap {
|
||||
<<abstract>>
|
||||
#ProtocolContext protocolContext
|
||||
#ProtocolCfg protocolCfg
|
||||
#Listener listener
|
||||
#Forwarder forwarder
|
||||
+init() void
|
||||
+destroy() void
|
||||
+health() Health
|
||||
#getProtocolName()* String
|
||||
#_init()* void
|
||||
#_destroy()* void
|
||||
#messageProcessor()* ProtocolMessageProcessor
|
||||
}
|
||||
class YunkuaichongV150ProtocolBootstrap {
|
||||
+String PROTOCOL_NAME
|
||||
+getProtocolName() String
|
||||
+_init() void
|
||||
+_destroy() void
|
||||
+messageProcessor() ProtocolMessageProcessor
|
||||
}
|
||||
class YunkuaichongV160ProtocolBootstrap {
|
||||
+String PROTOCOL_NAME
|
||||
+getProtocolName() String
|
||||
+_init() void
|
||||
+_destroy() void
|
||||
+messageProcessor() ProtocolMessageProcessor
|
||||
}
|
||||
class YunkuaichongV170ProtocolBootstrap {
|
||||
+String PROTOCOL_NAME
|
||||
+getProtocolName() String
|
||||
+_init() void
|
||||
+_destroy() void
|
||||
+messageProcessor() ProtocolMessageProcessor
|
||||
}
|
||||
class YunKuaiChongProtocolMessageProcessor {
|
||||
-ProtocolCommandRouter~YunKuaiChongUplinkCmdExe~ uplinkRouter
|
||||
-ProtocolCommandRouter~YunKuaiChongDownlinkCmdExe~ downlinkRouter
|
||||
-DownlinkCmdConverter downlinkCmdConverter
|
||||
+uplinkHandle(ListenerToHandlerMsg) void
|
||||
+doDownlinkHandle(SessionToHandlerMsg) void
|
||||
-exeCmd(YunKuaiChongUplinkMessage, TcpSession) void
|
||||
-exeCmd(YunKuaiChongDwonlinkMessage, TcpSession) void
|
||||
}
|
||||
ProtocolBootstrap <|-- YunkuaichongV150ProtocolBootstrap
|
||||
ProtocolBootstrap <|-- YunkuaichongV160ProtocolBootstrap
|
||||
ProtocolBootstrap <|-- YunkuaichongV170ProtocolBootstrap
|
||||
ProtocolBootstrap --> YunKuaiChongProtocolMessageProcessor
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L25-L127)
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L20-L48)
|
||||
|
||||
### 消息处理架构
|
||||
|
||||
协议的消息处理采用了事件驱动的设计模式,通过命令路由器实现消息的动态路由和处理:
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Client as 客户端设备
|
||||
participant Listener as TCP监听器
|
||||
participant Processor as 消息处理器
|
||||
participant Router as 命令路由器
|
||||
participant Executor as 命令执行器
|
||||
Client->>Listener : 发送原始字节流
|
||||
Listener->>Processor : uplinkHandle(ListenerToHandlerMsg)
|
||||
Processor->>Processor : 解析协议头
|
||||
Processor->>Processor : 校验和验证
|
||||
Processor->>Processor : 构建YunKuaiChongUplinkMessage
|
||||
Processor->>Router : getExecutor(protocolName, cmd)
|
||||
Router->>Executor : 返回对应的命令执行器
|
||||
Executor->>Executor : execute(session, message, context)
|
||||
Executor-->>Processor : 处理结果
|
||||
Processor-->>Listener : 处理完成
|
||||
Listener-->>Client : 响应消息
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L63-L120)
|
||||
- [ProtocolCommandRouter.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/routing/ProtocolCommandRouter.java#L80-L104)
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L27-L61)
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L40-L80)
|
||||
|
||||
## 协议栈初始化流程
|
||||
|
||||
### Bootstrap类的生命周期管理
|
||||
|
||||
YunkuaichongV150ProtocolBootstrap作为协议栈的入口点,继承自ProtocolBootstrap并实现了三个关键的抽象方法:
|
||||
|
||||
#### onInit方法实现
|
||||
|
||||
Bootstrap类的初始化过程遵循严格的生命周期管理原则:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([启动初始化]) --> GetProtocolName["获取协议名称<br/>getProtocolName()"]
|
||||
GetProtocolName --> LoadConfig["加载协议配置<br/>protocolContext.getProtocolsConfigProvider()"]
|
||||
LoadConfig --> CheckForwarder{"检查转发器类型"}
|
||||
CheckForwarder --> |内存模式| CreateMemoryForwarder["创建内存转发器<br/>MemoryForwarder"]
|
||||
CheckForwarder --> |Kafka模式| CreateKafkaForwarder["创建Kafka转发器<br/>KafkaForwarder"]
|
||||
CreateMemoryForwarder --> LoadTcpConfig["加载TCP配置<br/>protocolCfg.getListener().getTcp()"]
|
||||
CreateKafkaForwarder --> LoadTcpConfig
|
||||
LoadTcpConfig --> CreateTcpListener["创建TCP监听器<br/>TcpListener"]
|
||||
CreateTcpListener --> CallOnInit["_init()回调"]
|
||||
CallOnInit --> CreateMessageProcessor["创建消息处理器<br/>messageProcessor()"]
|
||||
CreateMessageProcessor --> Complete([初始化完成])
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L40-L80)
|
||||
- [YunkuaichongV150ProtocolBootstrap.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/v150/YunkuaichongV150ProtocolBootstrap.java#L25-L47)
|
||||
|
||||
#### onDestroy方法实现
|
||||
|
||||
协议栈的销毁过程同样遵循优雅关闭的原则:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([开始销毁]) --> LogDestroy["记录销毁日志<br/>log.info('{} destroy...', getProtocolName())"]
|
||||
LogDestroy --> CheckListener{"监听器存在?"}
|
||||
CheckListener --> |是| DestroyListener["销毁TCP监听器<br/>listener.destroy()"]
|
||||
CheckListener --> |否| CheckForwarder{"转发器存在?"}
|
||||
DestroyListener --> CheckForwarder
|
||||
CheckForwarder --> |是| DestroyForwarder["销毁转发器<br/>forwarder.destroy()"]
|
||||
CheckForwarder --> |否| CallOnDestroy["_destroy()回调"]
|
||||
DestroyForwarder --> CallOnDestroy
|
||||
CallOnDestroy --> Complete([销毁完成])
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolBootstrap.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolBootstrap.java#L82-L95)
|
||||
|
||||
### TCP监听器的创建与配置
|
||||
|
||||
TCP监听器是协议栈的核心组件,负责建立和维护客户端连接:
|
||||
|
||||
#### 监听器配置参数
|
||||
|
||||
| 配置项 | 类型 | 默认值 | 描述 |
|
||||
|------------------------|---------|-----------|---------------|
|
||||
| bindAddress | String | localhost | 绑定地址 |
|
||||
| bindPort | int | 8888 | 监听端口 |
|
||||
| bossGroupThreadCount | int | 1 | Boss线程池大小 |
|
||||
| workerGroupThreadCount | int | 4 | Worker线程池大小 |
|
||||
| soBacklog | int | 100 | 连接队列长度 |
|
||||
| soKeepAlive | boolean | true | TCP保活机制 |
|
||||
| nodelay | boolean | true | TCP_NODELAY选项 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [TcpListener.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/listener/tcp/TcpListener.java#L40-L70)
|
||||
- [TcpCfg.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/cfg/TcpCfg.java#L15-L46)
|
||||
|
||||
## 消息处理机制
|
||||
|
||||
### 原始字节流解析流程
|
||||
|
||||
YunKuaiChongProtocolMessageProcessor的核心职责是将接收到的原始字节流解析为结构化的消息对象:
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([接收原始字节流]) --> QuickFail["快速失败检查<br/>长度 & 起始标志"]
|
||||
QuickFail --> ParseHeader["解析协议头<br/>dataLength, seqNo, encryptFlag, frameType"]
|
||||
ParseHeader --> BoundaryCheck["边界检查<br/>dataLength & 可读字节数"]
|
||||
BoundaryCheck --> FieldParse["字段快速解析<br/>序列号、加密标志、帧类型"]
|
||||
FieldParse --> ChecksumCheck["校验和验证<br/>CRC LE & BE两种模式"]
|
||||
ChecksumCheck --> BuildMessage["构建消息对象<br/>YunKuaiChongUplinkMessage"]
|
||||
BuildMessage --> RouteMessage["消息路由<br/>ProtocolCommandRouter"]
|
||||
RouteMessage --> ExecuteCmd["执行命令<br/>命令执行器"]
|
||||
ExecuteCmd --> Complete([处理完成])
|
||||
QuickFail --> |失败| DropMessage["丢弃消息"]
|
||||
BoundaryCheck --> |失败| DropMessage
|
||||
ChecksumCheck --> |失败| DropMessage
|
||||
DropMessage --> Complete
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L63-L120)
|
||||
|
||||
### 消息头信息路由机制
|
||||
|
||||
协议通过消息头中的关键字段进行智能路由:
|
||||
|
||||
| 字段 | 用途 | 路由策略 |
|
||||
|-------------------|-----------|-----------|
|
||||
| 命令码(cmd) | 指定具体操作类型 | 命令路由器精确匹配 |
|
||||
| 序列号(seqNo) | 事务跟踪和响应关联 | 保持会话状态 |
|
||||
| 加密标志(encryptFlag) | 安全级别标识 | 决定后续处理流程 |
|
||||
| 数据长度(dataLength) | 消息完整性验证 | 边界检查和解析控制 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L85-L120)
|
||||
|
||||
## 协议版本管理
|
||||
|
||||
### 版本兼容性策略
|
||||
|
||||
云快充协议支持多个版本的并存和演进:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "协议版本"
|
||||
V150[V150<br/>基础功能]
|
||||
V160[V160<br/>并行启动]
|
||||
V170[V170<br/>交易记录增强]
|
||||
end
|
||||
subgraph "版本特性"
|
||||
V150Features[登录认证<br/>心跳检测<br/>充电控制]
|
||||
V160Features[远程并行启动<br/>批量操作]
|
||||
V170Features[增强交易记录<br/>详细状态上报]
|
||||
end
|
||||
V150 --> V150Features
|
||||
V160 --> V160Features
|
||||
V170 --> V170Features
|
||||
V150 -.-> V160
|
||||
V160 -.-> V170
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java#L20-L35)
|
||||
|
||||
### 协议常量定义
|
||||
|
||||
每个协议版本都有明确的命名规范和常量定义:
|
||||
|
||||
| 版本 | 常量名 | 协议名称 |
|
||||
|--------|-------------------|------------------|
|
||||
| v1.5.0 | YUNKUAICHONG_V150 | yunkuaichongV150 |
|
||||
| v1.6.0 | YUNKUAICHONG_V160 | yunkuaichongV160 |
|
||||
| v1.7.0 | YUNKUAICHONG_V170 | yunkuaichongV170 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolConstants.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolConstants.java#L20-L40)
|
||||
|
||||
## 基础设施组件注入
|
||||
|
||||
### ProtocolContext的作用
|
||||
|
||||
ProtocolContext作为基础设施组件的容器,负责向协议栈注入必要的服务:
|
||||
|
||||
```mermaid
|
||||
classDiagram
|
||||
class ProtocolContext {
|
||||
-StatsFactory statsFactory
|
||||
-ProtocolsConfigProvider protocolsConfigProvider
|
||||
-ProtocolSessionRegistryProvider protocolSessionRegistryProvider
|
||||
-ServiceInfoProvider serviceInfoProvider
|
||||
-PartitionProvider partitionProvider
|
||||
-AppQueueFactory appQueueFactory
|
||||
-ShardingThreadPool shardingThreadPool
|
||||
+init() void
|
||||
}
|
||||
class StatsFactory {
|
||||
+createMessagesStats() MessagesStats
|
||||
}
|
||||
class ProtocolsConfigProvider {
|
||||
+loadConfig(protocolName) ProtocolCfg
|
||||
}
|
||||
class ShardingThreadPool {
|
||||
+execute(shardingKey, runnable) void
|
||||
}
|
||||
ProtocolContext --> StatsFactory
|
||||
ProtocolContext --> ProtocolsConfigProvider
|
||||
ProtocolContext --> ShardingThreadPool
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [ProtocolContext.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolContext.java#L25-L65)
|
||||
|
||||
### 组件依赖关系
|
||||
|
||||
ProtocolContext注入的各个组件在协议栈中发挥着关键作用:
|
||||
|
||||
| 组件 | 用途 | 关键功能 |
|
||||
|-------------------------|------|-------------|
|
||||
| StatsFactory | 性能监控 | 消息统计、健康检查 |
|
||||
| ProtocolsConfigProvider | 配置管理 | 协议配置加载 |
|
||||
| ShardingThreadPool | 并发处理 | 请求分片和线程池管理 |
|
||||
| AppQueueFactory | 消息队列 | 异步消息传递 |
|
||||
| ServiceInfoProvider | 服务发现 | 微服务环境下的服务定位 |
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [ProtocolContext.java](file://jcpp-protocol-api/src/main/java/sanbing/jcpp/protocol/ProtocolContext.java#L35-L65)
|
||||
|
||||
## 性能优化策略
|
||||
|
||||
### 消息处理优化
|
||||
|
||||
协议栈采用了多种性能优化技术:
|
||||
|
||||
1. **异步处理**: 使用线程池处理上行消息,避免阻塞网络线程
|
||||
2. **零拷贝**: Netty框架的ByteBuf提供了高效的内存管理
|
||||
3. **批量处理**: 支持批量消息处理以提高吞吐量
|
||||
4. **缓存机制**: 命令路由器使用ConcurrentHashMap实现快速查找
|
||||
|
||||
### 内存管理优化
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
Start([消息到达]) --> FastPath["快速路径<br/>长度检查 & 起始标志"]
|
||||
FastPath --> SliceBuffer["切片缓冲区<br/>避免数据复制"]
|
||||
SliceBuffer --> PoolAllocation["池化分配<br/>减少GC压力"]
|
||||
PoolAllocation --> AsyncProcess["异步处理<br/>非阻塞"]
|
||||
AsyncProcess --> ReleaseBuffer["释放缓冲区<br/>及时回收"]
|
||||
ReleaseBuffer --> Complete([处理完成])
|
||||
```
|
||||
|
||||
**图表来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L63-L120)
|
||||
|
||||
## 故障排除指南
|
||||
|
||||
### 常见问题诊断
|
||||
|
||||
1. **连接失败**: 检查TCP监听器配置和防火墙设置
|
||||
2. **消息解析错误**: 验证协议版本兼容性和消息格式
|
||||
3. **性能问题**: 监控线程池使用率和内存分配情况
|
||||
4. **路由失败**: 检查命令执行器的注册状态
|
||||
|
||||
### 日志分析要点
|
||||
|
||||
协议栈提供了详细的日志记录,便于问题诊断:
|
||||
|
||||
- **连接日志**: 记录客户端连接和断开事件
|
||||
- **消息日志**: 记录消息解析和处理过程
|
||||
- **错误日志**: 记录异常情况和失败原因
|
||||
- **性能日志**: 记录处理时间和资源使用情况
|
||||
|
||||
**章节来源**
|
||||
|
||||
- [YunKuaiChongProtocolMessageProcessor.java](file://jcpp-protocol-yunkuaichong/src/main/java/sanbing/jcpp/protocol/yunkuaichong/YunKuaiChongProtocolMessageProcessor.java#L100-L120)
|
||||
|
||||
## 总结
|
||||
|
||||
云快充协议的核心架构设计体现了现代分布式系统的设计理念:
|
||||
|
||||
1. **模块化设计**: 清晰的分层架构和职责分离
|
||||
2. **可扩展性**: 支持多版本协议并存和演进
|
||||
3. **高性能**: 异步处理和零拷贝优化
|
||||
4. **可靠性**: 完善的错误处理和恢复机制
|
||||
5. **可观测性**: 全面的日志记录和监控指标
|
||||
|
||||
通过YunkuaichongV150ProtocolBootstrap的实现,我们看到了一个典型的协议引导器应该具备的功能:协议栈的初始化、资源管理和生命周期控制。而YunKuaiChongProtocolMessageProcessor则展示了如何高效地处理复杂的协议消息,通过智能路由和优化的解析算法实现高性能的消息处理。
|
||||
|
||||
这种设计不仅满足了当前的业务需求,也为未来的功能扩展和性能优化奠定了坚实的基础。
|
||||
Reference in New Issue
Block a user