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

创建Channel时Nextflow Workflow突然终止问题求助

Nextflow流程停滞问题排查与修复

我用Nextflow写了一个从SRA下载、转fastq、质控、比对的流程,流程结构是:

  • descargaSRA:读取SRA编号文件,下载对应的.sra文件
  • convert_to_fastq:将.sra转为成对的fastq文件(SRR*_1.fastq/SRR*_2.fastq)
  • filter_reads:用fastp质控成对reads
  • indexing_reference:构建参考基因组索引
  • mapping_to_reference:将质控后reads比对到参考基因组

但运行时,convert_to_fastq进程结束后流程就停滞了,必须重新运行才能得到后续结果。我尝试用fromFilePairs将fastq文件按ID配对,但推测是通道关联错误导致的问题。

原始流程代码

process descargaSRA {
    publishDir("sra_files", mode: 'copy')

    input:
        path sra_file

    output:
        file "SRR*"

    script:
        """
        ids_sra=$(cat ${sra_file})
        for id in ${ids_sra}
        do
            prefetch ${id} -O .
            mv "$id/$id.sra" .
            rm -r "$id"
        done
        """
}

process convert_to_fastq {
    input:
        file sra_content 

    output:
        path "fastq_files"

    script:
        """
        fastq-dump --split-files ${sra_content} -O fastq_files
        """
}

process filter_reads {
    publishDir("fastp_files", mode: 'copy')

    input:
        tuple val(id), val(fastq_files)
    output:
        tuple val(id), path("clean_*")

    script:
        """
        echo "Processing ${id} reads"

        fastp -i ${fastq_files[0]} -I ${fastq_files[1]} -o clean_${id}_1.fastq -O clean_${id}_2.fastq --cut_tail --cut_window_size 5 --cut_mean_quality 30
        """
}

process indexing_reference {
    publishDir("index_reference", mode: "copy")
    input:
        path reference
        path fastp

    output:
        path "genome*"

    script:
        """
        bwa index ${reference} -p genome
        """
}

process mapping_to_reference{
    publishDir("bam_files",mode:"copy")
    input:
        path genome
        tuple val(id), val(clean_fastq)
    output:
        path "*.sorted.bam"
    script:
        """
        bwa mem genome ${clean_fastq[0]} ${clean_fastq[1]} > ${id}.sam
        samtools view ${id}.sam -o ${id}.bam
        samtools sort ${id}.bam -o ${id}_sorted.bam
        samtools index ${id}_sorted.bam
        """
}

workflow {
    input_ch = Channel.fromPath(params.sra_file)
    reference = Channel.fromPath(params.reference)

    sra_ch = descargaSRA(input_ch)

    fastq_ch = convert_to_fastq(sra_ch)

    "This channel will group the fastq files by ID"
    fastq_paths = Channel.fromFilePairs("${fastq_ch}/SRR*_{1,2}.fastq")

    fastp_ch = filter_reads(fastq_paths)
    genome = indexing_reference(reference,fastp_ch)
    map_ch = mapping_to_reference(genome,fastp_ch)
}

问题根源

  1. 通道关联错误:Channel.fromFilePairs("${fastq_ch}/SRR*_{1,2}.fastq") 是错误用法。fastq_ch是convert_to_fastq的输出通道,它的元素是流程生成的fastq_files目录路径,不能直接用字符串插值${fastq_ch}来拼接路径。这种写法会导致fromFilePairs在convert_to_fastq还未生成文件时就执行,找不到目标文件,流程因此停滞。
  2. 不必要的依赖:indexing_reference的输入包含path fastp,但基因组索引构建只需要参考序列,和fastp的结果无关。这个多余的输入会让索引进程等待fastp_ch的输出,进一步加剧流程停滞问题。
  3. 输出定义不精准:convert_to_fastq的输出是整个fastq_files目录,但后续需要的是目录内的成对fastq文件,没有直接传递文件路径导致后续通道无法正确捕获文件。

修复后的流程代码

process descargaSRA {
    publishDir("sra_files", mode: 'copy')

    input:
        path sra_file

    output:
        file "*.sra" emit: sra_files

    script:
        """
        ids_sra=$(cat ${sra_file})
        for id in ${ids_sra}
        do
            prefetch ${id} -O .
            mv "$id/$id.sra" .
            rm -r "$id"
        done
        """
}

