Producer / Consumer
Desacopla quien genera trabajo de quien lo procesa, usando una cola en memoria como buffer intermedio.
Tenés un flujo que genera trabajo más rápido de lo que se puede procesar. Si quien genera espera a quien procesa, bloqueás el request. Si los acoplás directamente, cualquier lentitud del procesador se propaga hacia arriba.
Separás en dos roles independientes:
- Producer: genera trabajo y lo deposita en una cola. No sabe quién ni cuándo lo va a procesar.
- Consumer: lee de la cola a su propio ritmo y procesa cada item.
La cola actúa como buffer: absorbe la diferencia de velocidad entre ambos.
[Producer] → [ cola / Channel<T> ] → [Consumer]
rápido buffer lento
El consumer duerme cuando la cola está vacía — sin polling, sin CPU desperdiciada. Se despierta solo cuando llega un item.
#Cómo funciona internamente
La magia está en TaskCompletionSource<T>. Cuando el consumer llama WaitToReadAsync y la cola está vacía, el runtime suspende ese código y libera el thread completamente. Cuando el producer llama WriteAsync, internamente llama TrySetResult(true) — eso resuelve la Task que estaba esperando y el runtime agenda el consumer para continuar.
Producer: WriteAsync("item")
→ encola el item
→ TrySetResult(true) ← despierta al consumer
→ runtime agenda continuation
→ consumer recibe el item
→ procesa
→ vuelve a dormir
No hay polling. No hay Thread.Sleep. No hay CPU consumida esperando.
#Ejemplo en C#
// Registro (singleton — misma instancia para producer y consumer)
services.AddSingleton(Channel.CreateBounded<string>(new BoundedChannelOptions(500)
{
FullMode = BoundedChannelFullMode.Wait,
SingleReader = false
}));
services.AddHostedService<MiConsumer>();
// Producer — puede ser un middleware, un endpoint, cualquier servicio
public class MiProducer(Channel<string> channel)
{
public async Task EncolarAsync(string item)
{
await channel.Writer.WriteAsync(item);
}
}
// Consumer — corre en background, independiente del ciclo HTTP
public class MiConsumer(Channel<string> channel, ILogger<MiConsumer> logger)
: BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken ct)
{
await foreach (var item in channel.Reader.ReadAllAsync(ct))
{
try
{
await ProcesarAsync(item);
}
catch (Exception ex)
{
// try/catch ADENTRO del loop — si falla un item, el consumer sigue vivo
logger.LogError(ex, "Falló procesando {Item}", item);
}
}
}
}
El try/catch tiene que estar dentro del await foreach, no afuera. Si lo ponés afuera, una excepción mata el loop completo y el consumer nunca más procesa nada — sin ningún error visible.
#Backpressure
Con BoundedChannelOptions controlás qué pasa cuando la cola se llena:
| FullMode | Comportamiento |
|---|---|
Wait |
El producer espera hasta que haya slot. El request HTTP se frena. |
DropNewest |
Descarta el item entrante. El producer no se bloquea. |
DropOldest |
Descarta el item más viejo de la cola. |
Para endpoints HTTP que no pueden bloquearse, TryWrite es la alternativa:
if (!channel.Writer.TryWrite(item))
return Results.StatusCode(503); // saturado, reintentá
#Múltiples consumers en paralelo
Si el procesamiento es lento, podés levantar N consumers sobre el mismo channel:
// Registrás el mismo HostedService N veces
services.AddHostedService<MiConsumer>();
services.AddHostedService<MiConsumer>();
services.AddHostedService<MiConsumer>();
Todos leen del mismo Channel<T> singleton. Cada item es procesado por exactamente uno.
#Monitoreo
app.MapGet("/health/queue", (Channel<string> channel) =>
Results.Ok(new { pendientes = channel.Reader.Count }));
Si pendientes crece indefinidamente, el cuello de botella está en el consumer — necesitás más instancias o cachear.
| Pro | Contra |
|---|---|
| Producer y consumer completamente desacoplados | Los mensajes se pierden si el proceso muere (no hay persistencia) |
| Sin polling — 0 CPU cuando la cola está vacía | Sin dead-letter queue — un item que falla se descarta |
Backpressure natural con Bounded |
Sin reintentos automáticos |
| N consumers en paralelo sin cambiar el producer | No escala entre procesos — solo dentro del mismo proceso |
#Cuándo NO usarlo
- Si necesitás durabilidad: un proceso que muere pierde todos los mensajes encolados. Usá Outbox Pattern + Service Bus.
- Si necesitás reintentos automáticos con dead-letter.
Channel<T>no tiene eso. - Si necesitás escalar horizontalmente entre múltiples instancias.
Channel<T>es en memoria — cada proceso tiene su propia cola. Usá RabbitMQ, Azure Service Bus o Kafka. - Si el procesamiento es tan rápido que la latencia de encolado importa — ahí un call directo es más simple.
#Comparación con Service Bus
Channel<T> |
Service Bus / RabbitMQ | |
|---|---|---|
| Persistencia | No — RAM | Sí — disco/red |
| Dead-letter | No | Sí |
| Reintentos | Manual | Automático |
| Multi-proceso | No | Sí |
| Latencia | Microsegundos | Milisegundos |
| Infraestructura | Ninguna | Broker externo |
Channel<T> es un Service Bus en memoria. Sin broker, sin red, sin costo operativo — a cambio de durabilidad.
#Disponibilidad
Channel<T> está disponible desde .NET Core 3.0. No requiere NuGet adicional.
using System.Threading.Channels;
No existe en .NET Framework. El equivalente legacy es BlockingCollection<T>, pero bloquea threads reales en lugar de usar async/await.
#concurrency #async #channel #queue #background-service