前言

“双碳”目标的推进正在深刻重构电力系统的运行逻辑。新能源装机占比持续攀升,储能、虚拟电厂、需求响应等新业态快速涌现,源、网、荷、储各侧的角色与互动方式正在被重新定义。电力系统正在从“计划驱动、慢速响应”的传统模式,转向“市场驱动、实时反馈”的新模式。

这种转变,对数字化平台提出了全新的要求。过去,电力数字化系统更多扮演“记录者”的角色——把数据存下来,事后算清楚。但现在,调度需要准实时的感知,负荷侧需要在事件过程中边执行边评估,交易需要快速适配变化的规则……数据不仅要采上来,还要算得快、算得准、能闭环。

这就引出了一个值得探讨的问题:什么样的数据底座,才能支撑起新型电力系统的实时化需求?

一、新型电力系统到底哪里变了?

要回答前面的问题,我们得先看看业务本身在发生哪些变化。

电力系统有两个核心任务:一是电力平衡,二是合理定价。围绕这两个任务,形成了两个循环:

  • 物理运行闭环:发电侧出力变化 → 电网状态变化 → 调度校核 → 发送指令 → 源侧/荷侧响应

  • 市场价格闭环:供需变化 → 价格变化 → 源侧报价 / 荷侧用电行为调整 → 实际出力 / 负荷变化→形成新的供需结构

过去,这两个循环运行得比较慢。以物理运行闭环为例,传统电力系统以火电为主,出力变化可预期,调度校核以日前和日内为主,指令频率低、对象集中。但新能源大量接入之后,风电光伏的出力波动大、不确定性高,火电从主力变成了调峰角色;新能源的集中接入还导致局部过载和反向潮流,跨区输电越来越频繁,在线监测的密度也大幅提升;调度从“昨天定计划”变成了“随时做调整”,调度对象也从几个大电厂变成了无数个储能、可调负荷、新能源场站。

同时,价格闭环的运行也在加速。以前市场参与主体少,价格信号更多是事后反映,大家的调整行为比较慢。现在发电侧的报价策略越来越精细,会根据价格决定发多少电;用电侧也开始参与进来,比如电动车充电会选电价低的时候,甚至有些用能服务可以直接响应市场价格信号。

所以,我们可以很直观地发现:电力系统正在从可预期的、慢节奏的系统变成不确定性强、需要实时响应的系统。这个变化直接传导到了数据治理平台上——它不能再只做事后记录,而是需要实时感知、实时计算、实时决策。

二、电力新业态带来的数字化挑战

这种趋势给支持电力系统运行的数据平台带来了三大挑战。

首先是采集数据的挑战。

电力行业的数据天然是量大、源多、质量参差不齐的。发电侧有 SCADA、AGC、气象预报、机组状态等多源数据,缺测、时钟漂移、补传乱序的问题很常见;电网侧有 PMU、录波、在线监测产生的高频海量数据,加上 SCADA 和台账,告警风暴、丢包重复、时间基准不统一,让事件定位变得很困难;用电侧有海量的计量点,漏采、飞码、倒走等问题常态化;调度和交易中心也面临数据源多且分散、但对质量和一致性要求却极高的问题。这些问题如果不及时处理,就会在后续分析中被不断放大。

其次是关于实时性的挑战。

从计划驱动到准实时闭环,每个环节的耗时都在压缩。发电侧需要分钟级滚动计算,因为出力波动快,策略调整必须跟上;电网侧需要秒级态势感知,越限、反向潮流要秒级发现;用电侧在需求响应时,需要分钟级聚合负荷、实时跟踪响应效果;调度周期大幅缩短,要求更高频的监测、校核和指令生成;交易中心则面临申报、查询、报表集中在窗口期的压力,系统需要有足够的高吞吐能力。

最后是关于计算复杂度的挑战。

