Volver a la lista

Nextflow DSL2: Diseñar el flujo de datos en sí, no los archivos

Si el DAG de Snakemake se basaba en las relaciones entre archivos, Nextflow con DSL2 trata los propios flujos de datos llamados canales como objetos de primera clase. Se abordan la sintaxis de procesos y canales junto con ejemplos prácticos.

Intermedio
|
22min
|
Verificado (2026-07-29)
Nextflowworkflow managerdataflownf-core
Progreso0/120 (0%)

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.

groovy
// 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).

Snakemake:archivo existeactivacioˊn de reglaNextflow:llegada de datos al canalcreacioˊn de instancia de proceso\text{Snakemake}: \text{archivo existe} \to \text{activación de regla} \qquad \text{Nextflow}: \text{llegada de datos al canal} \to \text{creación de instancia de proceso}

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

10020×10min=5×10=50min\left\lceil \frac{100}{20} \right\rceil \times 10\text{min} = 5 \times 10 = 50\text{min}

Si la ejecución fuera secuencial (en serie), sería

100×10min=1,000min16.7horas100 \times 10\text{min} = 1{,}000\text{min} \approx 16.7\text{horas}

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

bash
# 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 -resume pertenece 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.

💬 Preguntas y comentarios

0 comentarios

Puedes publicar sin iniciar sesión. Los comentarios de invitados no pueden editarse ni eliminarse después.

0/2000

Cargando...