1. 背景与目标

1.1 背景

在工业智能控制领域,如综合能源管理、智能制造、过程控制等场景,系统需要实时处理时序数据,并执行一系列复杂的计算任务,包括预测(如负荷预测、光伏预测)、优化(如储能策略优化)、模型预测控制(MPC)等。这些任务往往存在数据依赖关系,构成一个有向无环图(DAG),且需要根据项目(如不同站点、不同客户)进行隔离计算。

典型的业务场景下,每个项目包含全年8760小时的逐小时数据,计算节点多、数据量大,且要求支持10~20个项目的并发处理。系统通常部署于单台服务器(如边缘网关、本地工作站),资源有限(CPU、内存),但必须保证高吞吐、低延迟和稳定性。

现有解决方案面临的主要挑战:

  • 并发模型选择:传统的多线程方案受GIL限制,无法充分利用多核;纯异步方案无法处理CPU密集型计算;多进程方案又面临进程间通信开销大、资源隔离复杂的问题。
  • 任务依赖管理:硬编码的串行流程难以扩展,且不支持灵活的重算(如修正某个中间节点后自动重算下游)。
  • 数据复用:相同输入(如8760时间序列)可能被多个下游节点反复使用,若重复计算或传递大量数据,将造成严重的性能浪费。
  • 资源消耗:在单机环境下,必须精细控制内存和CPU使用,避免因大量并发导致系统崩溃。
  • 数据库设计:需要同时支持缓存持久化(L3)和业务数据(如气象、设备参数)的存取,且要避免相互影响,实现异步非阻塞访问。

针对上述问题,业界已有多种探索。例如,CSDN上的一篇经典文章[1]提出了一种基于“异步+多进程”的轻量级数据汇聚框架,通过asyncio处理I/O并发,multiprocessing执行预测任务,以队列解耦各模块。该方案简洁高效,但缺乏对任务依赖管理和缓存复用的支持,且手动管理进程和队列存在健壮性隐患。

另一种思路是引入Actor模型,如使用Pykka框架,将每个项目封装为一个Actor,通过消息驱动并发。Actor模型天然支持状态隔离和容错,但Pykka的ThreadingActor每个Actor占用一个线程,高并发下线程开销大,且仍需自行处理DAG调度和缓存。

1.2 目标

本框架旨在设计并实现一个轻量级、高性能、可扩展的异步数据汇聚与并行计算调度引擎,满足以下核心需求:

1.2.1. 功能需求

  • 项目级隔离:每个项目拥有独立的计算状态和缓存,互不干扰。
  • 数据汇聚:按项目标识(如项目ID)自动汇聚多源输入数据,形成完整的计算上下文。
  • DAG任务调度:支持用户通过简单配置定义任务依赖关系,调度器自动检测就绪节点并触发执行,无需硬编码流程。
  • 三级缓存复用
    • L1(内存缓存):使用LRU策略存储近期结果,加速重复查询。
    • L2(共享内存缓存):针对大尺寸数组(如8760时间序列),存入共享内存实现零拷贝传递,避免进程间序列化开销。
    • L3(持久化缓存):基于SQLite存储计算结果,确保重启后可复用,并启用WAL模式支持并发读写。
  • 灵活重算:支持从任意节点强制重算,自动重置其下游节点状态,并重新调度执行。
  • 混合并发:利用asyncio处理高并发I/O(如HTTP请求、数据库访问),利用ProcessPoolExecutor执行CPU密集型计算,兼顾吞吐与计算性能。
  • 异步数据库池化:实现业务数据库与缓存数据库分离,各自拥有独立的读写连接池,避免相互阻塞。
  • 节点类型支持:区分CPU密集型、I/O密集型和本地计算型节点,采用不同的执行策略。
  • 外部消息集成:支持通过RabbitMQ等消息队列接收任务触发指令,实现与外部系统的无缝对接。
  • 优雅退出与资源清理:支持捕获终止信号,确保子进程、共享内存等资源正确释放,避免泄漏。

