异步数据管道 (设备变量归档)
DarraRT 的设备变量归档 (Historian) 采用异步管道, 避免 PLC 每个扫描周期都同步写数据库, 否则 1-2ms 级扫描会被 SQL IO 直接拖垮。
注: 本章的"数据管道"是 PLC 变量 → 时序数据库的设备级归档, 不是业务数据 ETL / 订单流 / CRM 同步等。目标是让工艺参数、产量、报警、能耗等设备级变量长期保存, 供事后分析与报表生成。
问题: 同步 INSERT 会阻塞
假设扫描周期 2ms, 一个 PLC 产生 1000 个变量, 每个变量 2ms 采样 1 次, 那就是:
1000 变量 × (1000 / 2) = 500,000 次采样/秒
如果每次采样都 INSERT, 即使最快的 SQL (SQLite WAL mode) 也撑不住 (磁盘 IO 约 1-2 万 INSERT/秒), 更关键的是每次 INSERT 会阻塞 PLC 扫描线程 — 扫描抖动秒级, 所有实时控制失效。
方案: 环形缓冲 + 异步 Writer 线程 + 批量 INSERT
PLC 扫描线程 (2ms) 异步 Writer 线程 (100ms) 时序数据库
| | |
| --(非阻塞 enqueue)---> RingBuffer | |
| | (定时) dequeue 5000 条 |
| | ---(BEGIN TRANSACTION)-----> |
| | --- 5000 个 INSERT ---------> |
| | ---(COMMIT)-----------------> |
| | |
v v v
环形缓冲区
内存数组, 容量 10-100 万条, 写入 O(1) 无锁 (CAS), 读出 O(N) 批量:
public sealed class HistorianRing {
private readonly Record[] _buf;
private long _head, _tail;
public HistorianRing(int capacity) {
_buf = new Record[NextPowOf2(capacity)];
}
public bool TryEnqueue(in Record rec) {
long head = Interlocked.Read(ref _head);
long tail = Interlocked.Read(ref _tail);
if (head - tail >= _buf.Length) return false; // 满, 丢弃最新
int idx = (int)(head & (_buf.Length - 1));
_buf[idx] = rec;
Interlocked.Increment(ref _head);
return true;
}
public int Drain(Record[] dst, int max) {
long head = Interlocked.Read(ref _head);
long tail = Interlocked.Read(ref _tail);
int avail = (int)Math.Min(head - tail, max);
for (int i = 0; i < avail; i++) {
int idx = (int)((tail + i) & (_buf.Length - 1));
dst[i] = _buf[idx];
}
Interlocked.Add(ref _tail, avail);
return avail;
}
}
public struct Record {
public long TsUnixMs;
public int VarId;
public double Value;
}
PLC 扫描线程调用 TryEnqueue 仅纳秒级, 对扫描抖动影响忽略不计。
异步 Writer 线程
public class HistorianWriter {
private readonly HistorianRing _ring;
private readonly SqliteConnection _conn;
private readonly Record[] _batch = new Record[5000];
private readonly CancellationTokenSource _cts = new();
public HistorianWriter(string dbPath, HistorianRing ring) {
_ring = ring;
_conn = new SqliteConnection("Data Source=" + dbPath + ";Cache=Shared;");
_conn.Open();
using (var cmd = _conn.CreateCommand()) {
cmd.CommandText = "PRAGMA journal_mode=WAL;"
+ "PRAGMA synchronous=NORMAL;"
+ "PRAGMA temp_store=MEMORY;"
+ "PRAGMA cache_size=-65536;";
cmd.ExecuteNonQuery();
}
Task.Factory.StartNew(Run, TaskCreationOptions.LongRunning);
}
private async Task Run() {
while (!_cts.IsCancellationRequested) {
int n = _ring.Drain(_batch, _batch.Length);
if (n > 0) {
FlushBatch(n);
} else {
await Task.Delay(100, _cts.Token); // 100ms 检查一次
}
}
}
private void FlushBatch(int count) {
using var tx = _conn.BeginTransaction();
using var cmd = _conn.CreateCommand();
cmd.Transaction = tx;
cmd.CommandText = "INSERT INTO points(ts, var_id, value) VALUES(@ts, @v, @val)";
var pTs = cmd.Parameters.Add("@ts", SqliteType.Integer);
var pVar = cmd.Parameters.Add("@v", SqliteType.Integer);
var pVal = cmd.Parameters.Add("@val", SqliteType.Real);
for (int i = 0; i < count; i++) {
pTs.Value = _batch[i].TsUnixMs;
pVar.Value = _batch[i].VarId;
pVal.Value = _batch[i].Value;
cmd.ExecuteNonQuery();
}
tx.Commit();
}
public void Stop() { _cts.Cancel(); _conn.Close(); }
}
SQLite WAL 模式 + 批量 5000 条 + 单事务 = 约 50-100 万行/秒 的写入速度, 远超 PLC 生成速度。
表结构
宽表 vs 窄表
窄表 (推荐): 变量定义表 + 数据表两张, 数据表按 (var_id, ts) 组合:
-- 变量字典 (启动时加载到内存哈希表)
CREATE TABLE variables (
id INTEGER PRIMARY KEY,
name TEXT UNIQUE NOT NULL, -- 例如 'MD.Temp'
unit TEXT,
type TEXT -- 'BOOL','INT','REAL',...
);
-- 数据表 (按日分区, 便于归档删除)
CREATE TABLE points (
ts INTEGER NOT NULL, -- Unix ms
var_id INTEGER NOT NULL,
value REAL NOT NULL
) WITHOUT ROWID;
-- 按 (var_id, ts) 查询最常见
CREATE INDEX idx_points_var_ts ON points(var_id, ts);
宽表: 每变量一列。不灵活, 新增变量要 ALTER TABLE, 不推荐用于大规模归档。
按日分区 (逻辑)
每天一张子表, 定期归档/删除:
CREATE TABLE points_20260417 (...);
CREATE TABLE points_20260418 (...);
-- ...
Writer 按当天日期选表, 查询时 UNION ALL 多张子表。SQLite 不支持原生分区, 用 VIEW 屏蔽细节:
CREATE VIEW points AS
SELECT * FROM points_20260417
UNION ALL
SELECT * FROM points_20260418
...;
可选数据库引擎
| 引擎 | 优点 | 缺点 | 适用规模 |
|---|---|---|---|
| SQLite WAL | 无需服务, 嵌入简单 | 单机 | <= 10 万变量, 30 天归档 |
| TimescaleDB | 时序原生支持, 自动分区 | 需 PG 服务器 | 百万变量, 跨节点 |
| InfluxDB 2.x | 专用时序, 压缩好 | LOR 学习曲线 | 需要 Flux 查询的场景 |
| ClickHouse | 海量 OLAP | 写入需批量 | 多年归档, 亿级点 |
默认 SQLite WAL, 对典型 PLC 场景 (<10k 变量) 足够且零运维。
降采样 (Downsampling)
历史数据要压缩保存, 不能永远保留毫秒级原始点。分层策略:
| 层 | 分辨率 | 保留期 | 说明 |
|---|---|---|---|
| 原始 | 扫描周期 | 7 天 | 精确回溯近期故障 |
| 1 分钟 | 1 min 均值 | 90 天 | 常规回看 |
| 1 小时 | 1 h 均值 | 2 年 | 长期趋势 |
| 1 天 | 1 d 均值 | 10 年 | 年度报表 |
降采样由定时任务执行:
-- 每小时执行: 把 1 小时前的原始点合并到 1 分钟表
INSERT INTO points_1m(ts, var_id, avg, min, max)
SELECT
(ts / 60000) * 60000 AS bucket,
var_id,
AVG(value) AS avg,
MIN(value) AS min,
MAX(value) AS max
FROM points
WHERE ts < (strftime('%s','now') - 3600) * 1000
GROUP BY bucket, var_id;
DELETE FROM points WHERE ts < (strftime('%s','now') - 7 * 86400) * 1000;
背压策略 (Buffer 满了怎么办)
环形缓冲被填满时的选项:
| 策略 | 行为 | 场景 |
|---|---|---|
| drop-new | 丢弃最新数据 | 优先保证历史完整 (很少见) |
| drop-old | 覆盖最旧数据 | 优先保证最近数据 (默认) |
| block | 阻塞扫描线程 | 危险, 会影响控制实时性, 禁用 |
| disk-spill | 溢出到本地磁盘文件 | 数据库短时不可达 |
推荐 drop-old + 异步报警 "归档队列满 N 秒"; 短时数据库异常时开启 disk-spill, 恢复后补回。
完整性校验
启动时校验
Writer 启动时检查:
- 上次停机前的最后一条
ts与环形缓冲最老的ts是否有断档 - 断档则打 warn 日志, 记
meta.gap_log表
周期性聚合校验
每天对比:
- 原始点的采样次数
?==?预期 (采样频率 × 86400) - 偏差
> 5%告警 (可能丢了数据)
查询 API
Service 暴露 Historian 查询给 HMI / SDK:
GET /api/historian/query
?vars=MD.Temp,MD.Pressure
&start=1713355200000
&end=1713441600000
&agg=avg|min|max|last
&interval=1m|5m|1h|1d
返回 JSON:
{
"points": [
{ "ts": 1713355200000, "MD.Temp": 85.3, "MD.Pressure": 3.21 },
{ "ts": 1713355260000, "MD.Temp": 85.5, "MD.Pressure": 3.22 },
...
]
}
性能基线
写入
| 配置 | 吞吐 | 延迟 |
|---|---|---|
| SQLite WAL + batch 1000 | 20 万点/秒 | <50ms |
| SQLite WAL + batch 5000 | 50 万点/秒 | <100ms |
| TimescaleDB + batch 10000 | 200 万点/秒 | <30ms |
查询
| 场景 | 延迟 |
|---|---|
| 单变量 24 小时 (原始) | <500ms |
| 单变量 30 天 (1 分钟聚合) | <200ms |
| 10 变量并发 1 天 | <1s |
运维清单
-
ls -l historian.db*观察 WAL 文件大小, 长时间不收缩说明 checkpoint 异常 - 每日备份
.db+.db-wal到远程 - 监控环形缓冲使用率
>=80% 触发告警 - 监控 Writer 线程存活 (心跳变量
$.Historian.HeartbeatTs) - 磁盘可用空间
<20% 告警 - 定期
VACUUM回收空间 (SQLite) - 降采样任务运行情况