Historian 时序数据库
概述
Historian 是工业过程历史数据库, 专门用来存储和查询带时间戳的测量值。DarraRT Historian 是一个变量级时序库:
对 PLC 中所有标记 Archive=True 的变量持续采样、压缩、归档, 并提供毫秒级的按时间范围查询与降采样聚合, 驱动趋势图、报表、SPC 回溯。
核心指标:
- 写入: 单机 20 万点/s 稳态, 峰值 50 万点/s
- 压缩: Swing Door 算法平均 20-100 倍
- 查询: 1 年数据 10 万点降采样到 1000 点
<200ms
适用场景
| 场景 | 频率 | 保留 | 典型查询 |
|---|---|---|---|
| 设备趋势 | 100ms-1s | 1-5 年 | 按时间段 + 降采样 |
| 报警回溯 | 事件驱动 | 3 年 | 报警前后 ±5 min 高分辨率 |
| 能耗统计 | 1 min | 10 年 | 按小时/日 求和 |
| 批次追溯 | 批次驱动 | 永久 | 按批次号查全量 |
| 工艺对比 | 多变量 | 2 年 | 多条趋势叠加 |
选型对比
| 选型 | 优势 | 劣势 | DarraRT 支持 |
|---|---|---|---|
| SQLite + WAL | 零依赖, 嵌入 | 并发写差, 单机 | 内置, 默认 |
| TimescaleDB | PG 生态, SQL, 分区 | 运维复杂 | 一等公民 |
| InfluxDB | 专为时序, 查询快 | 有授权限制 | Plugin |
| ClickHouse | 列存, 压缩狠 | 更新慢 | Plugin |
| 自研 Parquet + DuckDB | 存储廉价 | 自己维护 | 企业版 |
架构图
PLC Runtime Service Historian
┌──────────────┐ ┌──────────────────────────────┐
│ 变量值变化 │─delta─► │ ChangeCollector (订阅) │
│ (Archive=Y) │ │ │ │
└──────────────┘ │ ▼ │
│ Compressor (Swing Door / 死区)│
│ │ │
│ ▼ │
│ Writer (批量 COPY) │
│ │ │
│ ▼ │
│ 热层 (TimescaleDB 行存) │
│ │ │
│ ▼ 7 天后 │
│ 温层 (TimescaleDB 列存压缩) │
│ │ │
│ ▼ 1 年后 │
│ 冷层 (S3 Parquet) │
└──────────┬───────────────────┘
│ API
▼
┌──────────────────────────┐
│ /api/historian/query ... │
│ <darra-trend> │
└──────────────────────────┘
前置条件
- Service 版本 ≥ 2.4, 变量表包含
Archive列 - 磁盘: 热层
≥采样率 × 7 天 × 16 字节 / 压缩比 - TimescaleDB (推荐) 或 InfluxDB 2.x 已部署
详细步骤
1. 变量标注 Archive
# 变量表示例
variables:
- path: DB1.temp
type: REAL
archive: true
deadband: 0.2
min_interval_ms: 100
- path: DB1.pressure
type: REAL
archive: true
deadband: 0.05
compression: swing_door
epsilon: 0.1
2. 死区 / 摆动门压缩
死区: 简单, 变化 <deadband 则丢弃。
摆动门 (Swing Door Trending): 保证任意查询点插值与原始值误差 <ε。
// 简化版摆动门
public sealed class SwingDoor
{
private readonly double _eps;
private Point _last; // 上次发送点
private double _upSlope; // 上边界斜率 (最小)
private double _downSlope; // 下边界斜率 (最大)
private Point _pending; // 候选点
public SwingDoor(double epsilon) { _eps = epsilon; }
public bool Process(Point p, out Point toSend)
{
toSend = default;
if (_last.Equals(default))
{
_last = p;
_pending = p;
return true;
}
var dt = (p.T - _last.T).TotalSeconds;
var newUp = (p.V - _eps - _last.V) / dt;
var newDown = (p.V + _eps - _last.V) / dt;
_upSlope = Math.Max(_upSlope, newUp);
_downSlope = Math.Min(_downSlope, newDown);
if (_upSlope >= _downSlope)
{
// 门关了, 发送上一个候选点
toSend = _pending;
_last = _pending;
_upSlope = double.MinValue;
_downSlope = double.MaxValue;
_pending = p;
return true;
}
_pending = p;
return false;
}
}
典型压缩比 (温度/压力等平稳信号): 50-200 倍。
3. 存储 Schema (TimescaleDB)
CREATE TABLE historian (
ts TIMESTAMPTZ NOT NULL,
var_id INTEGER NOT NULL,
value DOUBLE PRECISION NOT NULL,
quality SMALLINT NOT NULL DEFAULT 100
);
SELECT create_hypertable('historian', 'ts',
chunk_time_interval => INTERVAL '1 day');
CREATE TABLE var_meta (
var_id SERIAL PRIMARY KEY,
var_path TEXT UNIQUE NOT NULL,
unit TEXT,
dtype TEXT,
archive BOOLEAN
);
-- 列存压缩 + 保留策略
ALTER TABLE historian SET (
timescaledb.compress,
timescaledb.compress_segmentby = 'var_id',
timescaledb.compress_orderby = 'ts DESC'
);
SELECT add_compression_policy('historian', INTERVAL '7 days');
SELECT add_retention_policy('historian', INTERVAL '3 years');
-- 连续聚合 (降采样物化视图)
CREATE MATERIALIZED VIEW historian_1min
WITH (timescaledb.continuous) AS
SELECT
time_bucket('1 minute', ts) AS bucket,
var_id,
avg(value) AS avg_v,
min(value) AS min_v,
max(value) AS max_v,
last(value, ts) AS last_v,
count(*) AS cnt
FROM historian
GROUP BY bucket, var_id;
SELECT add_continuous_aggregate_policy('historian_1min',
start_offset => INTERVAL '1 day',
end_offset => INTERVAL '1 minute',
schedule_interval => INTERVAL '1 minute');
同样的策略做 historian_1h / historian_1d, 查询时 Service 自动选最近粒度。
4. 查询 API
GET /api/historian/query?var=DB1.temp&from=2026-04-17T00:00Z&to=2026-04-18T00:00Z&agg=avg&bucket=1m
→ {
"var": "DB1.temp",
"unit": "℃",
"points": [
{ "ts": "2026-04-17T00:00:00Z", "v": 72.3 },
{ "ts": "2026-04-17T00:01:00Z", "v": 72.5 },
...
]
}
参数:
| 参数 | 默认 | 说明 |
|---|---|---|
var | 必填 | 变量路径, 支持通配 DB1.* |
from / to | 必填 | ISO8601 UTC |
agg | last | avg/min/max/sum/p95/first/last/count |
bucket | 自动 | 降采样窗口 1s/1m/5m/1h/1d |
fill | none | none/null/prev/linear 空桶填充 |
limit | 10000 | 最大返回点数 |
5. Service 实现要点
// Darra.PLC.Service / Historian / QueryService.cs
public async Task<QueryResult> QueryAsync(QueryRequest req)
{
// 1. 选择合适的物化视图
var table = SelectTable(req.From, req.To, req.Bucket);
// <24h 用 historian (原始)
// <30d 用 historian_1min
// <2y 用 historian_1h
// >2y 用 historian_1d
// 2. 生成 SQL
var sql = $@"
SELECT time_bucket(@bucket, ts) AS b,
{AggFunction(req.Agg)}(value) AS v
FROM {table}
WHERE var_id = @varId AND ts BETWEEN @from AND @to
GROUP BY b ORDER BY b";
// 3. 查询 + 空桶填充
var rows = await _conn.QueryAsync<BucketRow>(sql, new {
bucket = req.Bucket,
varId = _metaCache.GetVarId(req.Var),
from = req.From,
to = req.To
});
return req.Fill switch {
"prev" => FillPrevious(rows, req),
"linear" => FillLinear(rows, req),
_ => rows.ToResult()
};
}
6. HMI 趋势绑定
<darra-trend
api="/api/historian/query"
series='[
{"var":"DB1.temp","label":"温度","color":"#e74c3c","unit":"℃"},
{"var":"DB1.pressure","label":"压力","color":"#3498db","unit":"bar"}
]'
range="24h"
auto-bucket="true"
refresh-ms="5000">
</darra-trend>
组件内部用 Chart.js 或 ECharts, 根据用户缩放等级自动向 API 请求不同 bucket。
7. 数据归档到冷层
# Service 定时任务 (每月 1 号 02:00)
psql -c "
COPY (SELECT * FROM historian
WHERE ts < NOW() - INTERVAL '1 year')
TO PROGRAM 'aws s3 cp - s3://darrart-cold/$(date +%Y)/$(date +%m).parquet'
WITH (FORMAT parquet);
"
# 清理热层
psql -c "DELETE FROM historian WHERE ts < NOW() - INTERVAL '1 year';"
冷层查询由 DuckDB 直接读 S3 Parquet, 无需恢复到热库。
参数表
| 参数 | 默认 | 调优 |
|---|---|---|
deadband | 变量精度 × 2 | 过大丢信息, 过小不压缩 |
min_interval_ms | 100 | 快变变量设 10, 慢变设 1000 |
swing_door_epsilon | 精度 | 保证回放误差上限 |
chunk_size | 1 day | 按写入量调, 太小碎片多 |
compression_after | 7 days | 压缩降低查询速度 |
retention | 3 years | 按合规要求 |
排错
| 现象 | 原因 | 措施 |
|---|---|---|
| 查询慢 | 没用物化视图 | 检查 EXPLAIN ANALYZE |
| 磁盘爆 | 变量没做死区 | 全局 deadband 1% 量程 |
| 数据有缺口 | 采样被降频 | 查 min_interval_ms |
| 冷层不可读 | S3 凭据错 | 查 Service 日志 |
| 新变量没入库 | 没标 archive=true | 变量表更新 + 重启 Service |
高级技巧
- 变量分级: 关键 10 点原始存, 其他 500 点 1min 聚合存
- 多租户隔离: 不同产线用独立 hypertable 或 schema
- 预聚合加速: 常用查询 ( 24h/avg ) 作连续聚合, 毫秒返回
- 冷热查询透明: API 层自动路由, 客户端不感知分层
- 批次追溯: 把批次号作为 tag, 按批次直接按标签过滤