能否用Snakemake顺序执行单文件API查询任务以规避限流与重启风险
Snakemake实现单文件串行API调用+断点续跑方案
完全可以实现你的需求,核心是通过checkpoint动态生成单文件任务+自定义资源限制串行执行,既规避API限流,又支持断点续跑,失败后仅重跑未完成的单个文件。
具体实现步骤
1. 保留文件拆分的Checkpoint(原有逻辑复用)
假设你已有拆分大文件的split checkpoint(生成15万+单个json文件),继续保留该逻辑,确保所有待处理文件都被生成在resources/split-files/目录下。
2. 动态捕获所有待处理文件名
通过函数从拆分目录中获取所有文件名,用于后续生成单个任务:
import os from snakemake.io import glob_wildcards def get_all_split_filenames(wildcards): # 从split checkpoint的输出目录获取所有文件名 split_dir = checkpoints.split.get(**wildcards).output.output_dir return glob_wildcards(os.path.join(split_dir, "{filename}.json")).filename
3. 定义单文件API调用规则
将原批量API调用拆分为单个文件的任务,每个任务对应一个输入输出文件:
rule query_single_file: input: json_file="resources/split-files/{filename}.json" output: enriched_json="resources/enriched-files/{filename}.json" resources: mem="2GB", # 单个任务资源需求大幅降低 runtime="1h", api_slot=1 # 自定义资源,用于控制串行执行 retries: 3 # 可选:API请求失败自动重试 script: "scripts/query_single_api.py"
4. 定义最终触发规则
用expand生成所有单个任务的输出依赖,触发整个流程:
rule all: input: expand("resources/enriched-files/{filename}.json", filename=get_all_split_filenames)
脚本适配(query_single_api.py)
将原批量处理脚本修改为单文件逻辑:
import json # 读取输入的单个json文件 with open(snakemake.input.json_file, 'r') as f: data = json.load(f) # 调用REST API(替换为你的API调用逻辑) def call_api(data): # 此处编写你的API请求代码 return data enriched_data = call_api(data) # 写入输出文件 with open(snakemake.output.enriched_json, 'w') as f: json.dump(enriched_data, f)
运行方式
执行以下命令启动流程,通过--resources api_slot=1限制同一时间仅1个任务调用API:
snakemake --resources api_slot=1 -j 10 # -j可设更大值,但api_slot会强制串行
关键优势
- 断点续跑:每个任务对应独立输出文件,Snakemake会自动跳过已完成的任务,失败后仅重跑未完成的单个文件。
- 限流控制:通过自定义资源
api_slot灵活控制API调用并发数,完全匹配API的限流规则(如需允许2个并发,可设置api_slot=2并调整运行参数)。 - 资源优化:单个任务资源需求远低于批量任务,集群调度更高效,避免大任务占用资源过久。
内容的提问来源于stack exchange,提问作者el-xefe
相关产品推荐
相关产品推荐