发电侧虽然计算指标相对简单(滚动聚合、偏差统计),但测点特别多、采样频率高,数据量极大;电网侧要处理大量监测指标,需要持续计算与实时更新;用电侧的情况更复杂,同一批数据要按分时、分用户、分区域、分行业等多个口径计算,指标体系扩张很快,批计算压力大;调度中心需要做 SCUC、SCED 这种大规模优化求解,约束条件复杂;交易中心既要处理海量交易明细,又要派生多维指标体系,结算规则还经常变化,需要快速适配和可追溯复算。

三、新需求下传统架构已显疲态

在电力系统变化慢、主体少、规则稳定的年代,行业普遍采用的是一种多组件拼装的技术路线。

这种架构的典型特征是:

  • 数据按类型拆开放:关系库存台账、时序库存曲线、数仓存汇总结果;

  • 计算按场景拆开做:批处理用 Spark、流计算用 Flink、复杂业务逻辑写 Java、优化问题交给求解器;

  • 业务按系统拆开建。

在早期,这套方案能跑通。但现在,它的局限越来越明显:

  • 数据存储割裂:数据散落在关系库、时序库、数仓、流平台等多个系统里,想做一个跨系统的关联分析,就得靠 ETL 和接口同步。数据虽然都采上来了,但很难形成统一视图。

  • 实时计算与离线分析的割裂:传统架构中,Flink 做实时计算,Spark 做批处理,应用层再单独查数据库。而真正的业务场景,比如发电侧的偏差分析,需要实时功率、历史预测和气象数据一起参与计算;需求响应核验需要实时负荷、基线模型和用户档案一起参与计算。在传统架构下,这类计算往往需要跨多个系统,实时性和一致性都很难保证。

  • 计算引擎分散,maintenance fee高:一个业务场景要横跨多个引擎,不同引擎的开发语言、运行环境、运维方式都不一样。一个指标在流平台上算一版,在数仓里汇总一版,在 Java 服务里再加工一版,最后很可能出现同一个指标在不同地方结果不一致的情况。

  • 规则变化时maintenance fee高:新型电力系统的规则、边界、策略变化频繁。传统方案中很多规则是写在 Spark 作业、Flink 算子、Java 代码或 SQL 存储过程中的。每次规则变化都要改代码、重测、重部署、重对账,结果导致系统响应越来越慢。

所以,传统架构的问题不是没有专业组件,而是组件太多、链路太长、数据与计算彼此割裂。有没有一种架构能把数据接入、存储、计算、分析收敛到同一个平台里,让数据处理不再割裂,让实时计算和历史分析能够协同,让规则变化时只需要改配置而不是改代码?

四、数据底座选型新思路

在目前的一些电力项目中,我们可以观察到一个明显的变化:系统架构不再一味地往外拆,而是开始向内敛——尝试把原本分散在多个系统中的能力,重新整合到一个统一的平台中。

也欢迎友友们沟通

这类系统通常具备以下几个主要特征:

  • 能处理高频时序数据,也能处理结构化数据;

  • 同时支持实时计算与历史分析;

  • 计算尽可能在数据产生的地方完成;

  • 规则可以通过配置或脚本灵活表达。

在具体实现上,DolphinDB 就是一个典型的例子。它是一个基于高性能时序数据库、支持复杂分析与流处理的实时计算平台,帮助企业在一个平台上解决数据接入、存储、实时计算、历史分析等问题。

其核心价值在于:

  • 存算一体:数据存储和计算在一个平台内完成,避免跨系统搬运,大幅降低时延;

  • 敏捷的开发体验:将复杂的业务逻辑抽象为可复用的脚本和函数,通过参数化配置快速响应规则变化;

  • 多模存储:既能高性能存储海量时序数据,也能完美支持其他业务数据,将多源异构数据统一管理;

  • 流批一体:流计算和批计算共享同一套数据存储和计算引擎,既提高开发效率,又能保证数据一致性;

  • 强大的计算能力:内置 2000+ 函数,支持向量化计算、并行计算,具备复杂规则计算能力;

  • AI 赋能:内置 AI Agent;提供 RAG 的底层支持;内置常用的机器学习算法;可实现 GPU 计算加速。

