时序原生操作
EnerOS 将时序数据作为内核原生数据类型,无需外挂时序数据库即可支撑百万点/秒的写入与毫秒级查询。这是「时序原生」理念的直接体现:电网的电压、电流、功率天然是时序数据,操作系统应当原生理解而非通过外部插件支持。
电力系统是「时间密集型」领域:PMU 数据以 100-120 Hz 采样,故障录波达到 10 kHz 以上,状态估计每 100ms 一次。传统架构依赖外部 InfluxDB / TimescaleDB,引入网络往返与序列化开销。EnerOS 将时序引擎内嵌到内核,与拓扑、潮流共享同一内存空间,写入路径短至 1μs。
架构总览
┌─────────────────────────────────────────────────────┐
│ Agent / 应用 / 数字孪生 │
└──────────────────┬──────────────────────────────────┘
│ write / query / aggregate
┌──────────────────┴──────────────────────────────────┐
│ eneros-timeseries (内核时序引擎) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Writer │ │ Reader │ │ Compactor│ │
│ └──────────┘ └──────────┘ └──────────┘ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ CQ 引擎 │ │ Retention│ │ Downsamp │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└──────────────────┬──────────────────────────────────┘
│ 列式存储 + Gorilla 编码
┌──────────────────┴──────────────────────────────────┐
│ LSM-Tree (per-tenant) │
│ MemTable → SSTable → Compaction │
└─────────────────────────────────────────────────────┘
数据模型
时序数据由三要素构成:测量点(tag)、时间戳、字段值。
use eneros_timeseries::{Series, Sample, Timestamp, Quality};
let mut series = Series::new("bus_1.voltage");
series.push(Sample::new(ts1, 1.02));
series.push(Sample::new(ts2, 1.025));
// 带质量位的样本
series.push(Sample::with_quality(ts3, 1.03, Quality::Good));
series.push(Sample::with_quality(ts4, 0.98, Quality::Suspect));
Sample 字段说明
| 字段 | 类型 | 含义 | 备注 |
|---|---|---|---|
| timestamp | Timestamp | 纳秒精度时间戳 | UTC |
| value | f64 | 测量值 | 标量 |
| quality | Quality | 数据质量位 | Good / Suspect / Bad |
| flags | u8 | 用户自定义标志 | 0-255 |
Quality 枚举
| 变体 | 含义 | 数值 |
|---|---|---|
| Good | 数据有效 | 0 |
| Suspect | 数据可疑(如传感器告警) | 1 |
| Bad | 数据无效 | 2 |
| Missing | 数据缺失(补全) | 3 |
| Calculated | 计算得到(虚拟传感器) | 4 |
写入
支持单点、批量与流式写入:
use eneros_timeseries::TimeSeriesStore;
let store = TimeSeriesStore::open("/var/eneros/tsdb")?;
// 单点写入
store.write("bus_1.voltage", ts, 1.02).await?;
// 带质量位写入
store.write_with_quality("bus_1.voltage", ts, 1.02, Quality::Good).await?;
// 批量写入(推荐,吞吐高 5x)
let batch = vec![
Sample::new(ts1, 1.02),
Sample::new(ts2, 1.025),
Sample::new(ts3, 1.03),
];
store.write_batch("bus_1.voltage", batch).await?;
// 多测点批量写入
let multi_batch = vec![
("bus_1.voltage", Sample::new(ts, 1.02)),
("bus_1.frequency", Sample::new(ts, 50.01)),
("bus_1.load", Sample::new(ts, 50.0)),
];
store.write_multi(multi_batch).await?;
// 流式写入(用于实时采集)
let mut stream = store.stream("bus_1.voltage");
stream.send(Sample::now(1.02)).await?;
stream.flush().await?;
写入 API 一览
| 方法 | 入参 | 适用场景 | 吞吐 |
|---|---|---|---|
| write | tag, ts, value | 单点写入 | 2M samples/s |
| write_with_quality | tag, ts, value, q | 带质量位 | 2M samples/s |
| write_batch | tag, Vec | 批量写入 | 5M samples/s |
| write_multi | Vec<(tag, Sample)> | 多测点批量 | 5M samples/s |
| stream | tag | 流式写入 | 5M samples/s |
查询
use eneros_timeseries::{Query, Range, Aggregation, Filter};
let result = store.query(
Query::select("bus_1.voltage")
.range(Range::last(Duration::hours(1)))
.downsample(Duration::seconds(1), Aggregation::Avg)
.filter(Filter::quality(Quality::Good))
).await?;
for sample in result {
println!("{}: {:.4}", sample.timestamp, sample.value);
}
多测点查询
let multi = store.query(
Query::select_multi(&["bus_1.voltage", "bus_1.frequency", "bus_1.load"])
.range(Range::between(t_start, t_end))
.align(AlignStrategy::Interpolate)
).await?;
for (tag, samples) in multi {
println!("测点 {}: {} 条样本", tag, samples.len());
}
Range 类型
| 变体 | 含义 | 示例 |
|---|---|---|
| Last(d) | 最近 d 时长 | 最近 1 小时 |
| Between(t1, t2) | 闭区间 | [t1, t2] |
| Since(t) | 自 t 至今 | 自上次重启 |
| Relative(offset) | 相对当前 | 偏移 -1h |
| All | 全部历史 | 全量回放 |
Filter 类型
| 变体 | 含义 |
|---|---|
| quality(q) | 按质量位过滤 |
| value_range(min, max) | 按值范围过滤 |
| anomaly_only | 仅异常值 |
| time_mask(mask) | 按时间掩码(如工作时段) |
聚合与降采样
内置常用聚合函数,并支持滚动窗口与滑动窗口:
// 整体聚合
let avg = store.aggregate(
"bus_1.voltage",
Range::last(Duration::hours(24)),
Aggregation::Avg,
).await?;
// 滚动窗口聚合
let windowed = store.window_aggregate(
"bus_1.load",
Range::last(Duration::hours(24)),
Duration::hours(1),
Aggregation::Percentile(95),
).await?;
for window in &windowed {
println!("{} ~ {}: P95 = {:.2}", window.start, window.end, window.value);
}
// 滑动窗口聚合
let sliding = store.sliding_window(
"bus_1.load",
Range::last(Duration::hours(1)),
Duration::minutes(15), // 窗口大小
Duration::minutes(5), // 步长
Aggregation::StdDev,
).await?;
| 聚合函数 | 说明 | 适用场景 |
|---|---|---|
| Avg | 算术平均 | 电压、功率平均 |
| Sum | 求和 | 电度量累计 |
| Min / Max | 极值 | 极端电压检测 |
| Percentile(p) | 分位数 | P95/P99 负荷 |
| Count | 计数 | 故障次数 |
| First / Last | 首尾值 | 状态变化 |
| StdDev | 标准差 | 波动分析 |
| Rate | 变化率 | 频率变化率 |
| Integral | 积分 | 电度量(Wh) |
| Median | 中位数 | 抗干扰统计 |
连续查询
注册自动执行的连续查询,避免重复计算:
use eneros_timeseries::ContinuousQuery;
// SQL 风格 CQ
store.register_cq(
"voltage_5min_avg",
"SELECT avg(value) INTO bus_1.voltage_5m FROM bus_1.voltage GROUP BY time(5m)",
).await?;
// 类型化 CQ
let cq = ContinuousQuery::new("load_p95_hourly")
.source("bus_1.load")
.into("bus_1.load_p95_1h")
.aggregate(Aggregation::Percentile(95))
.window(Duration::hours(1))
.build();
store.register_cq_typed(cq).await?;
// 列出所有 CQ
let cqs = store.list_cqs().await?;
for cq in &cqs {
println!("{}: {} (last_run: {:?})", cq.name, cq.sql, cq.last_run);
}
// 删除 CQ
store.drop_cq("voltage_5min_avg").await?;
数据保留策略
支持按测量点配置保留期与精度降级:
use eneros_timeseries::RetentionPolicy;
store.set_retention("bus_1.voltage",
RetentionPolicy::new()
.raw_for(Duration::days(7)) // 原始数据保留 7 天
.downsampled(Duration::seconds(1), Duration::days(90)) // 1s 粒度保留 90 天
.downsampled(Duration::minutes(1), Duration::years(5)) // 1min 粒度保留 5 年
.downsampled(Duration::hours(1), Duration::years(20)) // 1h 粒度保留 20 年
).await?;
// 批量配置
let policy = RetentionPolicy::default()
.raw_for(Duration::days(30));
store.set_retention_multi(&["bus_1.*", "bus_2.*"], policy).await?;
RetentionPolicy 层级
| 层级 | 粒度 | 保留期 | 用途 |
|---|---|---|---|
| Raw | 原始 | 7 天 | 故障录波、实时控制 |
| High | 1s | 90 天 | 短期分析 |
| Medium | 1min | 5 年 | 月度报表 |
| Low | 1h | 20 年 | 长期趋势、规划 |
异常检测
内置时序异常检测算法,触发事件通知 Agent:
use eneros_timeseries::{AnomalyDetector, AnomalyKind};
let detector = AnomalyDetector::new("bus_1.voltage")
.method(AnomalyKind::Spike { threshold: 3.0 }) // 3σ 突变
.method(AnomalyKind::LevelShift { window: 60 }) // 60s 均值漂移
.method(AnomalyKind::Flatline { window: 30 }) // 30s 平直
.on_anomaly(|event| {
log::warn!("异常: {:?} at {}", event.kind, event.timestamp);
});
store.register_detector(detector).await?;
存储引擎
底层采用列式存储 + Gorilla 编码 + Delta-of-Delta 时间戳压缩,单节点可承载 10 年电网历史数据。
压缩算法
| 数据类型 | 算法 | 压缩比 |
|---|---|---|
| 时间戳 | Delta-of-Delta | 12x |
| 浮点值 | Gorilla XOR | 8-15x |
| 字符串 | Dictionary | 5-10x |
| 整数 | Simple8b | 6x |
| 布尔 | Bit-packing | 32x |
LSM-Tree 结构
MemTable (内存)
│ flush
▼
SSTable L0 (磁盘,10MB)
│ compaction
▼
SSTable L1 (磁盘,100MB)
│ compaction
▼
SSTable L2+ (磁盘,1GB+)
性能指标
| 操作 | 吞吐 | 延迟 |
|---|---|---|
| 单点写入 | 2M samples/s | < 1μs |
| 批量写入(10k 批) | 5M samples/s | < 50μs/批 |
| 流式写入 | 5M samples/s | < 1μs |
| 范围查询(1 天 1s 粒度) | - | < 5ms |
| 范围查询(1 年 1min 粒度) | - | < 50ms |
| 聚合查询(24h Avg) | - | < 10ms |
| 多测点查询(100 测点) | - | < 20ms |
| 连续查询执行 | - | < 100ms |
| 压缩比 | 8-15x | - |
| 单节点存储上限 | 10 年电网数据 | - |
测试环境:4 核 / 8GB / Ubuntu 22.04,NVMe SSD。
配置参数
eneros.toml 中时序相关配置:
[timeseries]
# 存储路径
data_dir = "/var/eneros/tsdb"
# MemTable 大小(MB)
memtable_mb = 64
# SSTable 压缩算法:snappy / lz4 / zstd
compression = "lz4"
# Compaction 线程数
compaction_threads = 2
# 默认保留期(天)
default_retention_days = 365
# 是否启用连续查询
enable_cq = true
# CQ 执行间隔(秒)
cq_interval_secs = 10
# 是否启用异常检测
enable_anomaly = true
# 写入缓存大小(MB)
write_buffer_mb = 128
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| data_dir | String | /var/eneros/tsdb | 存储路径 |
| memtable_mb | u32 | 64 | MemTable 大小 |
| compression | enum | lz4 | SSTable 压缩算法 |
| compaction_threads | u32 | 2 | Compaction 线程数 |
| default_retention_days | u32 | 365 | 默认保留期 |
| enable_cq | bool | true | 启用连续查询 |
| cq_interval_secs | u32 | 10 | CQ 执行间隔 |
| enable_anomaly | bool | true | 启用异常检测 |
| write_buffer_mb | u32 | 128 | 写入缓存大小 |
与其他能力的关系
| 关联能力 | 互动方式 |
|---|---|
| 数字孪生引擎 | 时序数据驱动孪生回放 |
| 物联网泛在接入 | IoT 设备时序采集 |
| 电网分析进阶 | 状态估计依赖时序数据 |
| 多智能体协作 | Agent 记忆基于时序数据 |
| 多租户与隔离 | 时序数据按租户 LSM 隔离 |
| 安全守卫 | 命令审计日志入库 |
| 实时双执行域 | PMU 数据流式写入 |