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)
- No mezcles APIs bloqueantes en tareas async sin spawn_blocking.
- Controla el grado de concurrencia (semaphore, bounded channels).
- Propaga timeouts y límites de retry explícitamente.
- Prefiere estructuras asíncronas (tokio::sync::Mutex) cuando el lock se puede aguantar en un future.
- 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.
¿Quieres comentar?
Inicia sesión con Telegram para participar en la conversación