1. 背景与设计目标

本框架面向边缘计算工业现场应用场景,在资源受限(CPU、内存、存储)的环境下,提供高并发、低延迟的计算能力。设计核心原则:

  • 依赖极少:仅依赖 aiosqlitenetworkxnumpypydantic 四个核心库,均轻量且稳定。
  • 计算效率高:通过进程池 + 共享内存,充分利用多核 CPU,避免 GIL 限制。
  • 简洁易维护:业务逻辑与调度框架分离,通过配置驱动,降低维护成本。
  • 适配 Python 3.11+:利用最新异步特性,代码风格现代,类型提示完整。

框架并非通用工作流引擎,而是专注于数据密集型、计算密集型的微电网仿真场景,同时可灵活扩展至其他领域。

2. 框架概述

2.1 核心能力

  • 流程编排:通过定义有向无环图(DAG)描述计算步骤及依赖关系。
  • 混合并行:CPU 密集型任务自动分发到进程池,IO 密集型任务在异步事件循环中执行。
  • 三级缓存:内存(LRU) + 共享内存(跨进程) + SQLite(持久化),兼顾性能与可靠。
  • 状态恢复:节点执行状态持久化,服务重启后可从中断点继续计算。
  • 项目隔离:每个项目拥有独立的 Actor 和缓存空间,支持多项目并行。

2.2 框架边界与限制

  • 不支持条件分支节点:DAG 是静态结构,无法在运行时动态选择路径。业务上可通过节点函数内部判断实现分支效果。
  • 不支持循环:节点不会重复执行,如需迭代需在业务函数内部实现。
  • 不支持动态修改 DAG:工作流定义在启动时加载,不可热更新。

3. 快速开始(无 FastAPI 示例)

以下是一个完整的独立测试脚本,演示了如何编写业务函数、定义工作流并启动计算。

# test_simple.py
import asyncio
import random
import logging
from core.system import ActorSystem
from core.function_registry import registry
from core.config import WORKFLOWS

logging.basicConfig(level=logging.INFO)

# ---------- 1. 编写业务函数 ----------
def generate_random(project_id: str, config: dict) -> list:
    """返回24个随机数(0~100)"""
    return [round(random.uniform(0, 100), 2) for _ in range(24)]

def sum_two_lists(project_id: str, config: dict, list_a: list, list_b: list) -> list:
    """将两个列表按位相加"""
    return [round(a + b, 2) for a, b in zip(list_a, list_b)]

async def store_result(project_id: str, data: list, db):
    """将结果存入数据库(需提前建表)"""
    # 此处省略具体存储逻辑,实际使用时注入数据库连接
    logging.info(f"Stored {len(data)} items for project {project_id}")
    return {"stored": len(data)}

# ---------- 2. 注册函数到框架 ----------
# 方式1:动态注册(测试用)
registry._functions['generate'] = generate_random
registry._functions['sum'] = sum_two_lists
registry._functions['store'] = store_result

# ---------- 3. 定义工作流(注入 WORKFLOWS) ----------
WORKFLOWS['test_flow'] = {
    "nodes": [
        ("config", None, "io", []),
        ("A", "generate", "cpu", ["config"]),
        ("B", "generate", "cpu", ["config"]),
        ("C", "sum", "cpu", ["config", "A", "B"]),
        ("save", "store", "io", ["C"]),
    ]
}

async def main():
    # 初始化系统(使用临时数据库)
    system = ActorSystem()
    await system.init(
        config_db_path="test_config.db",
        energy_db_path="test_energy.db"
    )
    # 启动项目
    params = {"name": "Test"}  # 项目配置
    await system.start_project("PROJ001", params=params, workflow="test_flow")
    
    # 等待计算完成(Actor 自动休眠)
    while "PROJ001" in system.actors:
        await asyncio.sleep(1)
    
    await system.shutdown()

if __name__ == "__main__":
    asyncio.run(main())

运行方式:直接在项目根目录执行 python test_simple.py。该示例完全独立,不依赖 FastAPI。

4. 核心原理

4.1 Actor 与 DAG

  • Actor:每个项目(project_id)对应一个 ProjectActor 实例,它拥有独立的 DAG、缓存键空间和运行状态。
  • DAG 调度DAGScheduler 维护节点图,只有所有依赖节点状态为 done 的节点才会被调度(get_ready())。节点状态:pending → running → done/failed

4.2 节点类型与执行环境

