Un paso más desde F27 — Del flujo de archivos al flujo de datos
Snakemake en F27 infería el orden del pipeline a partir de los patrones de nombres de entrada y salida de los archivos. Este enfoque es potente, pero en proyectos a gran escala donde las muestras aumentan de cientos a miles y la ejecución se distribuye en clústeres o la nube, la propia dependencia del sistema de archivos puede convertirse en un cuello de botella. Nextflow (2013~, Seqera·EMBL-EBI) con su DSL2 cambia el enfoque: diseña los canales (channel), por los cuales fluyen los datos entre procesos, como objetos de primera clase, en lugar de centrarse en los archivos.
Principios — Construcción del flujo de datos mediante canales y procesos
process — Unidad de cálculo independiente
El process de Nextflow define una unidad de cálculo similar al rule de Snakemake, pero intercambia entradas y salidas a través de canales en lugar de rutas del sistema de archivos.
// Ejemplo de main.nf (DSL2)
process ALIGN {
input:
tuple val(sample_id), path(r1), path(r2)
path ref
output:
tuple val(sample_id), path("${sample_id}.bam")
script:
"""
bwa mem ${ref} ${r1} ${r2} | samtools sort -o ${sample_id}.bam
"""
}
process MARK_DUPLICATES {
input:
tuple val(sample_id), path(bam)
output:
tuple val(sample_id), path("${sample_id}.dedup.bam")
script:
"""
gatk MarkDuplicates -I ${bam} -O ${sample_id}.dedup.bam -M metrics.txt
"""
}
workflow {
reads_ch = Channel.fromFilePairs("reads/*_R{1,2}.fastq.gz")
// Para reutilizar la misma referencia en múltiples lecturas, pase un valor en lugar de un canal de cola.
ref = file("ref/genome.fa")
ALIGN(reads_ch, ref)
MARK_DUPLICATES(ALIGN.out)
}Canal — Flujo de datos asíncrono
Un canal es una cola asíncrona que transporta rutas de archivo, valores o tuplas entre procesos. Cuando un proceso emite resultados a un canal, el siguiente proceso que toma ese canal como entrada puede comenzar la ejecución tan pronto como llegan los datos (sin necesidad de esperar a que termine todo un lote).
Si se procesan 100 muestras, Nextflow crea y ejecuta instancias de proceso independientes para cada muestra a medida que llega al canal. Esta es la forma natural de obtener paralelismo por muestra sin escribir explícitamente un bucle que diga "repite 100 veces".
Ejemplo de cálculo manual: paralelismo de canales y latencia
Supongamos que las 100 muestras llegan al canal de entrada en momentos distintos (debido a diferencias en la velocidad de transmisión). Si un solo proceso de alineación (ALIGN) tarda en promedio 10 minutos y el clúster permite ejecutar hasta 20 instancias en paralelo, el tiempo total de procesamiento es aproximadamente
Si la ejecución fuera secuencial (en serie), sería
Gracias al paralelismo que otorga naturalmente el flujo de datos basado en canales, se obtiene una mejora de velocidad de procesamiento de aproximadamente 20 veces sin tener que escribir manualmente la lógica de paralelización. Por supuesto, la ganancia real depende de los recursos disponibles del clúster y de los cuellos de botella de E/S en cada etapa.
Diferencias entre DSL1 y DSL2 — Modularidad
El Nextflow inicial (DSL1) describía proceduralmente todo el pipeline en un solo script. DSL2 introduce la estructura de separar process en módulos (module) reutilizables y combinarlos dentro del bloque workflow. Esta modularidad es el requisito previo que hace posible el ecosistema nf-core, que se tratará en F29 (una estructura donde cientos de pipelines reutilizan módulos comunes).
Práctica: configurar un canal de pares de archivos y verificar con dry-run
# SageMaker Studio Lab. Nextflow se ejecuta sobre el entorno de tiempo de ejecución de Java.curl -s https://get.nextflow.io | bash./nextflow -v
# Use la opción -preview para verificar solo la estructura del flujo de trabajo sin ejecutarlo realmente../nextflow run main.nf -preview
# Ejecución real. -resume comparte el mismo principio que la optimización de reejecución de Snakemake en F27,# reutilizando los resultados de ejecuciones anteriores en caché para continuar desde el punto de fallo../nextflow run main.nf -resume-resume La bandera almacena en caché el valor hash de cada ejecución de proceso y omite la reejecución de los procesos cuya entrada no ha cambiado. Aunque comparte el mismo objetivo que la optimización de reejecución basada en marcas de tiempo de Snakemake descrita en F27, su implementación se basa en hashes de contenido.
Mapeo de CS
- Programación de flujo de datos (dataflow programming): Representar los cálculos como nodos (procesos) y los datos que fluyen entre ellos (canales) constituye la definición misma del paradigma de programación de flujo de datos.
- Flujos reactivos (reactive streams): El mecanismo por el cual la llegada de datos a un canal desencadena la ejecución de la siguiente etapa sigue la misma lógica que el patrón Observable de bibliotecas de programación reactiva como RxJava y RxJS.
- Caché basado en contenido: La reutilización basada en hashes de
-resumepertenece a la misma línea estratégica que la utilizada por los sistemas de compilación de software (como Bazel) para determinar si hay un acierto en la caché de construcción mediante el hash del contenido de entrada.
Errores frecuentes
- Mezcla de sintaxis DSL1 y DSL2: Dado que muchos tutoriales y ejemplos antiguos están escritos con la sintaxis DSL1, pegarlos directamente en proyectos modernos DSL2 provoca errores. Es fundamental verificar siempre la versión de Nextflow y la configuración de
nextflow.enable.dsl=2. - Pasar varios canales queue a la entrada de un proceso: Varios canales queue pueden emparejarse de manera no determinista según el orden de llegada de los valores. Si un proceso requiere múltiples entradas, se debe utilizar solo un canal queue para los datos variables, mientras que las referencias comunes o los valores de configuración deben pasarse mediante
Channel.value()o como valores normales. En DSL2, cuando un canal es suscrito por varios procesos, se ramifica automáticamente; no confundir esto con la antigua regla de "consumo único".
Para profundizar más
El texto ha sido reescrito directamente por el equipo de investigación de BPD. Para una comprensión más profunda, consulte el artículo original y la documentación oficial.
- Artículo original de Nextflow: Di Tommaso et al. (2017), Nextflow enables reproducible computational workflows, Nature Biotechnology 35.
- Documentación oficial de Nextflow: Guía de migración a DSL2 y referencia de operadores de canales.
- Wellcome Sanger — Material de formación práctica de Nextflow: Cursos de formación oficiales.
En el próximo capítulo (F29), realizaremos un recorrido práctico por el ecosistema estándar de pipelines, nf-core, construido sobre esta modularización DSL2.