Procesamiento por lotes de FASTQ en la nube: de la computadora local a AWS
Al finalizar este tema
Podrá crear un flujo de trabajo que encapsule sus scripts de Python locales en contenedores para ejecutarlos de forma masiva en AWS Batch, integrando los conceptos de despliegue en AWS y almacenamiento en S3 aprendidos. Aunque este ejemplo utiliza archivos FASTQ, este patrón es aplicable a cualquier tipo de procesamiento de datos a gran escala.
Este artículo es un ejemplo didáctico. Los flujos de trabajo de secuenciación reales en producción utilizan lenguajes de flujo de trabajo como Nextflow, Snakemake o WDL.
"El disco de mi computadora está lleno": las limitaciones del procesamiento local
Supongamos que ha recibido un archivo FASTQ de 200 GB desde un centro de secuenciación. Planea realizar las siguientes tareas con este archivo:
- Recorte de calidad (eliminación de lecturas de baja calidad)
- Alineamiento con el genoma de referencia
- Llamada de variantes (variant calling)
Enfoque local:
import subprocess
subprocess.run(["fastp", "-i", "sample.fastq.gz", "-o", "trimmed.fastq.gz"])subprocess.run(["bwa", "mem", "reference.fa", "trimmed.fastq.gz", "-o", "aligned.sam"])subprocess.run(["samtools", "sort", "aligned.sam", "-o", "sorted.bam"])subprocess.run(["bcftools", "call", "sorted.bam", "-o", "variants.vcf"])Este enfoque presenta tres problemas prácticos.
Problema 1: Disco. El archivo original de 200 GB + varios archivos intermedios = más de 1 TB. El disco de su computadora portátil no tiene capacidad suficiente.
Problema 2: Memoria y tiempo. El alineamiento con BWA requiere 32 GB de memoria y varias horas. Su computadora portátil no podrá realizar otras tareas durante este proceso.
Problema 3: Escalabilidad. Esta vez es una sola muestra, pero el próximo proyecto tendrá 100. El procesamiento secuencial local tardaría semanas.
El enfoque real consiste en encapsular este pipeline en un contenedor y ejecutarlo en AWS Batch. Los archivos originales se mantienen en S3, y el job de Batch levanta el contenedor en una instancia EC2 de gran tamaño para procesar los archivos y luego guarda los resultados nuevamente en S3. Las 100 muestras se ejecutan como 100 jobs paralelos.
De caja negra a componentes: abriendo el pipeline en la nube
Hay cuatro componentes clave.
Componente 1: S3 — Almacenamiento ilimitado
S3 es un almacenamiento de objetos con capacidad prácticamente ilimitada. Un archivo puede pesar hasta 5 TB y no hay límite de archivos por bucket. Mediante políticas de ciclo de vida, se pueden mover automáticamente los archivos antiguos a una capa de bajo costo (Glacier).
Se accede mediante el SDK de Python (boto3).
import boto3
s3 = boto3.client("s3")
s3.upload_file("sample.fastq.gz", "my-bucket", "raw/sample.fastq.gz")
s3.download_file("my-bucket", "results/variants.vcf", "variants.vcf")
for obj in s3.list_objects_v2(Bucket="my-bucket", Prefix="raw/")["Contents"]: print(obj["Key"], obj["Size"])Atención: Los costos de S3 incluyen almacenamiento + transferencia + solicitudes. Si se accede desde una instancia EC2 en la misma región, no hay costos de transferencia. Ejecutar el proceso en la misma región donde se encuentran los archivos es el primer principio para optimizar costos.
Componente 2: Contenedor Docker
Agrupa todo lo que requiere su script de pipeline — Python + fastp + bwa + samtools + bcftools + su propio código — en una sola imagen.
Dockerfile Ejemplo:
FROM ubuntu:22.04
RUN apt-get update && apt-get install -y \
python3.12 python3-pip \
fastp bwa samtools bcftools \
&& rm -rf /var/lib/apt/lists/*
COPY requirements.txt /app/
RUN pip3 install -r /app/requirements.txt
COPY pipeline.py /app/
WORKDIR /app
ENTRYPOINT ["python3", "pipeline.py"]Compilar y subir a ECR (Elastic Container Registry):
docker build -t my-fastq-pipeline .aws ecr get-login-password | docker login --username AWS --password-stdin <account>.dkr.ecr.<region>.amazonaws.comdocker tag my-fastq-pipeline:latest <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latestdocker push <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latestComponent 3: AWS Batch Job Definition
Batch consists of three layers.
Compute environment: The layer that manages the actual EC2 instances. It specifies the minimum vCPU, maximum vCPU, and instance type.
Job queue: The queue where jobs wait. Priority can be specified.
Job definition: Defines which container image to run, with which resources, and how.
import boto3
batch = boto3.client("batch")
batch.register_job_definition( jobDefinitionName="fastq-pipeline", type="container", containerProperties={ "image": "<account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latest", "vcpus": 8, "memory": 32000, "command": [ "--input", "Ref::input_s3", "--output", "Ref::output_s3" ], "jobRoleArn": "arn:aws:iam::<account>:role/BatchJobRole" })Ref::input_s3 es un parámetro que se completa al enviar un trabajo.
Componente 4: Envío de trabajos
Envía un trabajo para una muestra.
def submit_sample_job(sample_id: str, input_s3: str, output_s3: str) -> str: response = batch.submit_job( jobName=f"fastq-{sample_id}", jobQueue="my-fastq-queue", jobDefinition="fastq-pipeline", parameters={ "input_s3": input_s3, "output_s3": output_s3 } ) return response["jobId"]
job_id = submit_sample_job( sample_id="sample01", input_s3="s3://my-bucket/raw/sample01.fastq.gz", output_s3="s3://my-bucket/results/sample01/")print(f"Submitted job: {job_id}")Uniendo las cuatro partes: la lógica interna del contenedor
Así es como se ve pipeline.py ejecutándose dentro del contenedor.
import argparseimport subprocessfrom pathlib import Pathimport boto3
def parse_s3_uri(uri: str) -> tuple[str, str]: parts = uri.replace("s3://", "").split("/", 1) return parts[0], parts[1] if len(parts) > 1 else ""
def download_from_s3(s3_uri: str, local_path: Path) -> None: bucket, key = parse_s3_uri(s3_uri) s3 = boto3.client("s3") s3.download_file(bucket, key, str(local_path))
def upload_to_s3(local_path: Path, s3_uri: str) -> None: bucket, key = parse_s3_uri(s3_uri) s3 = boto3.client("s3") s3.upload_file(str(local_path), bucket, key)
def run_pipeline(input_s3: str, output_s3: str, work_dir: Path) -> None: work_dir.mkdir(exist_ok=True)
raw = work_dir / "sample.fastq.gz" trimmed = work_dir / "trimmed.fastq.gz" aligned = work_dir / "aligned.bam" sorted_bam = work_dir / "sorted.bam" variants = work_dir / "variants.vcf"
print(f"[1/5] Downloading {input_s3}") download_from_s3(input_s3, raw)
print("[2/5] Quality trimming with fastp") subprocess.run(["fastp", "-i", str(raw), "-o", str(trimmed)], check=True)
print("[3/5] Alignment with BWA") with open(aligned, "w") as f: subprocess.run( ["bwa", "mem", "/reference/hg38.fa", str(trimmed)], stdout=f, check=True )
print("[4/5] Sort BAM") subprocess.run( ["samtools", "sort", str(aligned), "-o", str(sorted_bam)], check=True )
print("[5/5] Variant calling") subprocess.run( ["bcftools", "call", "-o", str(variants), str(sorted_bam)], check=True )
print(f"Uploading {variants} to {output_s3}variants.vcf") upload_to_s3(variants, f"{output_s3}variants.vcf")
if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument("--input", required=True) parser.add_argument("--output", required=True) args = parser.parse_args()
run_pipeline(args.input, args.output, Path("/tmp/work"))Este script se ejecuta tanto en tu portátil (con archivos pequeños) como en el contenedor Batch (con archivos grandes). La diferencia no es el código Python, sino el entorno de ejecución.
Desvanecimiento: los tres espacios en blanco que debes completar
Espacio en blanco 1: Orquestador de envío masivo
Envía automáticamente 100 muestras en lugar de una sola.
def submit_batch(sample_manifest: str, output_bucket: str) -> list[str]: """ sample_manifest: CSV with columns [sample_id, input_s3] Enviar un trabajo para cada fila y devolver una lista de IDs de trabajo """ import csv
job_ids = [] with open(sample_manifest) as f: reader = csv.DictReader(f) for row in reader: # TODO: llamar a submit_sample_job # Combinar output_s3 con output_bucket + sample_id # Agregar el job_id devuelto a job_ids pass
return job_idsPista: output_s3 = f"s3://{output_bucket}/results/{row['sample_id']}/"
Espacio en blanco 2: Sondeo de progreso
Verifique periódicamente el estado de los trabajos enviados.
def wait_for_jobs(job_ids: list[str], poll_interval: int = 60) -> dict: """ Devolver el estado final de cada trabajo. Estado: SUBMITTED, PENDING, RUNNABLE, STARTING, RUNNING, SUCCEEDED, FAILED """ import time
final_states = {} pending = set(job_ids)
while pending: # TODO: llamar a batch.describe_jobs(jobs=list(pending)) # Verificar el estado de cada trabajo, registrar los completados en final_states y eliminarlos de pending # Si quedan trabajos pendientes, esperar poll_interval segundos pass
return final_statesPista: response["jobs"] Verifique el campo status de cada elemento. Solo SUCCEEDED y FAILED se consideran estados finales.
Espacio en blanco 3: Lógica de reintento de fallos
Las instancias Spot son de bajo costo, pero pueden interrumpirse. Los trabajos fallidos se reintentan automáticamente.
def submit_with_retry(sample_id: str, input_s3: str, output_s3: str, max_attempts: int = 3) -> str: response = batch.submit_job( jobName=f"fastq-{sample_id}", jobQueue="my-fastq-queue", jobDefinition="fastq-pipeline", parameters={"input_s3": input_s3, "output_s3": output_s3}, # TODO: agregar retryStrategy # especificar attempts=max_attempts y la condición evaluateOnExit ) return response["jobId"]Pista:
retryStrategy={ "attempts": max_attempts, "evaluateOnExit": [ {"onStatusReason": "Host EC2*", "action": "RETRY"}, {"onExitCode": "0", "action": "EXIT"} ]}Reflexión — ¿En qué se diferencia este pipeline de un pipeline de secuenciación en producción?
Lenguaje de flujo de trabajo: En producción se utilizan lenguajes de flujo de trabajo como Nextflow, Snakemake o WDL. Mientras que vuestro script de Python ejecuta cada paso de forma secuencial, los lengendajes de flujo de trabajo expresan las dependencias entre pasos mediante un DAG (grafo acíclico dirigido), automatizando la ejecución en paralelo y el almacenamiento en caché.
Gestión de datos de referencia: Aunque la referencia hg38 está codificada (hardcoded) en vuestro contenedor, en producción se almacena en un EFS (Elastic File System) o en un volumen grande independiente para que sea compartido por múltiples trabajos.
Optimización de costes: En producción se aprovechan intensivamente las instancias Spot para reducir los costes hasta un 70%. No obstante, es imprescindible implementar lógica para gestionar la interrupción de las instancias Spot. Además, se separan los requisitos de CPU y memoria de cada paso, ejecutando los pasos más ligeros en instancias más pequeñas.
Seguridad: En producción son fundamentales el principio de mínimo privilegio en los roles de IAM, el cifrado de los buckets de S3, la red privada mediante VPC endpoints y los registros de auditoría (CloudTrail).
Observabilidad: Se recopilan en CloudWatch la salida estándar (stdout/stderr) de cada trabajo, el uso de recursos y el tiempo de ejecución. En caso de fallo, se envían notificaciones inmediatas mediante SNS.
Proyectos de expansión
1. Migración a Nextflow: Reescribir el pipeline de Python anterior en Nextflow DSL2. Separar cada paso en un process para verificar la paralelización.
2. Panel de consulta de resultados: Generar automáticamente un resumen de los archivos VCF resultantes de los trabajos completados y mostrarlos en un panel web. Con Streamlit es muy rápido.
3. Seguimiento de costes: Calcular los recursos reales utilizados y el coste de cada trabajo para generar informes de coste por muestra. Utilizar la API de Cost Explorer.
4. Modo de desarrollo local: Construir un entorno de desarrollo que simule S3 y Batch localmente mediante LocalStack. Desarrollar el pipeline sin incurrir en costes reales de AWS.
Mapa de componentes de este capítulo
- [F] Despliegue en AWS: Diseño del entorno de cómputo de Batch, cola de trabajos y definición de trabajo.
- [F] S3: Almacenamiento y acceso a grandes volúmenes de datos. Comprensión de regiones, costes y transferencia.
- [W] venv·file-io: Entornos virtuales de Python y manejo de archivos locales (utilizados dentro del contenedor).
- [W] Docker: Construcción de imágenes de contenedor y subida a ECR (se proporciona el script completo).
[F] = Implementado por vosotros directamente / [W] = Concepto de herramienta proporcionada con código completo.