创建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) }
问题根源
- 通道关联错误:
Channel.fromFilePairs("${fastq_ch}/SRR*_{1,2}.fastq")是错误用法。fastq_ch是convert_to_fastq的输出通道,它的元素是流程生成的fastq_files目录路径,不能直接用字符串插值${fastq_ch}来拼接路径。这种写法会导致fromFilePairs在convert_to_fastq还未生成文件时就执行,找不到目标文件,流程因此停滞。 - 不必要的依赖:
indexing_reference的输入包含path fastp,但基因组索引构建只需要参考序列,和fastp的结果无关。这个多余的输入会让索引进程等待fastp_ch的输出,进一步加剧流程停滞问题。 - 输出定义不精准: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) }
关键修复说明
- 精准传递文件与ID:在convert_to_fastq中提取SRA编号,输出时将ID和对应的fastq文件绑定为元组,避免后续手动查找文件。
- 正确处理成对文件:使用
map操作将每个ID对应的两个fastq文件整理为有序列表,确保后续进程能正确获取R1/R2文件。 - 移除不必要的依赖:indexing_reference仅接收参考基因组输入,无需等待fastp结果,可独立并行执行,提升流程效率。
- 避免直接读取文件系统:不再使用
Channel.fromFilePairs从文件系统读取,而是通过流程通道直接传递文件,保证异步执行的正确性。 - 清理中间文件:在mapping_to_reference中添加中间文件清理,减少磁盘占用。
内容的提问来源于stack exchange,提问作者Javier Biotech
相关产品推荐
相关产品推荐

