跳到主要内容

异步数据管道 (设备变量归档)

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 100020 万点/秒<50ms
SQLite WAL + batch 500050 万点/秒<100ms
TimescaleDB + batch 10000200 万点/秒<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)
  • 降采样任务运行情况

相关文档