类型 执行方式 适用场景 函数特征
cpu 进程池(ProcessPoolExecutor 计算密集型(如矩阵运算、优化算法) 同步函数(无 async
io 事件循环(asyncio 数据库查询、网络请求 异步函数async def
local 事件循环(asyncio 轻量级操作(如简单数据处理) 同步或异步均可(推荐同步)

4.3 数据传递与缓存

  1. 节点间传递:每个节点执行结果存入 self.dag.node_results 字典(键为节点名),下游节点按依赖顺序从该字典中获取。
  2. 三级缓存
    • L1(内存 LRU):缓存小对象或共享内存元数据。
    • L2(共享内存):存储大型 numpy 数组,支持跨进程零拷贝访问。
    • L3(SQLite):持久化所有结果,用于服务重启恢复。
  3. 引用计数:共享内存数组通过 SharedMemoryManager 管理引用计数,自动释放。

4.4 状态持久化与恢复

  • 每个节点完成后,框架自动调用 save_workflow_statenode_status 和结果引用存入 project_workflow_state 表。
  • 服务重启后,start_project 检测到该表中有记录,则调用 restore_state 重建 DAG 状态,并从缓存或数据库加载已完成节点的结果。

5. 配置文件与目录结构

5.1 推荐项目目录结构

your_project_root/
├── core/                     # 框架核心(不可修改)
│   ├── actor.py
│   ├── system.py
│   ├── cache.py
│   ├── dag.py
│   ├── worker.py
│   ├── function_registry.py
│   ├── config.py             # 框架配置(max_workers, l1_max_size等)
│   └── ...
├── config/                   # 业务配置目录(与core解耦)
│   ├── workflows.py          # 工作流定义(WORKFLOWS字典)
│   └── functions.json        # 函数注册映射
│   
├── business/                 # 业务函数实现(可根据领域命名)
│   ├── weather.py
│   ├── pv_power.py
│   └── storage.py
├── services/                 # 可选:API服务层(FastAPI等)
│   └── api/
├── db/                       # 数据库文件存放目录
├── logs/                     # 日志目录
├── main.py                   # 应用入口
└── requirements.txt

5.2 配置文件详解

5.2.1 工作流定义 (config/workflows.py)

该文件导出 WORKFLOWS 字典,键为流程名称,值为包含 nodes 列表的字典。每个节点是一个元组 (节点名, 注册函数名, 类型, 依赖列表)

# config/workflows.py
WORKFLOWS = {
    "default": {
        "nodes": [
            ("config", None, "io", []),
            ("fetch_weather", "fetch_weather", "io", ["config"]),
            ("pv_predict", "pv_predict", "cpu", ["config", "fetch_weather"]),
            ("store_pv", "store_pv", "io", ["pv_predict"]),
        ]
    },
    "feasibility": {
        "nodes": [
            ("config", None, "io", []),
            ("weather", "fetch_weather", "io", ["config"]),
            ("pv_simple", "pv_predict_simple", "cpu", ["config", "weather"]),
        ]
    }
}
  • 必须包含一个名为 config 的虚拟节点(函数名为 None),作为数据源根节点。
  • 节点名依赖列表 中的名称必须一致,且不能出现循环依赖。
  • 注册函数名 必须在 functions.json 中定义,或通过 registry._functions 动态注册(测试用)。

5.2.2 函数注册映射 (config/functions.json)

JSON 文件,映射函数名到具体的模块路径和函数名,支持别名、超时等元数据。

{
  "fetch_weather": {
    "module": "business.weather",
    "function": "fetch_typical_year_weather",
    "type": "io",
    "timeout": 60,
    "description": "获取典型年气象数据"
  },
  "pv_predict": {
    "module": "business.pv",
    "function": "calculate_pv_power",
    "type": "cpu",
    "timeout": 120
  },
  "store_pv": {
    "module": "business.storage",
    "function": "store_pv_result",
    "type": "io",
    "timeout": 30
  }
}
  • module:Python 模块路径,框架会动态导入。
  • function:模块内的函数名。
  • type:节点类型(cpu/io/local),若未指定,框架会根据函数是否为协程自动推断。
  • timeout:可选,节点执行超时秒数(仅对 CPU 节点有效)。

框架在启动时通过 registry.initialize("config/functions.json") 加载并注册这些函数。业务函数实现在 business/ 目录下,框架仅需知道模块路径即可。

5.3 配置加载机制

  • core/config.py 中的 Config 类管理框架级参数(进程数、缓存大小等),这些参数通常通过环境变量或 .env 文件配置。
  • 业务配置(工作流和函数映射)通过 core/config.py 中的动态导入加载(见前文“配置外置”方案)。
  • 项目运行时参数(如 params)由调用方(API 或脚本)提供,框架不主动读取文件。

5.4 与 core 的解耦原则

  • core 目录不包含任何业务代码或配置,仅包含框架核心类。
  • 所有业务相关的函数实现、工作流定义、函数映射均放在 config/business/ 中。
  • 框架通过依赖注入和配置路径获取业务信息,确保可独立升级和复用。

6. 典型流程模式

6.1 单路顺序

config → A → B → C

依赖关系:A 依赖 configB 依赖 AC 依赖 B。调度时依次执行。

配置示例

WORKFLOWS["sequence"] = {
    "nodes": [
        ("config", None, "io", []),
        ("A", "func_a", "cpu", ["config"]),
        ("B", "func_b", "cpu", ["A"]),
        ("C", "func_c", "cpu", ["B"]),
    ]
}

6.2 汇聚

config → A → C
config → B → C

C 依赖 AB,框架会等待两者都完成后才执行 Cargs 顺序按依赖列表顺序传入。

配置示例

WORKFLOWS["converge"] = {
    "nodes": [
        ("config", None, "io", []),
        ("A", "func_a", "cpu", ["config"]),
        ("B", "func_b", "cpu", ["config"]),
        ("C", "func_c", "cpu", ["config", "A", "B"]),
    ]
}

注意:C 的依赖列表显式包含 configAB,确保参数顺序为 (config, A_result, B_result)

6.3 并行分叉

config → A
config → B

AB 无依赖,会并行执行(受进程池大小限制)。

配置示例

WORKFLOWS["fork"] = {
    "nodes": [
        ("config", None, "io", []),
        ("A", "func_a", "cpu", ["config"]),
        ("B", "func_b", "cpu", ["config"]),
    ]
}

7. 高级主题

7.1 手动控制执行(重算)

通过 API(或直接调用 start_project)指定 recompute_from 参数,从某个节点开始重新计算子树,依赖该节点的下游节点全部重置为 pending

示例

await system.start_project("P001", params, workflow="default", recompute_from="weather")

该调用会清空 weather 及下游所有节点的状态,然后重新执行。

原理reset_subtree(node) 遍历 DAG 下游节点,清除 node_results 和共享内存引用,并将状态置为 pending,后续调度会重新执行这些节点。

7.2 自定义缓存策略

ThreeLevelCache 允许调整 L1 大小和共享内存阈值(通过 core/config.py 中的 l1_max_sizeshm_threshold)。缓存键由 project_id + node + input_hash 组成,input_hash 根据项目参数计算,当参数变更时缓存自动失效。

7.3 与 FastAPI 集成

框架本身与 Web 无关,可在 FastAPI 应用启动时创建 ActorSystem 实例,并挂载到 app.state。路由中通过依赖注入获取,调用 start_project 触发计算。

示例main.py):

@asynccontextmanager
async def lifespan(app: FastAPI):
    system = ActorSystem()
    await system.init()
    app.state.actor_system = system
    yield
    await system.shutdown()

路由中

@router.post("/compute/{project_id}")
async def compute(project_id: str, req: ComputeRequest, system: ActorSystem = Depends(get_actor_system)):
    params = await load_config(project_id)  # 业务层加载配置
    await system.start_project(project_id, params, req.recompute_from, req.workflow)
    return {"status": "started"}

这种集成方式保持框架纯计算能力,Web 层只负责触发和状态查询。

8. 常见问题

Q1:节点函数中如何访问数据库?
A:可通过 config 字典传递数据库管理器,或在函数内部使用单例模式获取。建议将 db 作为 params 的一部分传入(例如 params["db"] = db),这样所有节点都能通过 config 参数访问。

Q2:如何调试节点执行?
A:设置日志级别为 DEBUG,框架会输出详细的调度、缓存命中、共享内存操作等日志。关键日志点包括节点就绪、执行开始、缓存命中、共享内存加载/释放。

Q3:流程执行完后如何通知外部?
A:可在 ProjectActor._schedule 完成分支中调用自定义回调(需修改框架),或通过轮询 Actor 是否还在运行(如示例中的 while 循环)。生产环境建议使用消息队列或数据库状态标志。

Q4:能否支持动态分支?
A:不支持。若需条件执行,可在节点函数内判断上游结果并决定是否返回空数据,下游节点再根据空数据跳过处理。例如:

def conditional_func(project_id, config, data):
    if data is None or len(data) == 0:
        return {"skip": True}
    # 正常处理...

Q5:节点执行超时怎么处理?
A:CPU 节点可通过 WorkerPool.run(timeout=) 设置超时;IO 节点需在函数内部使用 asyncio.wait_for 控制。超时后节点标记为 failed,可人工重算。

Q6:如何清空项目缓存?
A:调用 cache.evict_project(project_id) 可清除该项目的 L1 缓存和共享内存引用,L3 数据仍保留。如需完全重置,需同时清除 project_workflow_state 表记录。

9. 总结

本框架通过高度模块化设计,将调度、执行、缓存、持久化分离,开发者只需关注业务函数和工作流定义,即可快速搭建并行计算应用。其轻量级依赖和简洁设计使其特别适合资源受限的边缘计算场景,同时保留了足够的灵活性以支持复杂业务。

遵循本文规范,即可充分利用框架能力,构建稳定高效的微电网规划仿真系统。如需扩展,可参考源码中的 ProjectActorDAGSchedulerThreeLevelCache 等核心类,它们均提供了清晰的接口和扩展点。

参考:
肖永威. Python轻量级异步数据汇聚与并行计算框架—AsyncAggActor. CSDN博客. 2026.03
肖永威. Python 异步数据汇聚与并行计算框架设计与实现. CSDN博客. 2025.12

更多推荐