You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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])

关键说明

  1. Checkpoint的作用:解决动态生成文件的依赖问题,让Snakemake在运行时自动发现拆分出的所有分片。
  2. 并行调度:Snakemake会自动并行处理多个分片的cmd2任务,提升整体运行效率。
  3. 顺序保障:合并时用ls -v按分片数字排序,确保合并后的文件和原大文件的内容顺序一致。

这个方案完全基于Snakemake原生能力,不需要额外嵌套Python流程,能完美适配你的需求。


内容的提问来源于stack exchange,提问作者Georg B

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 04:18:20