process convert_to_fastq {
    input:
        path sra_file

    output:
        tuple val(srr_id), path("fastq_files/*.fastq") emit: fastq_pairs

    script:
        """
        # 提取SRA编号
        srr_id=$(basename ${sra_file} .sra)
        fastq-dump --split-files ${sra_file} -O fastq_files
        """
}

process filter_reads {
    publishDir("fastp_files", mode: 'copy')

    input:
        tuple val(id), path(fastq_files)
    output:
        tuple val(id), path("clean_${id}_*.fastq") emit: cleaned_reads

    script:
        """
        echo "Processing ${id} reads"
        # 拆分成对的fastq文件
        R1=$(echo ${fastq_files} | tr ' ' '\n' | grep "_1.fastq")
        R2=$(echo ${fastq_files} | tr ' ' '\n' | grep "_2.fastq")
        
        fastp -i ${R1} -I ${R2} -o clean_${id}_1.fastq -O clean_${id}_2.fastq --cut_tail --cut_window_size 5 --cut_mean_quality 30
        """
}

process indexing_reference {
    publishDir("index_reference", mode: "copy")
    input:
        path reference

    output:
        path "genome*" emit: genome_index

    script:
        """
        bwa index ${reference} -p genome
        """
}

process mapping_to_reference{
    publishDir("bam_files",mode:"copy")
    input:
        path genome_index
        tuple val(id), path(clean_fastq)
    output:
        path "${id}_sorted.bam*" emit: mapped_bams

    script:
        """
        # 拆分质控后的成对文件
        R1=$(echo ${clean_fastq} | tr ' ' '\n' | grep "_1.fastq")
        R2=$(echo ${clean_fastq} | tr ' ' '\n' | grep "_2.fastq")
        
        bwa mem ${genome_index} ${R1} ${R2} > ${id}.sam
        samtools view -bS ${id}.sam -o ${id}.bam
        samtools sort ${id}.bam -o ${id}_sorted.bam
        samtools index ${id}_sorted.bam
        # 清理中间文件
        rm ${id}.sam ${id}.bam
        """
}

workflow {
    input_ch = Channel.fromPath(params.sra_file)
    reference_ch = Channel.fromPath(params.reference)

    # 下载SRA文件
    sra_files_ch = descargaSRA(input_ch)

    # 转换为fastq并按ID配对
    fastq_pairs_ch = convert_to_fastq(sra_files_ch)
                     // 将每个ID对应的两个fastq文件整理为元组
                     .map { id, files -> 
                         def paired_files = files.groupBy { it.name.contains("_1.fastq") ? "R1" : "R2" }
                         tuple(id, [paired_files.R1[0], paired_files.R2[0]]) 
                     }

    # 质控
    cleaned_reads_ch = filter_reads(fastq_pairs_ch)
                       // 同样整理质控后的成对文件
                       .map { id, files -> 
                           def paired_files = files.groupBy { it.name.contains("_1.fastq") ? "R1" : "R2" }
                           tuple(id, [paired_files.R1[0], paired_files.R2[0]]) 
                       }

    # 构建基因组索引(独立执行,不依赖其他进程)
    genome_index_ch = indexing_reference(reference_ch)

    # 比对
    mapped_bams_ch = mapping_to_reference(genome_index_ch, cleaned_reads_ch)
}

关键修复说明

  1. 精准传递文件与ID:在convert_to_fastq中提取SRA编号,输出时将ID和对应的fastq文件绑定为元组,避免后续手动查找文件。
  2. 正确处理成对文件:使用map操作将每个ID对应的两个fastq文件整理为有序列表,确保后续进程能正确获取R1/R2文件。
  3. 移除不必要的依赖:indexing_reference仅接收参考基因组输入,无需等待fastp结果,可独立并行执行,提升流程效率。
  4. 避免直接读取文件系统:不再使用Channel.fromFilePairs从文件系统读取,而是通过流程通道直接传递文件,保证异步执行的正确性。
  5. 清理中间文件:在mapping_to_reference中添加中间文件清理,减少磁盘占用。

内容的提问来源于stack exchange,提问作者Javier Biotech

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 08:32:04