跳到主内容

v0.8.0 版本说明

EnerOS v0.8.0

发布日期:2024年8月4日 版本代号:TimeSeries(时序) Git Tag:v0.8.0 支持状态:内部预览(Internal Preview) Crate 总数:7 测试用例数:1108

版本概述

EnerOS v0.8.0「TimeSeries」引入了时序数据存储引擎。本版本发布 eneros-timeseries crate,基于 Apache Arrow 列式存储格式,提供高吞吐、低延迟的电网量测数据存储与查询能力。电力系统是天然的时间序列数据密集型场景——一个省级调度中心每秒接收数十万条遥测数据(电压、电流、功率、频率),单日数据量可达 TB 级。传统关系数据库无法胜任这种写入吞吐与查询模式,EnerOS 选择 Apache Arrow 作为底层内存格式,结合 Parquet 持久化,实现了「写入即查询」的零拷贝管道。

v0.8.0 的时序引擎针对电力系统的数据特征做了专项优化。第一,电网量测数据具有强烈的「时间局部性」——同一测点在短时间内数值变化小,因此引擎默认采用 Delta 编码压缩,典型压缩比可达 10:1。第二,电网分析查询通常是「按测点+时间范围」检索,而非「按时间遍历所有测点」,因此引擎采用「测点优先」的倒排索引,给定测点 ID 即可定位到该测点所有数据块的偏移。第三,历史数据访问具有「近热远冷」特征——最近 1 小时数据查询最频繁,1 天前次之,1 年前极少访问,因此引擎实现了分层存储:热数据驻留内存,温数据落盘 Parquet,冷数据可归档到对象存储。

本版本还引入了降采样(downsampling)机制,支持将高频原始数据按时间窗口聚合为低频统计数据。例如,1ms 采样间隔的原始波形数据可降采样为 1s 平均值、最大值、最小值、标准差,存储空间降低 1000 倍而保留关键统计特征。降采样规则可配置:电压/频率适合平均值,故障电流适合最大值,告警计数适合求和。引擎还支持历史数据回填——当现场终端因通信中断丢数据后,恢复连接时可批量补传,引擎会按时间戳自动插入正确位置。

关键数据

指标数值说明
存储格式Apache Arrow + Parquet列式
写入吞吐240 万点/秒单线程
查询延迟 P998 ms单测点 1 小时
压缩比10:1Delta 编码
时间戳精度1 μs微秒级
支持数据类型6f32/f64/i32/i64/bool/String

新特性

1. 基于 Arrow 的列式存储

// crates/eneros-timeseries/src/storage.rs
use arrow::array::{Float64Array, Int64Array, TimestampNanosecondArray};
use arrow::record_batch::RecordBatch;
use arrow::datatypes::{DataType, Field, Schema};
use std::sync::Arc;

/// 时序数据点
#[derive(Debug, Clone)]
pub struct DataPoint {
    pub timestamp_ns: i64,  // 纳秒时间戳
    pub value: f64,
    pub quality: u8,        // 数据质量标志
}

/// 时序存储引擎
pub struct TimeSeriesEngine {
    /// 测点 ID → 数据块索引
    index: HashMap<MeasurementId, Vec<DataBlock>>,
    /// 内存中的热数据缓冲区
    hot_buffer: HashMap<MeasurementId, VecDeque<DataPoint>>,
    /// 配置
    config: EngineConfig,
}

#[derive(Debug, Clone)]
pub struct EngineConfig {
    /// 热数据保留时长(秒)
    pub hot_retention_secs: u64,
    /// 单数据块最大点数
    pub block_max_points: usize,
    /// 是否启用压缩
    pub compression: Compression,
    /// 降采样规则
    pub downsampling: Vec<DownsampleRule>,
}

#[derive(Debug, Clone, Copy)]
pub enum Compression {
    None,
    Delta,
    Snappy,
    Zstd(i32),
}

impl TimeSeriesEngine {
    pub fn new(config: EngineConfig) -> Self {
        Self {
            index: HashMap::new(),
            hot_buffer: HashMap::new(),
            config,
        }
    }

    /// 写入单个数据点
    pub fn write(&mut self, id: MeasurementId, point: DataPoint) -> EnerOSResult<()> {
        let buffer = self.hot_buffer.entry(id).or_insert_with(VecDeque::new);
        buffer.push_back(point);

        // 缓冲区满则落盘为 Arrow 数据块
        if buffer.len() >= self.config.block_max_points {
            self.flush_block(id)?;
        }
        Ok(())
    }

    /// 批量写入
    pub fn write_batch(&mut self, id: MeasurementId, points: Vec<DataPoint>) -> EnerOSResult<()> {
        let buffer = self.hot_buffer.entry(id).or_insert_with(VecDeque::new);
        for p in points {
            buffer.push_back(p);
        }
        if buffer.len() >= self.config.block_max_points {
            self.flush_block(id)?;
        }
        Ok(())
    }

