跳到主内容

v0.14.0 版本说明

EnerOS v0.14.0

发布日期:2025年01月26日 版本代号:Pipeline Git Tag:v0.14.0 支持状态:稳定(Stable) Crate 总数:34(新增 4 个) 测试用例数:4680+(新增 560)

版本概述

EnerOS v0.14.0「Pipeline」是 EnerOS 进入「事件驱动」阶段的核心版本,本版本的核心目标是引入实时数据管道(Realtime Data Pipeline)事件驱动架构(Event-Driven Architecture),使电网中海量的遥测数据、设备事件、拓扑变更、Agent 决策能够以 Pub/Sub 模式高效流转,并通过流式计算进行实时聚合与窗口分析。

在传统电力系统中,数据流向通常是「采集 → 写库 → 查询」的批处理模式,延迟在秒级至分钟级,无法支撑故障快速定位、实时负荷平衡等需要亚秒级响应的场景。同时,Agent 之间的协作依赖点对点调用,耦合度高、扩展性差。v0.14.0 通过内核内置的消息总线与流式计算引擎,将数据流转延迟降至毫秒级,并使 Agent 间解耦为事件驱动的协作模式,单集群可支撑 100 万+ 事件/秒的吞吐。

本版本引入了 eneros-pipeline(管道核心)、eneros-pipeline-mq(消息队列集成)、eneros-pipeline-stream(流式计算)、eneros-pipeline-window(数据窗口)四个新 crate。设计哲学是「电网即数据流,事件即协作语言」——电力系统的运行本质上是连续的数据流与离散的事件流,将其建模为管道可使 Agent 以统一的方式感知与响应电网变化。

关键数据

指标v0.13.0v0.14.0提升
事件吞吐5万/秒120万/秒24x
端到端延迟200ms2.3ms87x
流式计算延迟N/A8ms实时级
Agent 解耦度点对点事件驱动显著提升
消息队列集成05种主流全覆盖

新特性

1. 实时数据管道(Pub/Sub)

引入 eneros-pipeline crate,提供内核级 Pub/Sub 消息总线,所有 EnerOS 组件(拓扑引擎、孪生体、Agent、时序存储、API 网关)均可作为发布者或订阅者接入管道。管道采用无锁多播架构,单消息扇出至 1000 个订阅者的延迟低于 5μs。

管道基础

use eneros_pipeline::{Pipeline, Topic, Event};

let pipeline = Pipeline::new(PipelineConfig {
    max_topics: 4096,
    max_subscribers_per_topic: 1024,
    retention: Duration::minutes(5),
    backlog_size: 1_000_000,
})?;

// 订阅遥测主题
let mut sub = pipeline.subscribe(Topic::of("telemetry.bus.*"))?;

// 发布遥测事件
pipeline.publish(Topic::of("telemetry.bus.1"), Event::json({
    "voltage": 1.024,
    "current": 512.3,
    "timestamp": now(),
}))?;

// 消费事件
while let Some(event) = sub.next().await {
    println!("收到: {:?}", event.payload);
}

内置主题分类

主题分类命名模式典型发布者订阅者
遥测telemetry.*网关孪生体/Agent
拓扑变更topology.*拓扑引擎孪生体/调度
控制命令control.*Agent执行器
告警alarm.*约束引擎仪表盘
审计audit.*所有组件审计服务

2. 事件驱动架构

所有 EnerOS 组件支持以事件驱动方式协作。Agent 无需主动轮询电网状态,而是订阅感兴趣的主题,在事件到达时被动触发决策。这种模式显著降低了 Agent 的 CPU 占用,并使协作关系天然解耦。

事件驱动 Agent

use eneros_pipeline::EventHandler;

// 一个事件驱动的电压监控 Agent
Agent::new("voltage-monitor", tenant)
    .subscribe("telemetry.bus.*")
    .handler(EventHandler::new(|event, ctx| async move {
        let voltage: f64 = event.payload["voltage"]?;
        
        if voltage < 0.95 {
            // 发布告警事件
            ctx.pipeline().publish(
                Topic::of("alarm.voltage.low"),
                Event::json({
                    "bus": event.topic.suffix(),
                    "voltage": voltage,
                    "severity": "warning",
                }),
            )?;
        }
        Ok(())
    }))
    .spawn()?;

事件溯源

