Guía completa de concurrencia y async/await en Rust para desarrolladores

rust Guía completa de concurrencia y async/await en Rust para desarrolladores

Guía completa de concurrencia y async/await en Rust para desarrolladores

Esta guía explica los conceptos clave de la concurrencia asíncrona en Rust, patrones prácticos y un ejemplo completo: un pequeño CLI que descarga URLs concurrentemente con límite de concurrencia y timeouts. Incluye por qué se eligen ciertas abstracciones, errores comunes y técnicas de diagnóstico.

Conceptos clave (rápido)

  • Future: una tarea asíncrona perezosa. Solo avanza cuando el executor la poll-ea.
  • Executor: motor que ejecuta futures (tokio, async-std, smol).
  • Send y Sync: requisitos para enviar closures/futures a otros threads (tokio::spawn exige que la future sea 'Send + 'static).
  • Pin: garantiza que el objeto no se mueva en memoria (interesante en futures que contienen referencias).
  • await: punto de suspensión que permite al executor programar otras tareas.

Patrones prácticos y cuándo usarlos

  • tokio::spawn para tareas concurrentes ligeras y 'Send'.
  • Semaphore para limitar concurrencia (control de backpressure).
  • FuturesUnordered para procesar resultados a medida que llegan (streaming de resultados).
  • tokio::time::timeout para proteger operaciones que pueden bloquear indefinidamente.
  • mpsc channels para productor/consumidor y desacoplar etapas.
  • spawn_blocking para operaciones CPU-bound o que usan APIs bloqueantes.
  • Arc + tokio::sync::Mutex / RwLock para estado compartido. Evitar std::sync::Mutex dentro de async (puede bloquear el hilo del executor).

Errores comunes

  • Usar APIs bloqueantes en un context async sin spawn_blocking -> congelamiento de executor y rendimiento pobre.
  • Spawn de futures no 'Send (por ejemplo, referencias no 'static): compile-time error al usar tokio::spawn.
  • Olvidar manejar cancelación: si un task es cancelado, recursos pueden quedar sin liberar si no se usan finally/Drop/guards.
  • Contención excesiva al usar Mutex para sincronización fina -> reemplazar con estructuras concurrentes o rediseñar acceso.
  • Crear demasiados tasks sin límite de concurrencia -> OOM, throttling.

Proyecto práctico: fetcher (descarga concurrente con límite y timeouts)

Objetivo: construir un CLI mínimo que lea urls desde un archivo, las descargue concurrentemente con límite de concurrencia configurable y aplique timeout por petición. Uso: tokio + reqwest + semaphore.

Estructura de carpetas

fetcher/
├─ Cargo.toml
└─ src/
   └─ main.rs
urls.txt  (lista de URLs, una por línea)

Cargo.toml

[package]
name = "fetcher"
version = "0.1.0"
edition = "2021"

[dependencies]
tokio = { version = "1", features = ["full"] }
reqwest = { version = "0.11", features = ["gzip", "rustls-tls"] }
anyhow = "1.0"
futures = "0.3"

src/main.rs (código completo)

use anyhow::Result;
use reqwest::Client;
use std::sync::Arc;
use tokio::sync::Semaphore;

#[tokio::main]
async fn main() -> Result<()> {
    // Args: [path_to_urls] [concurrency]
    let args: Vec = std::env::args().collect();
    let path = args.get(1).map(|s| s.as_str()).unwrap_or("urls.txt");
    let concurrency: usize = args
        .get(2)
        .and_then(|s| s.parse().ok())
        .unwrap_or(10);

    let contents = std::fs::read_to_string(path)?;
    let urls: Vec = contents
        .lines()
        .map(|l| l.trim())
        .filter(|l| !l.is_empty())
        .map(|s| s.to_string())
        .collect();

    let client = Client::builder().user_agent("fetcher/1.0").build()?;
    let sem = Arc::new(Semaphore::new(concurrency));

    let timeout = std::time::Duration::from_secs(10);

    // Lanzamos tasks y recopilamos handles
    let mut handles = Vec::with_capacity(urls.len());

    for url in urls {
        let client = client.clone();
        let sem = sem.clone();
        let url_clone = url.clone();

        let handle = tokio::spawn(async move {
            // Adquirir permiso de concurrencia
            let permit = sem.acquire().await.expect("semaphore closed");

            // Timeout envolviendo la operación async
            let res = tokio::time::timeout(timeout, async {
                let resp = client.get(&url_clone).send().await?;
                let status = resp.status().as_u16();
                let body = resp.bytes().await?;
                Ok::<(u16, usize), reqwest::Error>((status, body.len()))
            })
            .await;

            // liberar permiso explícitamente (se libera cuando permit cae fuera de scope)
            drop(permit);

            match res {
                Ok(Ok((status, size))) => Ok::<_, anyhow::Error>((url_clone, status, size)),
                Ok(Err(e)) => Err(anyhow::anyhow!(e)),
                Err(_) => Err(anyhow::anyhow!("timeout")),
            }
        });

        handles.push(handle);
    }

    // Recolectar resultados
    for h in handles {
        match h.await {
            Ok(Ok((url, status, size))) => println!("{} -> {} ({} bytes)", url, status, size),
            Ok(Err(e)) => eprintln!("task error: {}", e),
            Err(e) => eprintln!("join error: {}", e),
        }
    }

    Ok(())
}

Por qué esta implementación

  • tokio::spawn ejecuta cada descarga de forma concurrente y aislada. Requiere que la future sea 'Send + 'static; por eso clonamos cualquier dato necesario (Client es barato de clonar y es Sync + Send).
  • Semaphore limita el número de descargas simultáneas y evita crear demasiadas conexiones o saturar memoria/CPU.
  • tokio::time::timeout evita tareas colgadas indefinidamente; se captura como error y la tarea retorna rápidamente.
  • La combinación de timeout + semaphore + spawn proporciona control fino sobre recursos sin bloquear el executor.

Variantes y mejoras

  • Usar futures::stream::FuturesUnordered para procesar resultados conforme llegan y reducir memoria al no acumular muchos handles.
  • Usar tokio::sync::mpsc para pipeline con productores (lectura del archivo) y consumidores (descarga), útil si el archivo es muy grande.
  • Usar reqwest::Client pool con configuración de conexión, timeouts globales y límites de conexiones para producción.
  • Para tareas CPU-bound (procesamiento de cuerpo) usar tokio::task::spawn_blocking para no bloquear hilos del runtime.

Diagnóstico y herramientas

  • tracing + tracing_subscriber para logs estructurados y correlacionar tasks.
  • tokio-console para inspeccionar tareas en tiempo real (si usas tokio con la feature correspondiente).
  • Flamegraphs y perf para identificar contención y hotspots CPU.

Resumen de buenas prácticas (rápido)

  1. No mezcles APIs bloqueantes en tareas async sin spawn_blocking.
  2. Controla el grado de concurrencia (semaphore, bounded channels).
  3. Propaga timeouts y límites de retry explícitamente.
  4. Prefiere estructuras asíncronas (tokio::sync::Mutex) cuando el lock se puede aguantar en un future.
  5. Usa un client compartido (reqwest::Client) y clónalo; evita crear uno por conexión.

Para seguir mejorando tu pipeline, instrumenta con tracing/tokio-console y prueba diferentes estrategias (bounded channels, FuturesUnordered, pánico controlado) para encontrar el mejor trade-off entre latencia y uso de recursos.

Comentarios
¿Quieres comentar?

Inicia sesión con Telegram para participar en la conversación


Comentarios (0)

Aún no hay comentarios. ¡Sé el primero en comentar!

Iniciar Sesión