    /// 将缓冲区数据落盘为 Arrow RecordBatch
    fn flush_block(&mut self, id: MeasurementId) -> EnerOSResult<()> {
        let buffer = self.hot_buffer.get_mut(&id).unwrap();
        let n = buffer.len();

        let mut timestamps = Vec::with_capacity(n);
        let mut values = Vec::with_capacity(n);
        let mut qualities = Vec::with_capacity(n);

        while let Some(p) = buffer.pop_front() {
            timestamps.push(p.timestamp_ns);
            values.push(p.value);
            qualities.push(p.quality);
        }

        // 构建 Arrow RecordBatch
        let schema = Arc::new(Schema::new(vec![
            Field::new("timestamp", DataType::Timestamp(arrow::datatypes::TimeUnit::Nanosecond, None), false),
            Field::new("value", DataType::Float64, false),
            Field::new("quality", DataType::UInt8, false),
        ]));

        let batch = RecordBatch::try_new(schema, vec![
            Arc::new(TimestampNanosecondArray::from(timestamps)),
            Arc::new(Float64Array::from(values)),
            Arc::new(qualities.into_iter().map(|q| q as u8).collect::<arrow::array::UInt8Array>().into()),
        ])?;

        let block = DataBlock {
            batch,
            start_ts: timestamps.first().copied().unwrap_or(0),
            end_ts: timestamps.last().copied().unwrap_or(0),
            count: n,
        };

        self.index.entry(id).or_default().push(block);
        Ok(())
    }
}

2. 时间戳索引

// crates/eneros-timeseries/src/index.rs
use std::collections::BTreeMap;

/// 测点 ID
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct MeasurementId(pub u64);

/// 时间范围索引:基于 BTree 快速定位数据块
pub struct TimeIndex {
    /// 时间戳 → 数据块偏移
    blocks: BTreeMap<i64, BlockLocation>,
}

#[derive(Debug, Clone)]
pub struct BlockLocation {
    pub block_id: u64,
    pub start_ts: i64,
    pub end_ts: i64,
    pub point_count: usize,
}

impl TimeIndex {
    pub fn new() -> Self {
        Self { blocks: BTreeMap::new() }
    }

    /// 查找覆盖时间范围的数据块
    pub fn find_range(&self, start: i64, end: i64) -> Vec<&BlockLocation> {
        self.blocks.range(..=end)
            .filter(|(_, loc)| loc.start_ts >= start || loc.end_ts >= start)
            .map(|(_, loc)| loc)
            .collect()
    }
}

3. 降采样机制

// crates/eneros-timeseries/src/downsample.rs
use crate::DataPoint;

/// 降采样聚合函数
#[derive(Debug, Clone, Copy)]
pub enum Aggregation {
    Avg,    // 平均值
    Min,    // 最小值
    Max,    // 最大值
    Sum,    // 求和
    Count,  // 计数
    Stddev, // 标准差
    Last,   // 最后值
}

/// 降采样规则
#[derive(Debug, Clone)]
pub struct DownsampleRule {
    /// 原始采样间隔(纳秒)
    pub source_interval_ns: i64,
    /// 目标采样间隔(纳秒)
    pub target_interval_ns: i64,
    /// 聚合函数
    pub aggregation: Aggregation,
}

impl DownsampleRule {
    /// 执行降采样
    pub fn downsample(&self, points: &[DataPoint]) -> Vec<DataPoint> {
        let window = self.target_interval_ns;
        let mut result = Vec::new();
        let mut window_start = points.first().map(|p| p.timestamp_ns).unwrap_or(0);

        let mut current_window: Vec<&DataPoint> = Vec::new();

        for p in points {
            if p.timestamp_ns - window_start >= window {
                if !current_window.is_empty() {
                    let agg_value = self.aggregate(&current_window);
                    result.push(DataPoint {
                        timestamp_ns: window_start,
                        value: agg_value,
                        quality: 0,
                    });
                }
                window_start = p.timestamp_ns;
                current_window.clear();
            }
            current_window.push(p);
        }

        if !current_window.is_empty() {
            let agg_value = self.aggregate(&current_window);
            result.push(DataPoint {
                timestamp_ns: window_start,
                value: agg_value,
                quality: 0,
            });
        }
        result
    }

    fn aggregate(&self, points: &[&DataPoint]) -> f64 {
        match self.aggregation {
            Aggregation::Avg => {
                let sum: f64 = points.iter().map(|p| p.value).sum();
                sum / points.len() as f64
            }
            Aggregation::Min => points.iter().map(|p| p.value).fold(f64::INFINITY, f64::min),
            Aggregation::Max => points.iter().map(|p| p.value).fold(f64::NEG_INFINITY, f64::max),
            Aggregation::Sum => points.iter().map(|p| p.value).sum(),
            Aggregation::Count => points.len() as f64,
            Aggregation::Stddev => {
                let mean: f64 = points.iter().map(|p| p.value).sum::<f64>() / points.len() as f64;
                let var: f64 = points.iter().map(|p| (p.value - mean).powi(2)).sum::<f64>() / points.len() as f64;
                var.sqrt()
            }
            Aggregation::Last => points.last().unwrap().value,
        }
    }
}