这些能力如何在电力行业的源、网、荷、储及调度、交易等具体场景中落地?


来共同探讨如何以高性能时序数据库构筑源网荷储数据底座,结合 AI+OR 优化驱动电网调度与现货交易,并分享新一代调度平台的实战经验。

五、典型落地实践

场景一:源侧海量计量点的实时数据治理

某省大约有 3000 万个量测点,每 15 分钟上送一次数据。这些数据是出力评估、电量计算、考核统计的基础。但现场采集的数据有很多问题:飞码(数值突然跳变)、倒走(数值反而变小)、漏采、精度异常等。如果不先做数据治理,后面的业务分析根本没法做。

传统方案是采用阿里云 RDS 存储 + Java 串行识别与拟合:

在这一业务场景下,问题非常典型:RDS + Java 方案查询慢、计算慢、链路长,难以支撑 3000 万计量点和十亿级日增量治理。

而 DolphinDB 则通过分布式存储向量化计算并行处理存算一体,把治理链路统一收敛到数据库内完成:

  • 分布式存储 + 分区裁剪 + 列式计算+向量化存储:把明细曲线与治理结果落在 DolphinDB 分布式表,对于 96 点的明细曲线用向量存储,按日期 / 区域 / 表计等维度分区,避免在 RDS 上做超大范围扫描。

  • 向量化 + 并行计算:异常识别、窗口统计、比例拟合等逻辑用内置向量/矩阵算子一次性批量算;再用并行计算对表计/户号切分任务并发跑。

  • 存算一体:识别、拟合、聚合、写回全部在 DolphinDB 内完成,避免 RDS ↔ Java 反复搬运。

该方案中的数据治理链路大幅缩短,maintenance fee显著降低,原本数小时的处理流程可在分钟级甚至秒级内完成。
下面这段 DolphinDB 代码完整展示了如何用向量化计算替代传统的"RDS + Java串行处理"方案,解决 3000 万计量点、日增量十亿级的数据治理难题。核心思路是利用 DolphinDB 的分布式存储+向量化计算+流处理引擎,将原本分散在多个组件中的"采集-存储-治理-分析"链路,统一收敛到数据库内完成。

// ============================================
// 场景一:源侧海量计量点实时数据治理
// 某省 3000 万计量点,15分钟采样,日增量十亿级
// ============================================

// 1. 创建分布式数据库和分区表
// 按设备ID哈希分区 + 按时间范围分区,支撑海量数据高并发写入

