Producer / Consumer

Desacopla quien genera trabajo de quien lo procesa, usando una cola en memoria como buffer intermedio.

Contexto

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.

Solución

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.

CSHARP
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#

CSHARP
// Registro (singleton — misma instancia para producer y consumer)
services.AddSingleton(Channel.CreateBounded<string>(new BoundedChannelOptions(500)
{
    FullMode = BoundedChannelFullMode.Wait,
    SingleReader = false
}));

services.AddHostedService<MiConsumer>();
CSHARP
// 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);
    }
}
CSHARP
// 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);
            }
        }
    }
}
Warning

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:

CSHARP
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:

CSHARP
// 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

CSHARP
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.

Tradeoffs
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
Reintentos Manual Automático
Multi-proceso No
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.

CSHARP
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