Nextflow:搭建生信工作流程

Published on 2026-07-27 17:33
Author: yuyy
56 views

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 — 官方交互式教程