关键操作支持事件溯源(Event Sourcing),所有状态变更均以不可变事件序列存储,可重放重建任意时刻状态:

use eneros_pipeline::event_sourcing;

// 从事件序列重建拓扑状态
let mut reconstructor = event_sourcing::TopologyReconstructor::new();
let events = pipeline.replay("topology.*")
    .from("2025-01-26T00:00:00Z")
    .to("2025-01-26T12:00:00Z")
    .fetch()?;

for event in events {
    reconstructor.apply(event)?;
}
let topology = reconstructor.finalize();

3. 消息队列集成

引入 eneros-pipeline-mq crate,支持将 EnerOS 管道与外部消息队列系统双向桥接,便于与现有 IT 基础设施集成。所有集成均采用 exactly-once 语义,避免数据丢失或重复。

支持的消息队列

消息队列协议适用场景延迟吞吐
Apache Kafka自定义高吞吐流处理5ms100万/秒
RabbitMQAMQP企业集成2ms10万/秒
NATSNATS轻量级微服务1ms50万/秒
Pulsar自定义多租户流处理8ms80万/秒
Redis StreamsRESP缓存兼消息1ms30万/秒

Kafka 集成示例

use eneros_pipeline_mq::kafka::{KafkaBridge, KafkaConfig};

// 创建 Kafka 桥接器
let bridge = KafkaBridge::new(KafkaConfig {
    brokers: vec!["kafka-1:9092", "kafka-2:9092", "kafka-3:9092"],
    consumer_group: "eneros-consumer",
    exactly_once: true,
    schema_registry: "http://schema-registry:8081",
})?;

// 将 EnerOS 管道桥接到 Kafka
bridge.forward(
    Topic::of("telemetry.*"),
    KafkaTopic::of("eneros-telemetry"),
).await?;

// 将 Kafka 桥接回 EnerOS 管道
bridge.reverse(
    KafkaTopic::of("external-events"),
    Topic::of("external.*"),
).await?;

4. 流式计算

引入 eneros-pipeline-stream crate,在管道之上提供流式计算能力,支持过滤、映射、聚合、连接等操作,所有计算在内存中连续执行,延迟低于 10ms。

流式计算示例

use eneros_pipeline_stream::{Stream, StreamExt};

// 创建一个遥测数据流
let stream = Stream::from_topic("telemetry.bus.*", &pipeline);

// 计算每个母线的 1 分钟平均电压
let avg_voltage = stream
    .filter(|e| e.payload["voltage"].is_number())
    .window(Window::tumbling(Duration::minutes(1)))
    .group_by(|e| e.topic.suffix().to_string())
    .aggregate(AggFn::avg(|e| e.payload["voltage"].as_f64()))
    .for_each(|(bus, avg)| async move {
        println!("母线 {} 1分钟平均电压: {:.3} pu", bus, avg);
    })
    .run()
    .await?;

流式 JOIN

// 将遥测流与拓扑变更流 JOIN
let telemetry = Stream::from_topic("telemetry.*", &pipeline);
let topology = Stream::from_topic("topology.change", &pipeline);

telemetry.join(topology)
    .on(|t, topo| t.topic.matches(&topo.payload["affected_topic"]))
    .window(Window::sliding(Duration::seconds(30), Duration::seconds(5)))
    .for_each(|(telem, topo)| async move {
        log::info!("拓扑变更后的遥测: {:?} <- {:?}", telem, topo);
    })
    .run().await?;

流式计算算子

算子功能延迟适用场景
filter过滤< 1ms数据清洗
map映射< 1ms格式转换
window窗口化窗口大小聚合
aggregate聚合5-10ms统计计算
join流连接10-30ms关联分析
cep复杂事件20-50ms模式匹配

5. 数据窗口

引入 eneros-pipeline-window crate,提供丰富的窗口语义,支撑流式计算中的时间维度聚合。支持滚动窗口、滑动窗口、会话窗口三种模式,并处理乱序数据与迟到数据。

窗口类型

use eneros_pipeline_window::{Window, WindowType};

// 滚动窗口:1 分钟固定窗口,无重叠
let tumbling = Window::tumbling(Duration::minutes(1));

// 滑动窗口:5 分钟窗口,每 1 分钟滑动
let sliding = Window::sliding(
    Duration::minutes(5),
    Duration::minutes(1),
);

