如何在Snakemake中处理可变数量的输入输出文件?
问题描述
我有一个基于Snakemake的简易流程,其中cmd2工具处理大文件时会出现内存溢出问题。我打算将文件拆分为1000行的片段依次执行,再合并结果继续流程,但尝试用dynamic/expand实现未果。现有流程代码(已修正语法错误)如下:
rule all: input: "results/{motif}--{file}.score" rule cmd1: input: "data/{file}.txt" output: "intermediate/{file}.fasta" shell: "bash cmd1.sh {input} > {output}" rule cmd2: input: "intermediate/{file}.fasta" output: "intermediate/{file}--{motif}.fasta" shell: "cmd2 {input} > {output}" rule cmd3: input: motif = "data/{motif}.mf", seq = "intermediate/{file}--{motif}.fasta" output: "results/{file}--{motif}.score" run: some_function(motif, seq, output[0])
请问如何编写规则实现拆分、分片执行cmd2及合并结果?该需求能否通过Snakemake合理实现,还是应改用Python函数实现嵌套流程?
解决方案
这个需求完全可以用Snakemake原生特性实现,不需要改用嵌套Python流程。核心思路是新增拆分规则、分片处理cmd2规则和合并规则,替换原有的cmd2规则,具体实现如下:
步骤1:用Checkpoint实现动态文件拆分
因为拆分前无法预知大文件会生成多少个分片,用Snakemake的checkpoint可以在运行时动态发现这些分片文件:
import os from glob import glob checkpoint split_fasta: input: "intermediate/{file}.fasta" output: directory("intermediate/{file}_parts") shell: # 创建分片目录,按1000行拆分大文件,输出命名为file_part_1.fasta、file_part_2.fasta... "mkdir -p {output} && split -l 1000 {input} {output}/{file}_part_ --additional-suffix .fasta" def get_split_parts(wildcards): # 从checkpoint输出目录中获取所有分片文件路径 return glob(os.path.join("intermediate", wildcards.file + "_parts", f"{wildcards.file}_part_*.fasta"))
步骤2:分片执行cmd2
针对每个小分片单独运行cmd2,避免内存溢出:
rule cmd2_part: input: fasta_part = "intermediate/{file}_parts/{file}_part_{part}.fasta" output: "intermediate/{file}_parts/{file}_part_{part}--{motif}.fasta" shell: "cmd2 {input.fasta_part} > {output}"
步骤3:合并分片处理结果
将所有cmd2处理后的分片按顺序合并,还原成原流程需要的文件格式:
rule merge_cmd2_results: input: parts = dynamic(get_split_parts), motif = "{motif}" output: "intermediate/{file}--{motif}.fasta" shell: # 按分片数字顺序合并,避免结果错乱 "ls -v {parts} | xargs cat > {output}"
步骤4:调整完整流程依赖
将原流程的依赖指向合并后的文件,完整流程代码如下:
import os from glob import glob # 替换为实际的motif和file列表 rule all: input: expand("results/{motif}--{file}.score", motif=["motifA", "motifB"], file=["sampleX", "sampleY"]) rule cmd1: input: "data/{file}.txt" output: "intermediate/{file}.fasta" shell: "bash cmd1.sh {input} > {output}" checkpoint split_fasta: input: "intermediate/{file}.fasta" output: directory("intermediate/{file}_parts") shell: "mkdir -p {output} && split -l 1000 {input} {output}/{file}_part_ --additional-suffix .fasta" def get_split_parts(wildcards): return glob(os.path.join("intermediate", wildcards.file + "_parts", f"{wildcards.file}_part_*.fasta")) rule cmd2_part: input: fasta_part = "intermediate/{file}_parts/{file}_part_{part}.fasta" output: "intermediate/{file}_parts/{file}_part_{part}--{motif}.fasta" shell: "cmd2 {input.fasta_part} > {output}" rule merge_cmd2_results: input: parts = dynamic(get_split_parts), motif = "{motif}" output: "intermediate/{file}--{motif}.fasta" shell: "ls -v {parts} | xargs cat > {output}" rule cmd3: input: motif = "data/{motif}.mf", seq = "intermediate/{file}--{motif}.fasta" output: "results/{file}--{motif}.score" run: some_function(input.motif, input.seq, output[0])
关键说明
- Checkpoint的作用:解决动态生成文件的依赖问题,让Snakemake在运行时自动发现拆分出的所有分片。
- 并行调度:Snakemake会自动并行处理多个分片的cmd2任务,提升整体运行效率。
- 顺序保障:合并时用
ls -v按分片数字排序,确保合并后的文件和原大文件的内容顺序一致。
这个方案完全基于Snakemake原生能力,不需要额外嵌套Python流程,能完美适配你的需求。
内容的提问来源于stack exchange,提问作者Georg B
相关产品推荐
相关产品推荐

