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.0 | v0.14.0 | 提升 |
|---|---|---|---|
| 事件吞吐 | 5万/秒 | 120万/秒 | 24x |
| 端到端延迟 | 200ms | 2.3ms | 87x |
| 流式计算延迟 | N/A | 8ms | 实时级 |
| Agent 解耦度 | 点对点 | 事件驱动 | 显著提升 |
| 消息队列集成 | 0 | 5种 | 主流全覆盖 |
新特性
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 | 自定义 | 高吞吐流处理 | 5ms | 100万/秒 |
| RabbitMQ | AMQP | 企业集成 | 2ms | 10万/秒 |
| NATS | NATS | 轻量级微服务 | 1ms | 50万/秒 |
| Pulsar | 自定义 | 多租户流处理 | 8ms | 80万/秒 |
| Redis Streams | RESP | 缓存兼消息 | 1ms | 30万/秒 |
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-mqKafka 集成在 rebalance 时数据丢失的问题(#1421) - 修复
eneros-pipeline-window会话窗口在长时间无数据时内存泄漏的问题(#1428)
破坏性变更
Agent::run默认改为事件驱动:原轮询模式迁移至Agent::run_pollingPipeline::subscribe返回类型:从Iterator改为Stream,需.awaitEvent::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?;