Nextflow:搭建生信工作流程
Nextflow 工作流框架使用指南
- Nextflow 是一款工作流管理框架,把一串分析步骤写成可复用、可并行、可断点续跑的流程。核心理念:只描述做什么,让框架管怎么跑。
- Nextflow 支持多种执行环境(本地、SLURM、SGE、PBS、K8s、云平台),通过配置文件切换执行器而无需修改流程代码。它还支持 Docker/Singularity 容器化,确保流程在不同环境下可复现。本文档系统介绍 Nextflow 的核心概念、语法结构、实践方法及初学者学习路径,并结合实际生信分析流程(土壤宏基因组 Streptomyces MAG 挖掘)中的代码示例进行讲解。
本文为初学者提供一条从零到能写流程的学习路线。在学习 Nextflow 之前,建议先掌握以下基础:
- Linux 基础命令:文件操作(ls、cd、cp、mv、rm)、文本处理(grep、wc、head、tail)、管道与重定向(|、>、2>&1)
- Shell 脚本基础:变量赋值(VAR=value)、条件判断(if)、循环(for)、函数定义
- 路径与环境变量:绝对路径与相对路径、PATH 环境变量、export、source
- 集群基本概念(如在 HPC 上使用):SLURM 调度(sbatch、squeue、sacct)、作业提交与查询、队列与 QOS
- Java 运行环境:Nextflow 依赖 Java 11+(推荐 Java 17 或 21),需确认 java -version 可用
- 不需要精通以上所有内容,但至少能看懂和修改 shell 脚本。
环境安装与配置
安装 Nextflow
方法一:官方一键安装(推荐新手)
# 下载安装脚本
curl -fsSL https://get.nextflow.io | bash
# 验证安装
./nextflow -version
# 移到 PATH 目录
mv nextflow ~/.local/bin/ # 或其他已在 PATH 中的目录
方法二:Conda 安装
conda install -c bioconda nextflow
创建配置文件创建配置文件
在项目目录下创建 nextflow.config:
// nextflow.config
params {
outdir = "./results"
}
process {
executor = "local"
cpus = 4
memory = "8 GB"
}
Nextflow核心能力
- 多样本自动并行:channel 放 N 个数据,自动跑 N 次
- 断点续跑:
-resume跳过已完成步骤 - 失败自动重试:
errorStrategy,可动态加资源 - 跨集群复用:改
config不动代码 - 自动日志追踪:每步耗时、状态、退出码
Nextflow核心概念
- process(进程):流水线上的一个"工位"。声明输入是什么、跑什么命令、输出是什么、要多少 CPU/内存。
- channel(通道):连接工位之间的"传送带"。承载的是数据,放 N 个数据则自动并行跑 N 次。
- workflow(工作流):"车间主任",把 process 按顺序连起来,只描述依赖关系。
Process:流水线的最小执行单元
Process 由三段组成:input(输入)→ shell(执行)→ output(输出)。
基本骨架:
process 名字 {
[directive] // 属性声明
input:
val x // 输入一个值
output:
stdout // 输出屏幕内容
shell:
"""
echo "!{x}" // 执行, !{} 插值
"""
}
实战示例:FASTP 质控 process
process FASTP {
tag { sample_id } // 任务显示名
publishDir { "${params.outdir}/${sample_id}/qc" } // 结果发布目录
input:
tuple val(sample_id), path(r1), path(r2)
output:
tuple val(sample_id), path("clean_R1.fastq.gz"),
path("clean_R2.fastq.gz"), emit: reads
path "fastp.json"
shell:
"""
fastp -i !{r1} -I !{r2} -o clean_R1.fastq.gz -O clean_R2.fastq.gz
"""
}
任何 process 都可拆成:属性 + 吃(input) + 干(shell) + 吐(output) 四部分。
input/output 类型详解
input 三种类型
- val:值类型,字符串/数字,不进 work 目录- path:文件类型,自动 stage 进 work 目录,跑完清理- tuple:元组,多个值组合,一次性传多个
val sample_id // 值类型
path r1 // 文件类型
tuple val(id), path(r1), path(r2) // 元组类型
output 三种类型:
- stdout:屏幕输出,命令的 stdout 变成 channel- path:文件,生成的文件作为输出- tuple+emit:元组 + 命名,多输出时用 .out.名字 精准引用
stdout // 屏幕输出
path "result.txt" // 文件输出
tuple val(id), path(f), emit: reads // 元组+命名
#关键区别:val 只传值不拷文件;path 自动 stage 文件;emit 给输出命名供下游精准引用。
Channel:数据管道
Channel 分为以下两类:
- queue channel(队列型):每个数据只被消费一次,样本数据典型场景,默认类型。
- value channel(广播型):可被多个 process 共享同一份数据,典型场景如参考数据库。
queue: 样本数据, 依次交由Process处理。
ch_samples = channel.of("SRR1","SRR2","SRR3")
value: 多 process 共享的数据,例如参考数据库,。
ch_db = Channel.value("/home/.../PlusPF")
KRAKEN2(ch_samples.combine(ch_db))
Channel的三种创建方式
Channel.of— 直接列值,测试或小规模场景Channel.fromPath— 通配符文件匹配,处理单端数据Channel.fromFilePairs— 双端自动配对,生成 (id, R1, R2) 元组
ch1 = channel.of("SRR1", "SRR2")
ch2 = Channel.fromPath("/data/*.fq.gz")
ch3 = Channel.fromFilePairs("/data/*_{1,2}.fq.gz")
自动并行机制:把样本放进 channel,Nextflow 自动并行启动,无需手动管理循环。
workflow {
ch = channel.of("SRR1", "SRR2", "SRR3")
DOWNLOAD(ch) // 自动启动 3 个独立任务, 同时跑
}
Workflow:编排执行顺序
Workflow 用 channel 把 process 串成流水线,只描述依赖关系,Nextflow 据此构建 DAG 调度执行。
workflow {
ch = channel.of("SRR1", "SRR2")
DOWNLOAD(ch) // 下载
FASTP(DOWNLOAD.out) // .out 接力
MEGAHIT(FASTP.out.reads) // .reads 是 emit 命名
}
.out 拿上游输出传给下游 —— 这就是接力机制。
实战示例:读样本表 → 造 channel → 接力串联
workflow {
Channel.fromPath(params.samples)
.splitCsv(header: true, sep: '\t')
.map { row -> tuple(row.sample_id, ...) }
.set { raw_ch }
DOWNLOAD(raw_ch)
FASTP(DOWNLOAD.out)
MEGAHIT(FASTP.out.reads)
BOWTIE2_MAP(MEGAHIT.out.join(FASTP.out.reads))
}
.splitCsv(header: true, sep: '\t'):按 tab 解析带表头的 TSV,每行变成字典.map:字典变成元组row.sample_id取列值。.set:命名
Directive:进程属性声明(六大类)
1. 标识类 — 给任务起名,方便在日志和监控中识别:
tag { sample_id } // 任务显示名(动态)
label 'high_mem' // 分组标签(静态)
2. 发布类2. 发布类 — 把 output 文件从 work/ 发布到人类可读目录:
publishDir { "${params.outdir}/${sample_id}/qc" }, mode: 'copy'
mode 可选:copy(复制)/ symlink(软链接)/ move(移动)。
3. 资源类 — 声明 CPU/内存/时间需求:
cpus 24
memory '120 GB'
time '72h'
disk '500 GB'
4. 执行类 — 指定执行器和队列:
executor 'slurm'
queue 'cu'
clusterOptions "--qos=normal --exclusive"
5. 错误处理类 — 失败时的行为策略:
errorStrategy 'retry' // retry/terminate/finish/ignore
maxRetries 2 // 最多重试次数
四种策略:retry 重试;terminate 立即终止(默认);finish 等已跑任务结束;ignore 忽略继续。
6. 环境类 — 声明软件来源:
container 'quay.io/biocontainers/fastp:0.23'
module 'metaSoil'
params:参数化配置
路径/资源不写死,更换计算硬件只改参数不动代码。可选参数:给默认值。
params.outdir = '/home/.../Streptomycesmeta'
// 没传命令行就用默认值
必填参数:null 占位
params.samples = null
if (params.samples) { ... } else { error "..." }
nextflow.config 是流程的总控制面板,放在 main.nf 同目录下。所有 process 的默认资源、执行器、参数默认值都在这里统一定义。
// nextflow.config —— 放 main.nf 同目录
// ===== 1. 全局参数 =====
params {
samples = null // 必填:样本列表 TSV
outdir = '/home/user/myproject'
mod_metasoil = 'metaSoil'
kraken2_db = '/home/mselab/biodbs/kraken2_db/PlusPF'
bakta_db = '/home/mselab/biodbs/bakta_db/db-full/db'
}
// ===== 2. process 默认设置(所有 process 都生效)=====
process {
executor = 'slurm' // 提交到 SLURM 集群
queue = 'normal' // 队列名
errorStrategy = { task.exitStatus in 137..140 ? 'retry' : 'terminate' }
maxRetries = 1
}
// ===== 3. 按 label 分组设置 =====
process {
withLabel: 'big_mem' {
cpus = 24
memory = '80 GB'
time = '48h'
}
withLabel: 'small' {
cpus = 4
memory = '8 GB'
}
}
// ===== 4. 按 process 名称精准设置(优先级最高)=====
process {
withName: 'MEGAHIT' {
cpus = 28
memory = '200 GB'
time = '72h'
clusterOptions = '--exclusive'
}
withName: 'FASTP' {
cpus = 8
memory = '16 GB'
}
}
- 配置优先级(从高到低)
- withName 覆盖 withLabel,withLabel 覆盖 process 全局默认
- process 代码内部写的 directive 优先级最低
- 命令行覆盖:
--outdir /other
插值机制
shell 块里 !{x} 是 Nextflow 变量(自动填),${} 留给 bash 变量(不冲突)。
shell:
"""
source !{projectDir}/bin/init_modules.sh # !{} NF变量
module load !{params.mod_metasoil} # !{} NF变量
t=!{task.cpus}; [ "$t" -gt 16 ] && t=16 # !{} NF; $t bash
fastp -i !{r1} -I !{r2} # !{} NF变量
"""
错误处理与断点续跑
errorStrategy + 动态资源:失败时自动重试,重试时翻倍资源。
process MEGAHIT {
errorStrategy 'retry'
maxRetries 2
memory { 8.GB * task.attempt } // 8G -> 16G -> 24G
cpus { 4 * task.attempt }
}
注意:shell 里 || true 会绕过 errorStrategy(吞错误),慎用。
断点续跑 -resume:利用 work/ 缓存,按输入哈希记忆,改了参数只重跑受影响步骤。
nextflow run main.nf --samples x.tsv -resume
work/ 删了无法 resume。
常用内置变量
- params.xxx:用户自定义参数,如 params.outdir
- projectDir:项目根目录(main.nf 所在目录),如 !{projectDir}/bin/init.sh
- workDir:work 缓存目录,即 work/
- launchDir:启动 nextflow 命令的目录
- task.cpus:当前 process 分配的 CPU 数,如 !{task.cpus}
- task.attempt:当前重试次数(第几次跑),如 memory { 8.GB * task.attempt }
- task.index:当前任务在 channel 中的序号
常用 Channel 操作函数
- .map { }:对每个元素做变换,如 .map { row -> tuple(row.id) }
- .filter { }:过滤元素,如 .filter { it.size() > 0 }
- .splitCsv():解析 CSV/TSV 文件,如 .splitCsv(header: true, sep: '\t')
- .set { name }:给 channel 命名,如 .set { raw_ch }
- .combine(ch):与另一个 channel 做笛卡尔积,如 ch_samples.combine(ch_db)
- .join(ch):按 key 做内连接,如 MEGAHIT.out.join(FASTP.out.reads)
- .collect():把所有元素收集成一个列表
- .view():打印 channel 内容(调试用)
进阶实践建议
- 模块化设计:把不同分析步骤的 process 写在不同 .nf 文件中,用 include 导入
- config 分离:把不同集群的配置写成独立的 conf/xxx.config,用 -profile 切换
- 容器化:用 Docker/Singularity 容器封装软件依赖,保证流程可移植性
- 测试驱动:先用小样本(1-2 个)测试流程,再扩大规模
- 日志管理:善用 .nextflow.log 和 work/ 下的 .command.sh、.command.err 排查问题
- 资源监控:用 trace 或 timeline 报告分析每步耗时和资源使用
- 学习 nf-core:参考社区流程 nf-core 的代码结构和最佳实践
适用场景
Nextflow 适用于任何需要将多个分析步骤串联、并行化、可复现的场景,尤其在生物信息学领域应用最为广泛:
- 宏基因组分析流程:质控 → 组装 → 分箱 → 物种注释 → 功能注释
- 转录组分析流程:质控 → 比对 → 定量 → 差异表达分析
- 变异检测流程:质控 → 比对 → 去重 → 变异 calling → 注释
- 基因组组装流程:质控 → 组装 → 抛光 → 评估
- 批量数据处理:数百个样本的标准化处理与质量控制
- 可重复研究:将分析流程打包,确保他人能用相同输入获得相同结果
- 不适合的场景:单步简单分析、交互式探索性分析(这些用 Jupyter/RStudio 更合适)。
社区资源
- 官方文档:https://www.nextflow.io/docs/latest/ — 最权威的参考,包含所有语法和 directive 说明
- GitHub 仓库:https://github.com/nextflow-io/nextflow — 源码、Issue 反馈
- nf-core 社区:https://nf-co.re — 社区共建的标准化流程集合,适合学习最佳实践
- Slack 社区:Nextflow 官方 Slack 频道,可提问和交流
- 培训资源:https://training.nextflow.io — 官方交互式教程