跳到主要内容

WebSocket 客户端驱动

概述

WebSocket 驱动以客户端身份接入外部 streaming server (ws / wss) 的 endpoint, 缓存收到的消息,并在每个采集周期 read() 时返回最新值,属于 push 模型。

项目
opc_typeWEBSOCKET
实现类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 / portWebSocket 服务器地址192.168.10.99 / 8765, stream.example.com / 443
options.pathendpoint path//stream/v1, /realtime/tag
options.tls是否使用 wssfalsetrue (wss) / false (ws)
options.subscribe-message连接建立后立即发送的消息 (可选){"op":"subscribe","topic":"line1.tempC"}
timecycleread 周期 (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}
tempC25.7
humid40.2

运行流程

  1. connect() 时通过 HttpClient.newWebSocketBuilder().buildAsync(...) 接入 endpoint。
  2. (可选) 若配置了 subscribe-message,则发送一次 sendText
  3. 服务器发送的所有文本消息在 onText 中累积 → fragment 结束时调用 handleMessage()
  4. handleMessage() 将消息全文保存到 _raw_ 键。若消息以 { 开头,则解析 JSON 后 按 top-level 键分别缓存。
  5. 每个采集周期 PLC collector 调用 read(plc_address) → 查找缓存 → 返回该值。
  6. 调用 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.WebSocket API:<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