WebSocket 客户端驱动
概述
WebSocket 驱动以客户端身份接入外部 streaming server (ws / wss) 的 endpoint,
缓存收到的消息,并在每个采集周期 read() 时返回最新值,属于 push 模型。
| 项目 | 值 |
|---|---|
opc_type | WEBSOCKET |
| 实现类 | plantpulse.driver.protocol.websocket.WebSocketDriver |
| 基础库 | java.net.http.HttpClient.WebSocket (Java 21 runtime 标准 API) |
| read | ✅ (消息缓存) |
| write | ✅ (sendText) |
| 安全 | TLS (wss) — options.tls=true |
如果说 HTTP 驱动是接收外部系统 push (REST POST) 的模型,那么 WebSocket 驱动就是
edge 以 客户端身份 outbound 连接 并接收 streaming 消息的模型。
OPC 注册表单 / 选项
| 字段 | 含义 | 默认值 | 示例 |
|---|---|---|---|
host / port | WebSocket 服务器地址 | — | 192.168.10.99 / 8765, stream.example.com / 443 |
options.path | endpoint path | / | /stream/v1, /realtime/tag |
options.tls | 是否使用 wss | false | true (wss) / false (ws) |
options.subscribe-message | 连接建立后立即发送的消息 (可选) | — | {"op":"subscribe","topic":"line1.tempC"} |
timecycle | read 周期 (ms) | — | 1000 |
实际 endpoint URL:<scheme>://<host>:<port><path> (scheme 为 ws / wss)
标签 plc_address 格式
| 写法 | 含义 |
|---|---|
空值或 _raw_ | 最后接收消息的 全文 (String) |
JSON top-level 键 (例如 tempC) | 消息为 JSON 时该键对应的值 (转换为 String) |
例)服务器 push {"tempC":25.7,"humid":40.2} 时:
plc_address | 结果 |
|---|---|
_raw_ | {"tempC":25.7,"humid":40.2} |
tempC | 25.7 |
humid | 40.2 |
运行流程
connect()时通过HttpClient.newWebSocketBuilder().buildAsync(...)接入 endpoint。- (可选) 若配置了
subscribe-message,则发送一次sendText。 - 服务器发送的所有文本消息在
onText中累积 → fragment 结束时调用handleMessage()。 handleMessage()将消息全文保存到_raw_键。若消息以{开头,则解析 JSON 后 按 top-level 键分别缓存。- 每个采集周期 PLC collector 调用
read(plc_address)→ 查找缓存 → 返回该值。 - 调用
write(value)时通过sendText原样发送到 server。
发生 onError 时按 connected=false 处理。自动重连交由上层 layer 的 reconnect 策略负责。
curl 注册示例
curl -X POST http://<edge-host>/api/v1/opc \
-H "Content-Type: application/json" \
-d '{
"opc_id": "OPC_WS_LINE1",
"opc_type": "WEBSOCKET",
"opc_name": "Line1 WebSocket Stream",
"opc_agent_ip": "192.168.10.99",
"opc_agent_port": "8765",
"site_id": "SITE_00001",
"auto_collect": true,
"timecycle": 1000,
"options": {
"path": "/stream/v1",
"tls": "false",
"subscribe-message": "{\"op\":\"subscribe\",\"topic\":\"line1\"}"
},
"tag_list": [
{"tag_id":"OPC_WS_LINE1_T01","tag_name":"TempC","plc_address":"tempC","data_type":"Float"},
{"tag_id":"OPC_WS_LINE1_T02","tag_name":"Humid","plc_address":"humid","data_type":"Float"},
{"tag_id":"OPC_WS_LINE1_T03","tag_name":"Raw", "plc_address":"_raw_","data_type":"String"}
]
}'
常见错误与解决
| 消息 / 现象 | 原因 | 解决 |
|---|---|---|
[WS] connect 실패: ...handshake... | URL / path 错误、服务器未运行 | 用 wscat -c ws://<host>:<port><path> 直接验证 |
[WS] connect 실패: ...timeout... | 5 秒内未完成握手 | 检查防火墙 / 端口 / TLS 配置 |
read() 始终为空字符串 | 服务器未发送消息 / 未设置 subscribe-message | 在服务器控制台确认 broadcast 流程,注册所需的 subscribe payload |
| 按 JSON 键 read 为空 | 消息不是 JSON / 为嵌套 (nested) 结构 | 先用 _raw_ 接收后再做后处理,或使用独立 parser。(当前 driver 仅提取 top-level) |
| TLS 证书错误 | 自签名 (self-signed) | 生产环境建议使用有效证书。临时可添加 cacerts |
限制 / 后续扩展
- 不支持 JSON path:当前未实现
a.b.c深度提取,仅支持 top-level key。必要时用${VALUE}后处理,或后续引入选项。 - 未内置自动重连:依赖 PLC collector 的 reconnect 周期 / 策略。独立退避机制待后续实现。
- 不处理 binary frame:仅处理 text frame (
onText),以 JSON 等文本 streaming 为主。 - per-message-deflate / 头部认证:仅使用标准
HttpClient.WebSocket的基本功能。自定义头部计划后续开放为选项。
参考
- Java 标准
java.net.http.WebSocketAPI:<https://docs.oracle.com/en/java/javase/17/docs/api/java.net.http/java/net/http/WebSocket.html> - 快速搭建 dummy server:Python
websockets库 —python -m websockets.server.run :8765