1.2.2. 非功能需求

  • 轻量级:代码量可控,无外部依赖(仅使用标准库及少量成熟第三方库如aiohttpnetworkxaiosqlite),便于部署和维护。
  • 高性能:在单机8核16GB环境下,可稳定支持20个项目并发,每个项目处理8760数据,端到端延迟可接受。
  • 可扩展性:新增计算节点只需在DAG中注册,无需修改核心逻辑;缓存层可灵活替换或升级。
  • 健壮性:通过引用计数管理共享内存生命周期,异常捕获与重试机制保障系统稳定运行。
  • 易调试:完善的日志系统,支持按模块和级别输出,便于问题追踪。
  • 不追求分布式能力,而是专注于单机性能最大化,适用于边缘计算、中小型监控系统或开发测试环境。

通过上述设计,本框架将为工业智能控制场景提供一个可靠、高效的本地计算底座,同时为未来可能的分布式扩展预留接口。

2. 设计理念

2.1. Actor 模型

Actor 模型是一种并发计算模型,其核心思想是将计算实体抽象为独立的“Actor”,每个 Actor 拥有私有状态、行为逻辑和消息队列,仅通过异步消息通信交互。这种模型天然适合构建高并发、可扩展、容错的系统。

在本框架中,每个项目对应一个 Actor(ProjectActor,拥有独立的消息队列和内部状态(DAG 节点状态、计算结果)。项目之间完全隔离,通过 ActorSystem 统一管理。消息驱动机制(如 STARTSHUTDOWN)使得项目计算流程清晰可控,无需显式加锁,大大简化了并发编程复杂性。

2.2. 异步 I/O 设计

框架基于 Python 的 asyncio 构建异步事件循环,充分利用其处理高并发 I/O 的能力:

  • 网络层:FastAPI 异步处理 HTTP 请求,非阻塞。
  • 数据库:使用 aiosqlite 异步访问 SQLite,避免阻塞事件循环。
  • Actor 内部:每个 Actor 的 run 协程循环从队列中获取消息,异步处理。
  • 进程池桥接:对于 CPU 密集型计算,通过 loop.run_in_executor 将任务提交给 ProcessPoolExecutor,执行完成后通过 Future 回调恢复协程,实现“异步 + 多进程”的无缝结合。
  • 消息队列集成:通过 asyncio.to_thread 将阻塞的 RabbitMQ 消费操作放入线程池,避免阻塞事件循环,实现异步消息驱动。

2.3. 共享内存缓存设计

针对 8760 小时时序数据(numpy 数组)的复用需求,设计了三级缓存架构:

  • L1 内存缓存:使用 OrderedDict 实现 LRU 淘汰策略,存储最近访问的标量结果或共享内存元数据。
  • L2 共享内存缓存:基于 multiprocessing.shared_memory 实现,将大数组存入共享内存,仅返回元数据(shm_name, shape, dtype)。后续节点可直接映射该内存,实现零拷贝读取,避免进程间序列化开销。
  • L3 持久化缓存:使用 SQLite 存储计算结果(标量或数组的 pickle 序列化),启用 WAL 模式支持读写并发,确保重启后缓存可复用。

共享内存生命周期管理是关键难点。通过 SharedMemoryManager 维护两个核心结构:

  • _handles:保存所有活跃的 SharedMemory 对象引用,防止 Python 垃圾回收误删。
  • _refcount:引用计数器,每增加一个持有者(L1 缓存或调用者如 node_results)就递增,释放时递减,归零时执行 close/unlink

2.4. 双数据库分离与异步连接池

系统同时管理两类数据库:

  • 缓存数据库(cache.db):用于 L3 持久化缓存,存储计算中间结果和最终结果。采用单文件、WAL 模式,读写并发。

  • 业务数据库(energy.db):存储项目配置、气象数据(如 PVGIS TMY)、设备参数等业务数据。通过 AsyncSQLiteManager 实现读写分离连接池:

    • 写连接:一个长期持有的写连接,配合 asyncio.Lock 确保事务串行化。
    • 读连接池:多个读连接,每次查询从池中借用一个连接,使用后归还,提高读并发能力。
    • 启用 WAL 模式,优化并发性能。

双库分离避免了缓存写入对业务查询的影响,且各自可根据负载独立调优。

2.5. 消息驱动扩展

为支持与外部系统集成,框架引入了 RabbitMQ 消息消费者模块。其设计理念是:

  • 解耦:任务触发与执行逻辑分离,外部系统只需向指定队列发送消息,无需直接调用 API。

  • 异步非阻塞:消息消费采用线程+异步协程结合的方式,避免阻塞主事件循环。

  • 可靠消费:通过手动确认机制和断线重连,确保消息不丢失。

3. 技术选型对比分析

3.1 三种架构方案对比

维度 异步数据汇聚 (CSDN方案) 异步Actor模型 (本方案) Pykka Actor方案
开发复杂度 非常高:需手动管理队列、超时、进程生命周期 高:需自研Actor框架、DAG调度和缓存,但核心逻辑清晰 :直接使用成熟库,API简单
代码量/简洁性 繁琐,代码量大(约300行核心,但功能单一) 适中(约800行),概念简洁但实现需精细 非常简洁(约200行即可实现基础Actor)
性能 最高(理论):无框架开销,极致优化 高:有消息传递和缓存开销,但通过共享内存优化 较高:线程级Actor有上下文切换和GIL限制
健壮性/容错 中:依赖手动实现错误处理和资源回收 非常高:通过Actor监督策略、引用计数和异常捕获,可构建健壮系统 非常高:内置监督机制,成熟稳定
可维护性/扩展性 差:功能耦合,新增节点需改多处代码 :模块化设计,DAG注册新节点即可 非常好:标准化Actor模式,易于理解和扩展
学习曲线 陡峭:需精通asynciomultiprocessing和队列管理 陡峭:需理解Actor模型、DAG和引用计数 平缓:只需学习Pykka API和Actor基本概念
适用场景 性能极致敏感、一次性或高度定制的流水线项目 复杂的、需要高容错和灵活重算的单机并发系统 绝大多数需要并发的项目,快速开发和维护

3.2 进程池实现方式对比

对比维度 原生 multiprocessing (CSDN方案) concurrent.futures.ProcessPoolExecutor (本方案)
抽象层次 低,直接操作Process和Queue 高,提供Future接口和任务提交
生命周期管理 手动启动、join、terminate,易遗漏 自动管理,通过with语句或shutdown()优雅关闭
异常传播 需自行捕获并通过队列传递 异常会自动封装在Future中,主进程可捕获
与asyncio集成 loop.run_in_executor配合,较麻烦 原生支持run_in_executor,无缝集成
资源回收 易产生僵尸进程 自动回收,保证进程退出
适用性 适合简单、一次性任务 适合生产环境,稳定可靠

4. 架构设计

4.1. 总体架构

系统对外提供两种触发方式:

  • REST API:通过 FastAPI 直接调用。
  • 消息队列:通过 RabbitMQ 消费者模块接收任务消息,转换为内部 Actor 消息。

4.1.1. HTTP API模式

Client → FastAPI → ActorSystem → ProjectActor → DAGScheduler → WorkerPool (ProcessPoolExecutor)
                                              ↕                ↕
                                          三级缓存         共享内存
                                              ↕
                                         SQLite持久化

4.1.2. 消息模式

外部消息源 (RabbitMQ)
       ↓
[MessageConsumer] (线程+异步)
       ↓
  ActorSystem → ProjectActor → DAGScheduler → WorkerPool (ProcessPoolExecutor)
       ↕              ↕              ↕
   三级缓存       共享内存       业务数据库池
       ↕
 SQLite缓存库 (cache.db)

4.2. 核心模块

  • ActorSystem:管理所有项目Actor,负责创建、路由消息。
  • ProjectActor:每个项目一个Actor,维护DAG状态,通过消息队列驱动。
  • DAGScheduler:基于networkx管理节点依赖,自动检测就绪任务。
  • ThreeLevelCache:三级缓存,支持共享内存引用计数。
  • WorkerPool:封装ProcessPoolExecutor,提交任务并支持超时。
  • SharedMemoryManager:管理共享内存生命周期。
  • MessageConsumer: RabbitMQ消息消费者,运行在独立线程中,将消息转换为Actor任务。
  • RabbitMqManager: 封装RabbitMQ连接、发送、接收等底层操作。
  • AsyncSQLiteManager: 异步数据库连接池,管理业务数据库的读写连接。

4.2.1. ActorSystem 与 ProjectActor

  • ActorSystem 维护一个字典 {project_id: ProjectActor},并提供 get_or_create_actor
  • ProjectActor 继承自普通类(非Pykka,简化),拥有自己的 asyncio.Queue 消息队列和 run 协程。处理 STARTSHUTDOWN 消息。
  • 支持从指定节点重算:_start(recompute_from) 调用 DAGScheduler.reset_subtree 并释放相关共享内存。

4.2.2. DAGScheduler

  • 基于 networkx.DiGraph 存储节点和依赖。
  • add_node(name, func, deps) 注册节点。
  • get_ready() 返回所有依赖已完成的节点。
  • reset_subtree(node) 重置节点及其下游状态,并清除结果。

4.2.3. 三级缓存 ThreeLevelCache

  • L1 (RAM)OrderedDict 实现LRU淘汰,存储标量结果或共享内存元数据。
  • L2 (共享内存):通过 SharedMemoryManager 管理,大数组(>1KB)存入共享内存,返回元数据(shm_name, shape, dtype)。
  • L3 (SQLite):持久化存储,启用WAL模式,异步写入。
  • 引用计数:每个元数据由L1和调用者(如 node_results)共同持有,通过 store(L1引用)、load(调用者引用)、try_add_ref(检查存活)和 release 管理。
  • 关键修复:SharedMemoryManager 持有 _handles 字典,防止共享内存对象被GC回收。

4.2.4. WorkerPool

  • 封装 ProcessPoolExecutor,提供 run(func, *args, timeout) 方法,支持超时和异常捕获。
  • 进程池大小由配置决定,避免过多进程。

4.2.5. 异步数据库连接池 AsyncSQLiteManager

  • 单例模式:确保全局只有一个管理器实例。
  • 写连接:一个长期持有的写连接,通过 asyncio.Lock 保证串行写入。
  • 读连接池:预创建多个读连接,支持并发查询。
  • 初始化:在 ActorSystem.init() 中调用 AsyncSQLiteManager.init(),建立写连接和读池。
  • 操作接口
    • write(sql, params):支持单条或批量写入,自动提交。
    • query(sql, params):从池中借读连接执行查询,返回字典列表。
    • query_df(sql, params):返回 pandas DataFrame。
  • WAL配置:启用 journal_mode=WALsynchronous=NORMAL 等优化参数,提升并发性能。

4.2.6. 节点类型与执行策略

  • CPU节点:通过 WorkerPool 提交到进程池执行,利用多核。
  • IO节点:直接在异步事件循环中执行(或通过 asyncio.to_thread 处理同步IO),如 fetch_weather 使用 asyncio.to_thread 调用同步PVGIS获取,同时利用连接池查询本地缓存。
  • Local节点:直接在 Actor 线程中执行,适合简单计算。

4.2.7. 消息消费者 MessageConsumer

  • 作用:监听RabbitMQ队列,接收外部任务消息,调用 ActorSystem.start_project 触发计算。
  • 设计
    • 采用后台线程运行RabbitMQ消费循环(因为 pika 是同步库,不能直接集成到 asyncio)。
    • 消费到的消息通过线程安全的 queue.Queue 传递给异步协程 _consume_loop
    • 异步协程从队列中取出消息,调用 actor_system.start_project 执行。
    • 支持优雅退出:通过 exit_event 通知线程和协程停止,并等待资源释放。

4.2.8. RabbitMqManager

  • 封装RabbitMQ连接参数和基础操作:
    • getConnectionMQ():建立带认证的阻塞连接。
    • consume_from_rabbitmq_and_enqueue():在后台线程中运行的消费函数,将消息放入目标队列。
    • sendMSGtoRabbitmq():向指定队列发送消息(可用于结果回传,目前未集成)。
    • getMqMSG_count():查询队列中的消息数量,用于监控。

4.2.9. 工作流管理

为适应不同业务场景(如全流程计算、可行性研究、柔性控制等),框架支持多工作流模板。每个工作流定义了一组计算节点及其依赖关系,项目在创建时通过 workflow 参数指定要使用的模板。

每个节点由 (name, func_name, node_type, dependencies) 定义:

  • name:节点标识符。
  • func_name:对应计算函数的字符串名称,需在 actor.pyFUNCTION_MAP 中映射到实际函数。
  • node_type:节点类型(cpu/io/local)。
  • dependencies:依赖的节点列表。

动态 DAG 构建ProjectActor 在初始化时根据 self.workflowWORKFLOWS 中读取节点定义,动态构建 DAG,实现逻辑与配置分离,新增工作流无需修改核心代码。

5. 关键技术实现

5.1. 共享内存引用计数

# 核心逻辑
store()   # 增加L1引用
load()    # 增加调用者引用
release() # 减少引用,归零时unlink
_handles  # 持有对象防止GC

5.2. DAG重算支持

  • recompute_from参数传入后,使用nx.descendants获取下游节点。
  • 释放下游节点在node_results中的共享内存引用。
  • 调用reset_subtree重置状态,调度器自动重新执行。

5.3. 异步与进程桥接

在异步混合并发架构中,异步(asyncio)与多进程(multiprocessing)分别擅长处理I/O密集型与CPU密集型任务。异步依赖单线程事件循环,若直接执行计算将阻塞循环,导致所有并发任务等待;而多进程世界独立于主进程,可并行利用多核。通过 loop.run_in_executor() 将CPU任务提交给 ProcessPoolExecutor,实现了异步与进程的“桥接”:异步调用方以 await 等待,事件循环不被阻塞;计算在独立进程中进行,结果通过Future返回。这一机制既保持了异步高吞吐,又充分利用了多核计算资源,且由标准库自动管理进程生命周期,规避了手动IPC的复杂性与风险。

result = await loop.run_in_executor(self.pool, func, *args)

5.4. 节点类型分发

ProjectActor._execute_node 中,根据 node_type 选择执行方式:

  • "cpu"await self.worker_pool.run(...)
  • "io"await func(...)(若内部有阻塞调用,应使用 asyncio.to_thread
  • "local" → 直接调用 func(...)

5.5. 业务数据库异步访问

fetch_weather 为例:

  • 生成区域ID,查询本地 pvgsis_tmy 表(通过 AsyncSQLiteManager.query)。
  • 若缓存未命中,通过 asyncio.to_thread 调用同步PVGIS API获取数据。
  • 将结果写入数据库(AsyncSQLiteManager.write)并返回DataFrame。
  • 整个过程完全非阻塞,且利用了连接池复用连接。

5.6. 消息队列异步消费集成

  • 线程与协程协作
    • 后台线程运行同步的 consume_from_rabbitmq_and_enqueue,将消息放入 queue.Queue
    • 异步协程 _consume_loop 通过 loop.run_in_executor 等待队列(避免阻塞事件循环),取出消息后调用 actor_system.start_project
  • 优雅退出:通过 threading.Event 通知线程退出,协程检测到退出事件后停止循环。
  • 错误处理:消费线程遇到连接错误会记录日志并退出,上层可重启(当前未实现自动重启,可通过系统守护进程管理)。

5.7. 日志系统

  • 使用 TimedRotatingFileHandler 实现按天轮转的日志文件,保留最近7天。
  • 日志格式包含时间、模块名、线程名、级别,便于追踪。
  • aiosqlite 单独设置日志级别,避免其调试信息污染主日志。

6. 运行与测试

6.1. 环境准备

pip install fastapi uvicorn aiosqlite networkx numpy pydantic-settings

6.2. 启动服务

uvicorn main:app --reload

6.3. API测试

# 从头计算项目 proj_001
curl -X POST http://localhost:8000/projects/proj_001/compute -H "Content-Type: application/json" -d "{\"params\": {}, \"workflow\": \"default\"}"

# 查询状态
curl http://localhost:8000/projects/proj_001/status

# 从 storage 节点重算
curl -X POST http://localhost:8000/projects/proj_001/compute -H "Content-Type: application/json" -d "{\"recompute_from\": \"storage\", \"workflow\": \"default\"}"

为方便自动化测试或集成,可使用 Python 的 requests 库编写客户端脚本,通过 HTTP API 触发项目计算并轮询状态直至完成。以下是一个完整示例,演示了如何创建项目、启动计算并等待所有节点执行完毕。

6.3.1 Python 客户端示例

import requests
import time

BASE_URL = "http://127.0.0.1:8000"
project_id = "test_pv_001"

# 项目参数(与消息队列测试示例一致)
params = {
    "latitude": 29.660784,
    "longitude": 121.444672,
    "altitude": 100,
    "tz": "Asia/Shanghai",
    "pvsystem": [
        {
            "array": [
                {"surface_tilt": 2.06, "surface_azimuth": 198, "pdc0": 92700},
                {"surface_tilt": 2.06, "surface_azimuth": 18, "pdc0": 92700}
            ],
            "albedo": 0.25,
            "gamma_pdc": -0.0034,
            "temp_ref": 25,
            "name": "main_array"
        }
    ]
}

# 1. 发起计算(使用默认工作流)
resp = requests.post(f"{BASE_URL}/projects/{project_id}/compute", json={"params": params})
print("计算启动响应:", resp.json())
# 预期输出:{"project_id":"test_pv_001","nodes":{...},"results":{}}

# 2. 轮询状态,直到所有节点完成
while True:
    status = requests.get(f"{BASE_URL}/projects/{project_id}/status").json()
    print(f"当前节点状态: {status['nodes']}")
    # 检查是否所有节点状态都是 "done"
    if all(v == "done" for v in status["nodes"].values()):
        print("计算完成!")
        # 获取最终结果
        results = status["results"]
        print("光伏结果元数据:", results.get("pv"))
        # 若需要读取共享内存中的数组,可根据元数据中的 shm_name、shape、dtype 进行映射
        break
    time.sleep(2)  # 等待2秒后再次查询

6.3.2 运行说明

  • 确保服务已启动(uvicorn main:app --reload)。
  • 将上述脚本保存为 test_api.py 并运行。
  • 脚本会每2秒查询一次项目状态,直至所有节点完成。最终输出各节点的结果(标量或共享内存元数据)。

6.3.3 预期输出示例

计算启动响应: {'project_id': 'test_pv_001', 'nodes': {'pv': 'pending', 'load': 'pending', 'price': 'pending', 'storage': 'pending', 'strategy': 'pending', 'loss': 'pending'}, 'results': {}}
当前节点状态: {'pv': 'running', 'load': 'pending', 'price': 'pending', 'storage': 'pending', 'strategy': 'pending', 'loss': 'pending'}
当前节点状态: {'pv': 'done', 'load': 'done', 'price': 'done', 'storage': 'running', 'strategy': 'pending', 'loss': 'pending'}
当前节点状态: {'pv': 'done', 'load': 'done', 'price': 'done', 'storage': 'done', 'strategy': 'running', 'loss': 'pending'}
当前节点状态: {'pv': 'done', 'load': 'done', 'price': 'done', 'storage': 'done', 'strategy': 'done', 'loss': 'running'}
当前节点状态: {'pv': 'done', 'load': 'done', 'price': 'done', 'storage': 'done', 'strategy': 'done', 'loss': 'done'}
计算完成!
光伏结果元数据: {'shm_name': 'wnsm_0d0de4a0', 'shape': [8760], 'dtype': 'float32'}

6.4. RabbitMQ 消息触发测试

为验证消息队列集成功能,使用提供的测试脚本 TestSendMSG.py 向 RabbitMQ 发送任务消息,触发项目计算。

6.4.1 测试脚本说明

测试脚本 TestSendMSG.py 主要功能:

  • 从命令行读取 JSON 配置文件路径。
  • 读取配置文件中的消息内容(需包含 project_idworkflowparams 字段)。
  • 使用 pika 库连接 RabbitMQ,将消息发送到指定交换机(energyStorageStrategy.fanout)和队列(energyStorageStrategy.queue),路由键为 typc-fpd-tysh
# TestSendMSG.py (关键代码片段)
def sendMSG(message, exchange_name, queue_name, routing_key):
    credentials = pika.PlainCredentials('rabbit', '*******')
    connection = pika.BlockingConnection(pika.ConnectionParameters(
        '192.168.17.52', port=55671, virtual_host='/pvet-dev', credentials=credentials))
    channel = connection.channel()
    channel.queue_declare(queue=queue_name, durable=True)
    channel.exchange_declare(exchange=exchange_name, exchange_type='fanout', durable=True)
    channel.queue_bind(queue=queue_name, exchange=exchange_name, routing_key=routing_key)
    channel.basic_publish(
        exchange=exchange_name,
        routing_key=routing_key,
        body=json.dumps(message, ensure_ascii=False),
        properties=pika.BasicProperties(delivery_mode=2)
    )
    connection.close()
    print(f"发送消息:{message}")

注意:连接参数(主机、端口、虚拟主机、凭据)需根据实际环境修改,与 RabbitMqManager.py 中的配置保持一致。

6.4.2 测试消息格式

消息需为 JSON 格式,包含以下字段:

  • project_id:项目标识(字符串)。
  • workflow:工作流类型(字符串,如 "default""feasibility"),未提供时默认为 "default"
  • params:项目参数字典,包含计算所需的所有配置(如气象坐标、光伏系统参数等)。

示例消息文件 testdat.json

{
    "project_id": "test_pv_001",
    "workflow": "default",
    "params": {
        "latitude": 29.660784,
        "longitude": 121.444672,
        "altitude": 100,
        "tz": "Asia/Shanghai",
        "pvsystem": [
            {
                "array": [
                    {"surface_tilt": 2.06, "surface_azimuth": 198, "pdc0": 92700},
                    {"surface_tilt": 2.06, "surface_azimuth": 18, "pdc0": 92700}
                ],
                "albedo": 0.25,
                "gamma_pdc": -0.0034,
                "temp_ref": 25,
                "name": "main_array"
            }
        ]
    }
}

6.4.3 测试步骤

  1. 确保 RabbitMQ 服务已启动,且连接参数与测试脚本一致。
  2. 启动框架服务
    uvicorn main:app --reload
    
    观察日志,确认消费者线程已启动:
    INFO - RabbitMQ consumer thread started
    INFO - Async message consumer task started
    
  3. 运行测试脚本
    python TestSendMSG.py
    
    根据提示输入消息文件名(如 test_pv_001.json),脚本将发送消息到 RabbitMQ。

6.5. 压力测试

可使用 aiohttp 编写简单并发脚本,模拟20用户同时请求。

7. 代码结构

core/
├── main.py                  # FastAPI应用入口,生命周期管理
├── config.py                # 配置类
├── models.py                # Pydantic数据模型
├── routes.py                # API路由定义
├── system.py                # ActorSystem 实现
├── actor.py                 # ProjectActor 实现
├── dag.py                   # DAGScheduler 实现
├── cache.py                 # 三级缓存(含SharedMemoryManager)
├── worker.py                # 进程池封装 WorkerPool
├── project_model.py         # 业务计算模型及数据获取
├── AsyncSQLiteManager.py    # 异步数据库连接池
├── mq_consumer.py           # RabbitMQ 消息消费者管理
├── RabbitMqManager.py       # RabbitMQ 底层连接与消费/发送函数
└── requirements.txt         # 依赖列表

gitee 仓库地址:https://gitee.com/xiaoyw71/AsyncAggActor
分支:main
使用说明:请参考 README.md 中的安装和运行步骤。

8. 风险与完善方向

8.1 当前风险

  • 共享内存引用计数:在极端并发下仍可能有竞态条件。
  • SQLite写入锁:高并发写L3时可能出现database is locked,需增加重试机制。
  • 内存占用:L1缓存若配置过大,可能占用过多内存。
  • 进程池大小:需根据CPU核数和任务类型动态调整。
  • 数据库连接池耗尽:若读池大小不足,高并发查询可能导致等待,需监控并调整。
  • RabbitMQ连接可靠性:当前消费者线程遇到连接错误会退出,未实现自动重连,可能因网络抖动导致服务中断。
  • 消息确认机制:当前使用 auto_ack=True,若处理过程中崩溃可能导致消息丢失。需改为手动确认,并在处理完成后确认。
  • 工作流配置冲突:需确保不同工作流中的节点名称不重复(除非意图共享),否则可能引起依赖混乱。建议节点命名遵循约定,如 storage_simplestorage 区分。

8.2 后续完善计划

  1. 增加监控与可观测性:集成Prometheus,暴露队列长度、缓存命中率、任务耗时等指标。
  2. 动态资源调整:根据系统负载自动调整生成频率和进程池大小。
  3. 增强错误处理:引入重试机制和死信队列。
  4. Web控制面板:提供实时状态查询和参数调整接口。
  5. 支持分布式扩展:预留接口,未来可对接消息队列(如Redis Streams)实现多机部署。
  6. 完善单元测试:覆盖核心模块的边界条件和异常场景。
  7. 优化数据库连接池:增加连接健康检查、自动重连机制。

9. 结论

本方案在参考经典异步数据汇聚方案的基础上,引入了Actor模型、DAG调度和三级缓存,解决了复杂任务依赖、大数组复用、灵活重算、高效数据存取和外部系统集成等多个核心需求。通过与Pykxa方案的对比,明确了自研Actor框架在性能和控制力上的优势,同时也承认其在开发复杂度上的代价。最终形成的轻量级框架,完全符合单机高并发、低外部依赖的设计目标,为工业智能控制等场景提供了可靠的计算基础。


[1] Python 异步数据汇聚与并行计算框架设计与实现

更多推荐