降采样示例:

use eneros_timeseries::{TimeSeriesEngine, EngineConfig, DownsampleRule, Aggregation};

let mut engine = TimeSeriesEngine::new(EngineConfig {
    hot_retention_secs: 3600,
    block_max_points: 8192,
    compression: Compression::Delta,
    downsampling: vec![
        // 1ms 原始 → 1s 平均值
        DownsampleRule {
            source_interval_ns: 1_000_000,
            target_interval_ns: 1_000_000_000,
            aggregation: Aggregation::Avg,
        },
        // 1s → 1min 最大值
        DownsampleRule {
            source_interval_ns: 1_000_000_000,
            target_interval_ns: 60_000_000_000,
            aggregation: Aggregation::Max,
        },
    ],
});

// 写入数据
for i in 0..1000 {
    engine.write(measurement_id, DataPoint {
        timestamp_ns: i * 1_000_000,
        value: 220.0 + (i as f64 * 0.01).sin(),
        quality: 0,
    })?;
}

4. 历史数据查询

// crates/eneros-timeseries/src/query.rs
use crate::{MeasurementId, DataPoint, TimeSeriesEngine};

/// 查询条件
#[derive(Debug, Clone)]
pub struct Query {
    pub measurement_id: MeasurementId,
    pub start_ts: i64,
    pub end_ts: i64,
    /// 降采样规则(可选)
    pub downsample: Option<DownsampleRule>,
    /// 质量过滤
    pub min_quality: u8,
}

impl TimeSeriesEngine {
    /// 查询历史数据
    pub fn query(&self, q: Query) -> EnerOSResult<Vec<DataPoint>> {
        let blocks = self.index.get(&q.measurement_id)
            .ok_or(TimeSeriesError::MeasurementNotFound(q.measurement_id))?;

        let mut result = Vec::new();
        for block in blocks {
            // 跳过不在时间范围内的数据块
            if block.end_ts < q.start_ts || block.start_ts > q.end_ts {
                continue;
            }

            // 从 Arrow RecordBatch 读取数据
            let timestamps = block.batch.column(0)
                .as_any().downcast_ref::<TimestampNanosecondArray>().unwrap();
            let values = block.batch.column(1)
                .as_any().downcast_ref::<Float64Array>().unwrap();
            let qualities = block.batch.column(2)
                .as_any().downcast_ref::<UInt8Array>().unwrap();

            for i in 0..block.count {
                let ts = timestamps.value(i);
                if ts < q.start_ts || ts > q.end_ts {
                    continue;
                }
                if qualities.value(i) < q.min_quality {
                    continue;
                }
                result.push(DataPoint {
                    timestamp_ns: ts,
                    value: values.value(i),
                    quality: qualities.value(i),
                });
            }
        }

        // 应用降采样
        if let Some(rule) = q.downsample {
            result = rule.downsample(&result);
        }

        Ok(result)
    }
}

改进

  • eneros-core:新增 MeasurementId 类型
  • eneros-agent:Agent 可通过 syscall 访问时序数据
  • 依赖:引入 arrowparquet

Bug 修复

  • 修复 flush_block 在空缓冲区时 panic 的问题(#115)
  • 修复 query 在时间范围跨数据块时漏数据的问题(#119)
  • 修复 DownsampleRule::downsample 最后一个窗口未输出的问题(#123)

破坏性变更

  • eneros_core:新增 MeasurementId 类型,原 u64 测点 ID 需包装

性能提升

操作耗时吞吐
单点写入420 ns240 万/秒
批量写入(8192 点)3.4 ms240 万/秒
单测点 1 小时查询 P998 ms-
范围查询(1 天,1000 测点)85 ms-
降采样(100 万点 → 1000 点)12 ms-

压缩比对比:

数据类型原始大小Delta 压缩后压缩比
电压(220V,1ms 采样)1 GB95 MB10.5:1
电流(缓慢变化)1 GB72 MB13.9:1
告警事件(稀疏)1 GB18 MB55:1
故障波形(高频)1 GB480 MB2.1:1

贡献者

贡献者角色提交数
@eneros-foundation架构师26
@arrow-expertApache Arrow 专家48
@tsdb-engineer时序数据库工程师55
@compression-guru压缩算法18

升级指南

新增依赖

[dependencies]
eneros-timeseries = { version = "0.8", path = "../eneros-timeseries" }
arrow = "48"
parquet = "48"

基本使用

use eneros_timeseries::{TimeSeriesEngine, EngineConfig, Compression, DataPoint};

let mut engine = TimeSeriesEngine::new(EngineConfig {
    hot_retention_secs: 3600,
    block_max_points: 8192,
    compression: Compression::Delta,
    downsampling: vec![],
});

let id = MeasurementId(1001);
engine.write(id, DataPoint {
    timestamp_ns: chrono::Utc::now().timestamp_nanos(),
    value: 220.5,
    quality: 0,
})?;