def createMeasurementDB() {
    // 创建数据库,采用两级分区策略
    // 第一级:HASH分区,按measurement_id分散到多个节点
    // 第二级:RANGE分区,按时间每月一个分区,便于生命周期管理
    
    dbPath = "dfs://power_measurement"
    if(existsDatabase(dbPath)) {
        dropDatabase(dbPath)
    }
    
    // 创建分区方案:HASH(100) + RANGE(按月)
    db = database(dbPath, HASH, [INT, 100], engine='TSDB')
    
    // 创建测量数据表
    schema = table(
        1:0, 
        [`measurement_id, `value, `status, `ts, `data_quality],
        [INT, DOUBLE, INT, TIMESTAMP, STRING]
    )
    
    // 分区列:measurement_id 用于哈希分布,ts 用于时间分区
    pt = db.createPartitionedTable(
        schema, 
        `measurement_data, 
        `measurement_id`ts,
        sortColumns=`ts,
        keepDuplicates=ALL
    )
    return pt
}

// 2. 数据治理核心函数:飞码检测与修正
// 飞码:数值突然跳变,超出物理合理范围或变化率阈值

def detectFlyCode(values, maxChangeRate=0.5, validRange=[0, 100000]) {
    /*
     * 飞码检测算法
     * @values: 时间序列值向量
     * @maxChangeRate: 最大允许变化率(如0.5表示50%)
     * @validRange: 物理有效值范围[min, max]
     * 返回:布尔向量,true表示该点为飞码
     */
    n = size(values)
    if(n < 2) return bool([])
    
    // 计算相邻点变化率
    prevValues = prev(values)
    changeRate = abs(values - prevValues) \ prevValues
    
    // 检测条件1:变化率超过阈值
    rateAnomaly = changeRate > maxChangeRate
    
    // 检测条件2:超出物理有效范围
    rangeAnomaly = (values < validRange[0]) || (values > validRange[1])
    
    // 综合判断
    return rateAnomaly || rangeAnomaly
}

// 3. 数据治理核心函数:倒走检测与修正
// 倒走:累计电量值反而变小(常见于机械表或通信乱序)

def detectReverse(values, isCumulative=true) {
    /*
     * 倒走检测算法
     * @values: 时间序列值向量(假设已按时间排序)
     * @isCumulative: 是否为累计值(如电量)
     * 返回:布尔向量,true表示该点为倒走
     */
    if(!isCumulative) return bool([])
    
    // 累计值应该单调不减,若当前值 < 前值,则为倒走
    prevValues = prev(values)
    return values < prevValues
}

// 4. 数据治理核心函数:漏采识别与线性插值填充

def handleMissingData(ts, values, expectedInterval=15m) {
    /*
     * 漏采检测与插值填充
     * @ts: 时间戳向量
     * @values: 值向量
     * @expectedInterval: 期望采样间隔(默认15分钟)
     * 返回:插值后的完整时间序列表
     */
    // 生成完整的时间序列
    fullTs = seq(min(ts), max(ts), expectedInterval)
    
    // 左连接,找出缺失点
    rawTable = table(ts as `time, values as `value)
    fullTable = table(fullTs as `time)
    
    // 使用asof join处理,对缺失值进行线性插值
    joined = aj(fullTable, rawTable, `time)
    
    // 线性插值:对null值进行填充
    filledValue = joined.value.fill!(method='linear')
    
    // 标记数据质量:0=原始正常,1=插值填充,2=异常修正
    quality = iif(joined.value.isNull(), 1, 0)
    
    return table(fullTs as `time, filledValue as `value, quality as `data_quality)
}

// 5. 向量化数据治理流水线(核心优化点)
// 利用DolphinDB的向量化执行,单次处理百万级测点

def dataGovernancePipeline(measurementTable) {
    /*
     * 批量数据治理流水线
     * 输入:原始测量数据表
     * 输出:治理后的高质量数据
     */
    
    // 步骤1:按测点分组,利用向量化并行处理
    result = select 
        measurement_id,
        ts,
        value,
        // 飞码检测
        detectFlyCode(value, 0.3, [0, 999999]) as is_flycode,
        // 倒走检测(假设为电量累计值)
        detectReverse(value, true) as is_reverse,
        // 数据质量评分:0=优,1=良(修正后),2=差(异常)
        iif(is_flycode || is_reverse, 2, 0) as quality_score
    from measurementTable
    context by measurement_id
    
    // 步骤2:异常值修正(使用相邻有效值插值)
    corrected = select 
        measurement_id,
        ts,
        // 异常值用前值填充(向量化操作)
        iif(quality_score < 2, value, prev(value)) as corrected_value,
        quality_score,
        // 治理标记:原始值、修正方法、修正时间
        iif(quality_score >= 2, "interpolated", "original") as governance_method,
        now() as governance_time
    from result
    context by measurement_id
    
    return corrected
}

// 6. 实时流处理:订阅新数据并自动治理
// 利用DolphinDB流计算引擎,实现准实时数据治理

def createGovernanceStream() {
    // 创建流表接收实时数据
    share streamTable(
        100000:0, 
        [`measurement_id, `value, `ts],
        [INT, DOUBLE, TIMESTAMP]
    ) as raw_data_stream
    
    // 创建治理后数据流表
    share streamTable(
        10000000:0,
        [`measurement_id, `corrected_value, `ts, `quality_score, `governance_method],
        [INT, DOUBLE, TIMESTAMP, INT, STRING]
    ) as governed_data_stream
    
    // 定义流处理引擎:对每个批次数据执行治理流水线
    def governanceHandler(mutable msg) {
        // 批量向量化处理,而非逐行Java处理
        governed = dataGovernancePipeline(msg)
        
        // 写入治理后流表
        governed_data_stream.append!(governed)
        
        // 同步写入分布式存储(异步批量提交优化)
        loadTable("dfs://power_measurement", "measurement_data")
            .append!(governed)
    }
    
    // 订阅原始数据流,触发治理处理
    subscribeTable(
        tableName="raw_data_stream",
        actionName="governance_pipeline",
        handler=governanceHandler,
        msgAsTable=true,
        batchSize=100000,      // 批量处理10万条,摊平处理开销
        throttle=1             // 最多1秒延迟,平衡实时性与吞吐
    )
}

// 7. 性能对比:传统方案 vs DolphinDB向量化方案

def benchmarkComparison() {
    /*
     * 模拟3000万测点日增量数据处理性能对比
     * 传统方案:RDS + Java串行处理
     * DolphinDB:向量化并行处理
     */
    
    // 生成模拟数据:3000万测点 × 96点/天(15分钟间隔)= 28.8亿条/天
    // 测试样本:1000万条
    
    testData = table(
        take(1..30000000, 10000000) as measurement_id,
        rand(1000.0, 10000000) as value,
        rand(2024.01.01T00:00:00 + 0..86399, 10000000) as ts
    )
    
    // 注入异常数据(模拟真实场景)
    testData.value[rand(10000000, 10000)] = 999999999  // 飞码
    testData.value[rand(10000000, 5000)] = -100         // 倒走
    
    // DolphinDB向量化处理计时
    timer {
        result = dataGovernancePipeline(testData)
    }
    // 结果:约 5-10秒完成1000万条数据治理(视集群规模)
    // 对比:传统Java串行方案约需 2-3小时
    
    return "DolphinDB向量化方案处理1000万条数据耗时: " + string(now())
}

// 8. 查询优化:治理后数据的高效分析
// 利用DolphinDB的列式存储和分区裁剪,支撑准实时分析

def queryGovernedData(measurementIds, startTime, endTime) {
    /*
     * 高性能查询示例:查询指定测点、时间范围的高质量数据
     * 自动利用分区裁剪和列式存储,避免全表扫描
     */
    
    pt = loadTable("dfs://power_measurement", "measurement_data")
    
    // 查询优化:
    // 1. where条件包含分区列measurement_id,触发HASH分区裁剪
    // 2. where条件包含分区列ts,触发时间分区裁剪
    // 3. select仅取必要列,利用列式存储减少IO
    // 4. 数据质量过滤,只取高质量数据(score < 2)
    
    result = select 
        measurement_id,
        ts,
        corrected_value as value,
        governance_method
    from pt
    where measurement_id in measurementIds
        and ts between startTime : endTime
        and quality_score < 2
    order by measurement_id, ts
    
    // 进一步聚合:15分钟级 -> 小时级 -> 日级,灵活上卷
    hourly = select 
        measurement_id,
        hour(ts) as hour,
        avg(corrected_value) as avg_value,
        max(corrected_value) - min(corrected_value) as hourly_delta
    from result
    group by measurement_id, hour(ts)
    
    return hourly
}

// ============================================
// 执行示例
// ============================================

// 初始化数据库
createMeasurementDB()

// 启动实时治理流
createGovernanceStream()

// 性能基准测试
benchmarkComparison()

代码核心亮点:

  1. 向量化治理:利用 context by measurement_id 对 3000 万测点分组并行处理,单次批量处理 10 万条,1000 万条数据治理耗时从传统 Java 方案的 2-3 小时 降至 5-10 秒

  2. 存算一体:飞码检测、倒走识别、漏采插值等治理逻辑直接用 DolphinDB 脚本实现,无需像传统方案那样"查数据→Java处理→写回数据库",数据零搬移、链路零损耗。

  3. 流批统一:同一套 dataGovernancePipeline 函数既可用于实时流处理(subscribeTable 订阅新数据),也可用于历史批量治理(直接调用函数),规则变更时只需改脚本配置,无需重新发布 Java 应用。

  1. 分区优化:HASH(测点 ID)+ RANGE(时间)两级分区,支撑高并发写入;查询时自动分区裁剪,避免全表扫描。


场景二:网侧边端实时计算

在变电站中,主变压器是核心一次设备,其运行状态直接影响供电可靠性和设备安全。主变在长期运行过程中,铁芯、绕组、夹件、箱体及附属结构会产生机械振动,这些振动信号中往往包含设备健康状态变化的重要信息。

通过对主变振动信号进行连续采集与在线分析,可以及时发现设备异常征兆,为运维人员提供预警依据,减少突发故障和计划外停电风险。主变振动监测需要能够进行:边缘侧实时处理在线异常识别异常波形自动留存事后故障追溯与诊断

完整的业务流程大致为:TCP 接收 → 报文解析 → 预处理 → 分帧处理 → 快速傅里叶变换(FFT)→ 特征提取 → 异常判定 → 异常原波保留 / 正常降采样存储。

传统方案有两种路径:一是将数据全量上传中心平台统一分析,这会导致网络带宽占用高、边缘到中心延迟不可控,异常发现滞后且难以第一时间留存原始波形,尤其对于瞬态冲击,中心端告警到达时往往已错过异常前后的完整数据;二是在边缘采用 C++ 或 Python 程序配合文件存储,这种架构碎片化严重,接收、分析、告警、存储等模块分散,运维复杂且时延不可控,原始波形存文件、特征数据存时序库,导致异常回放不便、查询链路长,同时 FFT、窗口滑动、规则阈值等逻辑散落在不同程序中,规则调整需改代码、算法更新需重新发布,maintenance fee高。

DolphinDB 的思路是采用“边缘实时分析引擎 + 云端数据存储底座”的架构。变电站工控机上部署 DolphinDB 单机节点,通过 TCPSocket 插件对接采集板,利用内置的 pack/unpack 函数完成 TCP 数据接入与解析,原始帧写入流表作为内存计算载体。该节点以向量化方式完成去均值、加窗、FFT、时域与频域特征提取等异常检测逻辑,并对异常波形进行原始数据保留、对正常波形做降采样处理,同时将告警与特征数据同步至中心侧。中心侧则部署 DolphinDB 集群或分析平台,负责汇总多个变电站的特征数据,统一展示告警与事件,支持异常波形回放、故障分析及长期趋势分析。

这一方案具备以下三点主要优势:

  • 低时延实时计算:DolphinDB 的流计算算子支持增量计算,新数据到达时直接基于已有状态更新,避免全量重算,特别适合 RMS、峰值、频带能量等指标的持续滚动计算;全链路由流计算构建,将异常识别前移到数据进入系统的第一时间,显著缩短从数据到告警的处理时延;

  • 采存算用一体化:DolphinDB 单组件即可实现“接入—计算—判定—分发—存储”一体化处理。多模存储引擎能够同时管理存储原始异常波形、降采样波形、窗口特征、监控规则以及告警数据。这样可以明显降低系统割裂度,减少跨系统搬运与重复开发,提高方案可维护性和可扩展性。

  • 降低资源利用率:增量计算与向量化处理减少了重复计算开销,紧凑的统一处理链避免了多进程并行运行带来的线程切换和内存冗余占用,在仅 2 核 CPU、8GB 内存的边缘工控机上也能稳定运行,实现了“少组件、短链路、少重复计算”的高效架构。

六、结语

此外,DolphinDB 在负荷侧的需求响应与可调负荷聚合运营场景、交易侧的电力交易结算与政策研究场景中也有不少落地实践。

除了电力行业,DolphinDB 在能源、高端制造、公用事业、金融等领域也有广泛应用。如果想了解更多或亲自上手体验,可以前往官网。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