Python轻量级异步数据汇聚与并行计算框架—AsyncAggActor
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. 非功能需求
- 轻量级:代码量可控,无外部依赖(仅使用标准库及少量成熟第三方库如
aiohttp、networkx、aiosqlite),便于部署和维护。 - 高性能:在单机8核16GB环境下,可稳定支持20个项目并发,每个项目处理8760数据,端到端延迟可接受。
- 可扩展性:新增计算节点只需在DAG中注册,无需修改核心逻辑;缓存层可灵活替换或升级。
- 健壮性:通过引用计数管理共享内存生命周期,异常捕获与重试机制保障系统稳定运行。
- 易调试:完善的日志系统,支持按模块和级别输出,便于问题追踪。
- 不追求分布式能力,而是专注于单机性能最大化,适用于边缘计算、中小型监控系统或开发测试环境。
通过上述设计,本框架将为工业智能控制场景提供一个可靠、高效的本地计算底座,同时为未来可能的分布式扩展预留接口。
2. 设计理念
2.1. Actor 模型
Actor 模型是一种并发计算模型,其核心思想是将计算实体抽象为独立的“Actor”,每个 Actor 拥有私有状态、行为逻辑和消息队列,仅通过异步消息通信交互。这种模型天然适合构建高并发、可扩展、容错的系统。
在本框架中,每个项目对应一个 Actor(ProjectActor),拥有独立的消息队列和内部状态(DAG 节点状态、计算结果)。项目之间完全隔离,通过 ActorSystem 统一管理。消息驱动机制(如 START、SHUTDOWN)使得项目计算流程清晰可控,无需显式加锁,大大简化了并发编程复杂性。
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模式,易于理解和扩展 |
| 学习曲线 | 陡峭:需精通asyncio、multiprocessing和队列管理 |
陡峭:需理解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协程。处理START和SHUTDOWN消息。- 支持从指定节点重算:
_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=WAL、synchronous=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通知线程和协程停止,并等待资源释放。
- 采用后台线程运行RabbitMQ消费循环(因为
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.py的FUNCTION_MAP中映射到实际函数。node_type:节点类型(cpu/io/local)。dependencies:依赖的节点列表。
动态 DAG 构建:ProjectActor 在初始化时根据 self.workflow 从 WORKFLOWS 中读取节点定义,动态构建 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_id、workflow、params字段)。 - 使用
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 测试步骤
- 确保 RabbitMQ 服务已启动,且连接参数与测试脚本一致。
- 启动框架服务:
观察日志,确认消费者线程已启动:uvicorn main:app --reloadINFO - RabbitMQ consumer thread started INFO - Async message consumer task started - 运行测试脚本:
根据提示输入消息文件名(如python TestSendMSG.pytest_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_simple与storage区分。
8.2 后续完善计划
- 增加监控与可观测性:集成Prometheus,暴露队列长度、缓存命中率、任务耗时等指标。
- 动态资源调整:根据系统负载自动调整生成频率和进程池大小。
- 增强错误处理:引入重试机制和死信队列。
- Web控制面板:提供实时状态查询和参数调整接口。
- 支持分布式扩展:预留接口,未来可对接消息队列(如Redis Streams)实现多机部署。
- 完善单元测试:覆盖核心模块的边界条件和异常场景。
- 优化数据库连接池:增加连接健康检查、自动重连机制。
9. 结论
本方案在参考经典异步数据汇聚方案的基础上,引入了Actor模型、DAG调度和三级缓存,解决了复杂任务依赖、大数组复用、灵活重算、高效数据存取和外部系统集成等多个核心需求。通过与Pykxa方案的对比,明确了自研Actor框架在性能和控制力上的优势,同时也承认其在开发复杂度上的代价。最终形成的轻量级框架,完全符合单机高并发、低外部依赖的设计目标,为工业智能控制等场景提供了可靠的计算基础。
更多推荐



所有评论(0)