时序原生
Time-Series Native 是 EnerOS 的核心设计理念:时序数据不是外部数据库的附属功能,而是操作系统内核的原生存储引擎。电力系统的所有运行状态(电压、电流、功率、频率、温度)本质上都是时序数据,EnerOS 将时序存储引擎内建为内核组件,提供纳秒级精度、百万级吞吐的写入与查询,无需依赖 InfluxDB、TimescaleDB、Prometheus 等外部时序数据库。
设计动机
传统时序方案的痛点
电力系统使用外部时序数据库存在以下问题:
| 痛点 | 描述 | 影响 |
|---|
| 网络延迟 | 应用层通过网络访问数据库 | 增加决策延迟 1-50ms |
| 部署复杂 | 需独立部署、运维数据库 | 运维成本高 |
| 数据不一致 | 内核状态与数据库状态异步 | 决策基于过时数据 |
| 双写问题 | 内核事件同时写日志和数据库 | 资源浪费 |
| 多租户限制 | 外部数据库多租户支持弱 | 难以隔离不同 Agent |
| 故障扩散 | 数据库故障影响所有应用 | 系统可用性下降 |
时序原生的解决思路
EnerOS 将时序存储引擎内建为内核组件,所有 Agent 通过系统调用访问:
┌─────────────────────────────────────────┐
│ Agent / 应用层 │
├─────────────────────────────────────────┤
│ 时序系统调用接口 (TimeSeriesSyscall) │
├─────────────────────────────────────────┤
│ 时序引擎 (TimeSeriesEngine) │
│ ├─ 写入路径 (WAL + LSM-Tree) │
│ ├─ 查询路径 (索引 + 向量化执行) │
│ ├─ 聚合引擎 (流式 + 预计算) │
│ └─ 降采样引擎 (LTTB / M4 / MinMax) │
├─────────────────────────────────────────┤
│ 存储后端 (Storage Backend) │
│ ├─ SQLite + WAL (默认) │
│ ├─ 内存缓存 (LRU) │
│ └─ 列式压缩 (Gorilla / Delta) │
└─────────────────────────────────────────┘
时序引擎架构
写入路径
测量数据 → 内存缓冲 → WAL 日志 → LSM-Tree → 压缩合并
↓ ↓ ↓
校验 索引更新 持久化
查询路径
查询请求 → 时间索引 → 数据加载 → 向量化过滤 → 聚合计算 → 结果返回
↓ ↓ ↓ ↓
B+树查找 列式读取 SIMD 加速 流式聚合
数据模型
Point(数据点)
时序数据的最小单元,包含测量名、时间戳、值和可选标签:
use eneros_timeseries::{Point, Timestamp, Value};
let point = Point::new(
"bus_1_voltage", // 测量名 (measurement)
Timestamp::now(), // 纳秒级时间戳
Value::Float(1.024), // 测量值
)
.with_tag("region", "east") // 标签
.with_tag("voltage_level", "110kV")
.with_tag("substation", "ss-001");
Series(时间序列)
一个测量名 + 一组标签的唯一组合构成一条时间序列:
use eneros_timeseries::{Series, SeriesKey};
let series_key = SeriesKey::new("bus_1_voltage")
.with_tag("region", "east")
.with_tag("voltage_level", "110kV");
let series = ts.get_series(&series_key)?;
println!("序列 ID: {}", series.id());
println!("数据点数: {}", series.count());
println!("时间范围: {} - {}", series.start(), series.end());
Tag(标签)
标签用于标识和过滤时间序列,支持多维查询:
| 标签键 | 示例值 | 用途 |
|---|
region | east / west / north | 调度区域 |
voltage_level | 500kV / 220kV / 110kV | 电压等级 |
substation | ss-001 / ss-002 | 变电站 |
device_type | breaker / transformer / line | 设备类型 |
phase | A / B / C / N | 相别 |
measurement_type | instantaneous / accumulated | 测量类型 |
写入 API
单点写入
use eneros_timeseries::{TimeSeriesEngine, Point, Value};
let ts = TimeSeriesEngine::open("data/eneros.db")?;
// 单点写入
ts.write(Point::new(
"bus_1_voltage",
Timestamp::now(),
Value::Float(1.024),
))?;
批量写入
// 批量写入(推荐,吞吐高 10 倍)
let batch: Vec<Point> = vec![
Point::new("bus_1_voltage", Timestamp::now(), Value::Float(1.024)),
Point::new("bus_1_frequency", Timestamp::now(), Value::Float(50.001)),
Point::new("bus_1_p", Timestamp::now(), Value::Float(45.5)),
Point::new("bus_1_q", Timestamp::now(), Value::Float(12.3)),
Point::new("bus_2_voltage", Timestamp::now(), Value::Float(0.998)),
Point::new("bus_2_frequency", Timestamp::now(), Value::Float(49.998)),
];
ts.batch_write(&batch)?;
流式写入
use eneros_timeseries::StreamWriter;
// 流式写入(适用于 SCADA 实时数据)
let mut writer = ts.stream_writer("bus_1_voltage")?;
for measurement in scada_stream {
writer.write(measurement.timestamp, measurement.value)?;
}
writer.flush()?; // 显式刷新
异步写入
use eneros_timeseries::AsyncWriter;
// 异步写入器(后台批量提交)
let writer = AsyncWriter::new(ts.clone())
.batch_size(1000) // 每批 1000 点
.flush_interval(Duration::milliseconds(100)); // 或 100ms 刷新
for point in measurement_stream {
writer.send(point).await?; // 立即返回,后台写入
}
writer.flush().await?; // 等待所有数据写入
查询 API
基础查询
use eneros_timeseries::{Query, Duration, Aggregation, Downsample};
// 1. 时间范围查询
let data: Vec<Point> = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::hours(24), now!())
.execute()?;
// 2. 多测量点查询
let multi: Vec<SeriesResult> = ts.query()
.points(vec!["bus_1_voltage", "bus_2_voltage", "bus_3_voltage"])
.range(now!() - Duration::hours(1), now!())
.execute()?;
// 3. 标签过滤查询
let filtered: Vec<Point> = ts.query()
.point("voltage")
.tag("region", "east")
.tag("voltage_level", "110kV")
.range(now!() - Duration::hours(1), now!())
.execute()?;
聚合查询
// 4. 聚合查询(15 分钟平均值)
let avg: Vec<Point> = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(7), now!())
.aggregation(Aggregation::Avg, Duration::minutes(15))
.execute()?;
// 5. 多聚合查询
let multi_agg: Vec<Point> = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(1), now!())
.aggregations(vec![
Aggregation::Avg,
Aggregation::Min,
Aggregation::Max,
Aggregation::StdDev,
], Duration::minutes(15))
.execute()?;
降采样查询
// 6. 降采样(取最近 1000 个点)
let downsampled: Vec<Point> = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(30), now!())
.downsample(Downsample::LTTB, 1000)
.execute()?;
// 7. 聚合 + 降采样组合
let combined: Vec<Point> = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(30), now!())
.aggregation(Aggregation::Avg, Duration::minutes(5))
.downsample(Downsample::M4, 2000)
.execute()?;
聚合函数
EnerOS 支持丰富的聚合函数:
| 聚合函数 | 说明 | 典型用途 |
|---|
Avg | 平均值 | 电压/频率均值 |
Min | 最小值 | 最低电压分析 |
Max | 最大值 | 最高负荷分析 |
Sum | 求和 | 累计电量 |
Count | 计数 | 数据完整率 |
StdDev | 标准差 | 波动分析 |
Percentile(p) | 分位数 | P99 延迟分析 |
First | 第一个值 | 状态初始值 |
Last | 最后一个值 | 当前状态 |
Median | 中位数 | 抗异常值均值 |
Range | 极差 | 波动范围 |
// 自定义聚合窗口
let result = ts.query()
.point("total_load")
.range(now!() - Duration::days(1), now!())
.aggregation(Aggregation::Percentile(99.0), Duration::minutes(15))
.execute()?;
降采样算法
| 算法 | 说明 | 适用场景 |
|---|
LTTB | Largest-Triangle-Three-Buckets | 可视化(保留趋势) |
M4 | Min-Min-Max-Max | K 线图 |
MinMax | 极值采样 | 保护逻辑(保留极值) |
Average | 平均采样 | 平滑曲线 |
First | 首点采样 | 状态序列 |
Every | 等间隔采样 | 简单降采样 |
// 可视化场景:LTTB 保留趋势
let viz_data = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(30), now!())
.downsample(Downsample::LTTB, 1000)
.execute()?;
// 保护场景:MinMax 保留极值
let protection_data = ts.query()
.point("bus_1_voltage")
.range(now!() - Duration::days(1), now!())
.downsample(Downsample::MinMax, 5000)
.execute()?;
数据保留策略
EnerOS 支持多级数据保留策略,自动管理数据生命周期:
use eneros_timeseries::{RetentionPolicy, TieredStorage};
let policy = RetentionPolicy::new()
// 原始数据:保留 7 天
.tier("raw", Duration::days(7), Resolution::Raw)
// 1 分钟聚合:保留 90 天
.tier("1min", Duration::days(90), Resolution::Minutes(1))
// 15 分钟聚合:保留 1 年
.tier("15min", Duration::days(365), Resolution::Minutes(15))
// 1 小时聚合:永久保留
.tier("1hour", Duration::infinity(), Resolution::Hours(1));
ts.set_retention_policy("bus_.*_voltage", policy)?;
保留策略配置
# config/timeseries.toml
[[retention]]
pattern = "bus_.*_voltage" # 测量名匹配
raw_days = 7 # 原始数据保留天数
tiers = [
{ name = "1min", days = 90, resolution = "1min" },
{ name = "15min", days = 365, resolution = "15min" },
{ name = "1hour", days = -1, resolution = "1hour" }, # -1 = 永久
]
[[retention]]
pattern = ".*_frequency"
raw_days = 30
tiers = [
{ name = "1sec", days = 7, resolution = "1sec" },
{ name = "1min", days = 365, resolution = "1min" },
]
[[retention]]
pattern = ".*_power"
raw_days = 14
tiers = [
{ name = "1min", days = 180, resolution = "1min" },
{ name = "15min", days = 1825, resolution = "15min" }, # 5 年
]
与 Agent 集成
Agent 直接查询时序
use eneros_timeseries::{Duration, Aggregation};
impl Agent for ForecastAgent {
async fn run(&mut self, ctx: &mut AgentContext) -> AgentResult<()> {
// 1. 查询历史负荷数据(通过系统调用)
let history = ctx.query_timeseries()
.point("total_load")
.range(now!() - Duration::days(7), now!())
.aggregation(Aggregation::Avg, Duration::minutes(15))
.execute()?;
// 2. 查询天气数据
let weather = ctx.query_timeseries()
.point("temperature")
.range(now!() - Duration::days(7), now!())
.aggregation(Aggregation::Avg, Duration::hours(1))
.execute()?;
// 3. 基于历史数据预测
let forecast = self.forecast(&history, &weather)?;
// 4. 写入预测结果
for point in forecast {
ctx.write_timeseries(
"forecast_load",
point.timestamp,
point.value,
).await?;
}
Ok(())
}
}
实时订阅
Agent 可以订阅时序数据流,实现事件驱动:
use eneros_timeseries::Subscription;
impl Agent for ProtectionAgent {
async fn run(&mut self, ctx: &mut AgentContext) -> AgentContext<()> {
// 订阅电压越限事件
let mut sub = ctx.subscribe_timeseries()
.point("bus_1_voltage")
.condition(|v| v < 0.95 || v > 1.05)
.build()?;
while let Some(event) = sub.recv().await {
log::warn!("电压越限: {:?}", event);
self.handle_voltage_violation(event).await?;
}
Ok(())
}
}
存储后端
默认后端:SQLite + WAL
use eneros_timeseries::StorageConfig;
let config = StorageConfig::sqlite("data/eneros.db")
.wal_mode(true) // 启用 WAL
.cache_size_mb(256) // 缓存 256MB
.page_size(4096) // 页大小
.checkpoint_interval(Duration::minutes(5));
let ts = TimeSeriesEngine::open_with_config(config)?;
内存后端
// 全内存存储(适用于测试或临时数据)
let config = StorageConfig::memory()
.max_size_mb(512);
let ts = TimeSeriesEngine::open_with_config(config)?;
压缩算法
| 算法 | 适用数据类型 | 压缩率 | 解压速度 |
|---|
| Gorilla | 浮点数(电压/电流) | 8-12x | 极快 |
| Delta-of-Delta | 时间戳 | 10-20x | 极快 |
| RLE | 状态量(开关位置) | 100x+ | 极快 |
| Snappy | 通用 | 2-4x | 快 |
| LZ4 | 通用 | 2-5x | 极快 |
性能指标
写入性能
| 写入模式 | 吞吐 | 延迟 |
|---|
| 单点同步 | 10万点/秒 | < 10μs |
| 批量同步(1000点) | 100万点/秒 | < 1ms |
| 异步批量 | 200万点/秒 | < 1μs(入队) |
| 流式写入 | 100万点/秒 | < 1μs |
查询性能
| 查询类型 | 数据量 | 延迟 |
|---|
| 单点查询 | 1 点 | < 1μs |
| 时间范围查询 | 1万点 | < 5ms |
| 时间范围查询 | 100万点 | < 100ms |
| 聚合查询 | 1000万点 | < 50ms |
| 降采样查询 | 1亿点 → 1000点 | < 200ms |
| 多序列查询 | 100序列 × 1万点 | < 50ms |
存储指标
| 指标 | 数值 |
|---|
| 原始数据大小 | 1.2 GB/亿点 |
| Gorilla 压缩后 | 100-150 MB/亿点 |
| 压缩率 | 8-12x |
| 索引大小 | 50 MB/亿点 |
与外部数据库对比
与 InfluxDB 对比
| 维度 | InfluxDB | EnerOS Time-Series Native |
|---|
| 部署方式 | 独立进程 | 内核组件 |
| 访问方式 | HTTP/gRPC | 系统调用 |
| 写入延迟 | 1-5ms(网络往返) | < 1μs(内核态) |
| 查询延迟 | 5-50ms | < 5ms |
| 吞吐 | 50万点/秒 | 100万点/秒 |
| 多租户 | 弱 | 内核原生 |
| 数据一致性 | 最终一致 | 强一致(与内核状态) |
| 运维成本 | 高(独立运维) | 零(内核管理) |
| 故障影响 | 独立故障 | 与内核同生命周期 |
与 TimescaleDB 对比
| 维度 | TimescaleDB | EnerOS Time-Series Native |
|---|
| 基础 | PostgreSQL 扩展 | 内核组件 |
| SQL 支持 | 完整 SQL | 查询 DSL |
| 写入吞吐 | 30万点/秒 | 100万点/秒 |
| 查询延迟 | 10-100ms | < 5ms |
| 事务支持 | ACID | ACID |
| 部署复杂 | 高(需 PostgreSQL) | 零 |
| 内存占用 | 高(500MB+) | 低(50MB) |
与 Prometheus 对比
| 维度 | Prometheus | EnerOS Time-Series Native |
|---|
| 设计目标 | 监控指标 | 电力时序数据 |
| 数据模型 | 指标 + 标签 | 测量 + 标签 |
| 写入方式 | Pull 模型 | Push(内核直写) |
| 精度 | 毫秒 | 纳秒 |
| 保留策略 | 固定窗口 | 多级分层 |
| 聚合能力 | PromQL | 内置聚合 + 降采样 |
| 适用场景 | IT 监控 | 电力系统 |
完整使用示例
use eneros_timeseries::{TimeSeriesEngine, Point, Value, Aggregation, Duration, Downsample};
fn main() -> Result<(), Box<dyn std::error::Error>> {
// 1. 打开时序引擎
let ts = TimeSeriesEngine::open("data/eneros.db")?;
// 2. 配置保留策略
let policy = eneros_timeseries::RetentionPolicy::new()
.tier("raw", Duration::days(7), eneros_timeseries::Resolution::Raw)
.tier("15min", Duration::days(365), eneros_timeseries::Resolution::Minutes(15));
ts.set_retention_policy("bus_.*", policy)?;
// 3. 写入测量数据
let now = eneros_timeseries::Timestamp::now();
ts.batch_write(&[
Point::new("bus_1_voltage", now, Value::Float(1.024)),
Point::new("bus_1_frequency", now, Value::Float(50.001)),
Point::new("bus_1_p", now, Value::Float(45.5)),
Point::new("bus_1_q", now, Value::Float(12.3)),
])?;
// 4. 查询最近 1 小时电压
let recent = ts.query()
.point("bus_1_voltage")
.range(now - Duration::hours(1), now)
.execute()?;
println!("最近 1 小时电压点数: {}", recent.len());
// 5. 聚合查询:过去 7 天的 15 分钟均值
let avg = ts.query()
.point("bus_1_voltage")
.range(now - Duration::days(7), now)
.aggregation(Aggregation::Avg, Duration::minutes(15))
.execute()?;
println!("过去 7 天 15 分钟均值点数: {}", avg.len());
// 6. 降采样:过去 30 天数据降到 1000 点
let downsampled = ts.query()
.point("bus_1_voltage")
.range(now - Duration::days(30), now)
.downsample(Downsample::LTTB, 1000)
.execute()?;
println!("降采样后点数: {}", downsampled.len());
// 7. 统计信息
let stats = ts.stats("bus_1_voltage")?;
println!("序列统计: {:?}", stats);
Ok(())
}
限制与权衡
| 权衡点 | 说明 | 缓解策略 |
|---|
| 单机存储 | 默认不分布式 | 支持多区域同步 |
| SQL 不支持 | 使用查询 DSL | 提供丰富 API |
| 内存占用 | 内核常驻缓存 | 可配置上限 |
| 历史数据迁移 | 跨引擎迁移需工具 | 提供 CSV/Parquet 导入导出 |
| 跨节点查询 | 需要多区域模块 | eneros-multiregion 支持 |
下一步