Azure Durable Functions:如何让活动共享大型数据对象?
解决方案:让多个Azure Activity共享大型字典避免内存复制
针对你遇到的Azure Durable Functions Activity无法共享大型字典、导致内存占用过高的问题,以下是几个实用的解决思路:
1. 借助本地共享缓存(如Redis)
将大型字典存入共享缓存,仅向每个Activity传递缓存键和target参数,Activity通过键读取缓存中的字典,全程无需复制完整字典。
实现步骤:
- 在Orchestrator中,先将大字典序列化后存入Redis(或其他内存缓存服务),生成唯一缓存键。
- 每个并行Activity的调用参数仅包含
target和缓存键,打包为单个JSON对象传递。 - Activity接收参数后,通过缓存键从Redis中读取并反序列化字典,完成扫描逻辑。
代码示例:
Orchestrator 侧
import redis import json from azure.durable_functions import DurableOrchestrationContext async def orchestrator_function(context: DurableOrchestrationContext): # 初始化Redis连接 r = redis.Redis(host="your-redis-host", port=6379, db=0) big_dict = {"key1": "value1", ...} # 你的大型字典 cache_key = "shared_big_dict" # 将字典存入Redis r.set(cache_key, json.dumps(big_dict)) # 生成30个并行Activity任务,仅传递target和缓存键 targets = [f"target_{i}" for i in range(30)] tasks = [] for target in targets: input_data = {"target": target, "cache_key": cache_key} tasks.append(context.call_activity("ScanBigDictActivity", input_data)) results = await context.task_all(tasks) # 任务完成后清理缓存 r.delete(cache_key) return results
Activity 侧
import redis import json async def main(input_data: dict) -> str: target = input_data["target"] cache_key = input_data["cache_key"] # 从Redis读取字典 r = redis.Redis(host="your-redis-host", port=6379, db=0) big_dict_str = r.get(cache_key) big_dict = json.loads(big_dict_str) # 执行扫描逻辑 result = f"Scanned {target}: found {big_dict.get(target, 'no match')}" return result
2. 利用单进程Worker的全局变量共享
如果你的Function App配置为单进程运行(设置WEBSITE_WORKER_PROCESS_COUNT=1),可以在Worker启动时将大字典加载到全局变量中,所有Activity实例直接访问该全局变量,无需传递字典参数。
实现步骤:
- 在Activity所在的Python模块中,定义全局变量存储大字典,通过
@app.on_startup装饰器在Worker启动时加载字典。 - Activity的
main函数仅接收target参数,直接使用全局变量中的大字典。
代码示例:
import azure.functions as func import json # 全局变量存储大型字典 big_dict = {} @app.on_startup def load_big_dict(): global big_dict # 从文件/数据库等加载大型字典 with open("big_dict.json", "r") as f: big_dict = json.load(f) async def main(target: str) -> str: # 直接使用全局的big_dict result = f"Scanned {target}: found {big_dict.get(target, 'no match')}" return result
注意:此方案仅适用于单进程Worker场景,若配置多进程,每个进程会独立加载一份字典,无法实现共享。
3. 写入临时文件/Blob存储,按需读取
将大字典序列化后写入本地临时文件或Azure Blob Storage,Activity通过文件路径/Blob SAS URL读取字典,避免内存中复制多份。
实现步骤:
- Orchestrator将大字典序列化(推荐用
pickle支持复杂Python类型),写入本地临时目录或Blob Storage。 - 向每个Activity传递
target和文件路径/Blob SAS URL,Activity读取文件并反序列化字典。 - 任务完成后清理临时文件/Blob。
代码示例(本地临时文件):
Orchestrator 侧
import pickle import tempfile from azure.durable_functions import DurableOrchestrationContext async def orchestrator_function(context: DurableOrchestrationContext): big_dict = {"key1": "value1", ...} # 创建临时文件存储字典 with tempfile.NamedTemporaryFile(delete=False, suffix=".pkl") as f: pickle.dump(big_dict, f) temp_file_path = f.name targets = [f"target_{i}" for i in range(30)] tasks = [] for target in targets: input_data = {"target": target, "file_path": temp_file_path} tasks.append(context.call_activity("ScanBigDictActivity", input_data)) results = await context.task_all(tasks) # 清理临时文件 import os os.unlink(temp_file_path) return results
Activity 侧
import pickle async def main(input_data: dict) -> str: target = input_data["target"] file_path = input_data["file_path"] # 读取临时文件中的字典 with open(file_path, "rb") as f: big_dict = pickle.load(f) result = f"Scanned {target}: found {big_dict.get(target, 'no match')}" return result
内容的提问来源于stack exchange,提问作者Ksenia Semenova
相关产品推荐
相关产品推荐