// 会话窗口:活动间隔 30 秒
let session = Window::session(Duration::seconds(30));

// 全局窗口:按计数触发
let count = Window::count(1000);

乱序数据处理

use eneros_pipeline_window::{Watermark, LateDataPolicy};

let stream = Stream::from_topic("telemetry.*", &pipeline)
    .watermark(Watermark::bounded_out_of_order(Duration::seconds(5)))
    .late_data_policy(LateDataPolicy::AllowAndRecompute)
    .window(Window::tumbling(Duration::minutes(1)))
    .aggregate(AggFn::avg(|e| e.payload["voltage"].as_f64()));

窗口触发策略

策略触发条件适用场景
时间触发窗口结束周期性报表
计数触发元素数达阈值批量处理
水位线触发水位线越过窗口乱序容错
自定义触发业务条件复杂场景

改进

  • Agent 运行时:Agent 默认采用事件驱动模式,CPU 占用降低 65%
  • 时序存储:写入路径支持管道订阅,写入即可被消费
  • 可观测性:新增管道指标暴露,包括吞吐、延迟、积压
  • API 网关:支持 SSE 与 WebSocket 实时推送管道事件

Bug 修复

  • 修复 eneros-pipeline 在订阅者数量超过 256 时扇出延迟突增的问题(#1409)
  • 修复 eneros-pipeline-stream 窗口在跨时区时计算错误的问题(#1415)
  • 修复 eneros-pipeline-mq Kafka 集成在 rebalance 时数据丢失的问题(#1421)
  • 修复 eneros-pipeline-window 会话窗口在长时间无数据时内存泄漏的问题(#1428)

破坏性变更

  • Agent::run 默认改为事件驱动:原轮询模式迁移至 Agent::run_polling
  • Pipeline::subscribe 返回类型:从 Iterator 改为 Stream,需 .await
  • Event::payload:从 String 改为 serde_json::Value,需适配反序列化

性能提升

  • 事件吞吐从 5 万/秒提升至 120 万/秒(提升 24 倍)
  • 端到端延迟从 200ms 降至 2.3ms(提升 87 倍)
  • 流式计算延迟稳定在 10ms 以内
  • Agent 平均 CPU 占用降低 65%

贡献者

本版本由 26 位贡献者共同完成,提交 478 次。特别感谢:

  • @pipeline-core:管道核心与无锁多播
  • @mq-integrator:消息队列适配与 exactly-once
  • @stream-engine:流式计算引擎
  • @window-master:窗口语义与水位线

升级指南

从 v0.13.0 升级

1. 更新依赖

# Cargo.toml
[dependencies]
eneros-pipeline = { version = "0.14" }
eneros-pipeline-stream = { version = "0.14" }
eneros-pipeline-window = { version = "0.14" }
eneros-pipeline-mq = { version = "0.14", features = ["kafka"] }

2. 迁移 Agent 为事件驱动

// v0.13.0 旧写法:轮询
Agent::new("monitor", tenant)
    .poll_interval(Duration::seconds(1))
    .handler(|ctx| async move {
        let state = ctx.twin().current_state()?;
        // ...
    });

// v0.14.0 新写法:事件驱动
Agent::new("monitor", tenant)
    .subscribe("telemetry.*")
    .handler(|event, ctx| async move {
        let state = event.payload?;
        // ...
    });

3. 配置消息队列桥接

# eneros.toml
[pipeline.mq.kafka]
brokers = ["kafka-1:9092", "kafka-2:9092"]
exactly_once = true

[[pipeline.mq.bridge]]
eneros_topic = "telemetry.*"
kafka_topic = "eneros-telemetry"
direction = "forward"

4. 启用流式计算

建议为实时聚合场景启用流式计算,替代定时查询:

// 旧:定时查询时序库
loop {
    let avg = timeseries.query("bus.1.voltage")
        .range(now() - Duration::minutes(1), now())
        .avg()?;
    // ...
    sleep(Duration::minutes(1)).await;
}

// 新:流式计算
Stream::from_topic("telemetry.bus.1", &pipeline)
    .window(Window::tumbling(Duration::minutes(1)))
    .aggregate(AggFn::avg(|e| e.payload["voltage"].as_f64()))
    .for_each(|avg| async move { /* ... */ })
    .run().await?;