跳到主内容

时序原生

核心概念

时序原生

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(标签)

标签用于标识和过滤时间序列,支持多维查询:

标签键示例值用途
regioneast / west / north调度区域
voltage_level500kV / 220kV / 110kV电压等级
substationss-001 / ss-002变电站
device_typebreaker / transformer / line设备类型
phaseA / B / C / N相别
measurement_typeinstantaneous / 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()?;

降采样算法

算法说明适用场景
LTTBLargest-Triangle-Three-Buckets可视化(保留趋势)
M4Min-Min-Max-MaxK 线图
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 对比

维度InfluxDBEnerOS Time-Series Native
部署方式独立进程内核组件
访问方式HTTP/gRPC系统调用
写入延迟1-5ms(网络往返)< 1μs(内核态)
查询延迟5-50ms< 5ms
吞吐50万点/秒100万点/秒
多租户内核原生
数据一致性最终一致强一致(与内核状态)
运维成本高(独立运维)零(内核管理)
故障影响独立故障与内核同生命周期

与 TimescaleDB 对比

维度TimescaleDBEnerOS Time-Series Native
基础PostgreSQL 扩展内核组件
SQL 支持完整 SQL查询 DSL
写入吞吐30万点/秒100万点/秒
查询延迟10-100ms< 5ms
事务支持ACIDACID
部署复杂高(需 PostgreSQL)
内存占用高(500MB+)低(50MB)

与 Prometheus 对比

维度PrometheusEnerOS 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 支持

下一步