Bajo presión: Cómo los sistemas de cola manejan la contrapresión con ejemplos en C# (Español (Spanish))

Bajo presión: Cómo los sistemas de cola manejan la contrapresión con ejemplos en C#

Sunday, 23 November 2025

//

16 minute read

La contrapresión es el héroe desconocido de los sistemas distribuidos. Es lo que evita que tus colas estallen en las costuras cuando los productores están disparando mensajes más rápido de lo que los consumidores pueden masticar a través de ellos. Dicho simplemente: es el sistema diciendo "espera un momento" cuando las cosas se ponen demasiado ocupadas.

Esta es la cosa: las técnicas de este artículo se aplican prácticamente a cada cola de mensajes y bus de servicio—RabbitMQ, Kafka, Azure Service Bus, AWS SQS, NATS, lo que sea. Los detalles varían, pero los principios son universales. Usaré RabbitMQ para la mayoría de los ejemplos porque es lo que sé mejor, pero te mostraré cómo estos patrones se traducen a través de plataformas.

Una confesión: Incluso la mayoría de los desarrolladores senior no implementan un manejo adecuado de la contrapresión. Construyen sistemas de camino feliz que funcionan bien en el desarrollo y la puesta en escena, y luego se preguntan por qué la producción se cae durante el Viernes Negro. El manejo de la contrapresión es una de esas técnicas que separa "funciona" de "escala". Si no estás pensando en ello, estás construyendo un sistema que eventualmente fallará bajo carga.

Qué es la contrapresión

En su núcleo, la contrapresión es un bucle de retroalimentación que ralentiza a los productores cuando los consumidores se quedan atrás. Piense en ello como semáforos en una carretera resbaladiza, no se puede simplemente apilar en la autopista siempre que lo desee. Las luces controlan el flujo, dejando que los coches se fusionen de forma segura sin causar una acumulación.

Sin contrapresión, un productor rápido abrumará a un consumidor lento. Los mensajes se acumulan en las colas, la memoria se agota, y eventualmente su sistema se cae. La contrapresión dice "en marcha" antes del desastre.

flowchart LR
    P[Producer] --> Q[Queue]
    Q --> C[Consumer]
    C -. "Slow down!" .-> P

    style P stroke:#f59e0b,stroke-width:2px
    style Q stroke:#0ea5e9,stroke-width:2px
    style C stroke:#10b981,stroke-width:2px

La belleza de la contrapresión es que es un conversación El consumidor señala: "Estoy lleno, dame un minuto" y el productor responde: "No hay problema, esperaré". Es cortés, colaborativo y evita que todo el mundo se caiga.

Cómo RabbitMQ maneja la contrapresión

RabbitMQ tiene varios mecanismos incorporados para manejar la contrapresión, y comprenderlos es crucial si estás construyendo sistemas que necesitan mantenerse en posición vertical bajo carga. Documentación del ConejoMQ es excelente, voy a enlazar a páginas específicas a medida que avanzamos.

Control de flujo

Cuando el uso de memoria de RabbitMQ o la profundidad de la cola excede los umbrales configurados, se activa control de flujo. Esto bloquea temporalmente las conexiones de los editores: los editores no pueden enviar nuevos mensajes hasta que el bróker haya eliminado suficiente retraso. alarmas de memoria y alarmas de disco que activan el control de flujo.

flowchart TD
    subgraph "RabbitMQ Flow Control"
        A[Publisher Sends Message] --> B{Memory/Queue<br/>Threshold OK?}
        B -->|Yes| C[Message Accepted]
        C --> D[Add to Queue]
        B -->|No| E[Connection Blocked]
        E --> F[Publisher Waits]
        F --> G{Threshold<br/>Cleared?}
        G -->|No| F
        G -->|Yes| H[Connection Unblocked]
        H --> A
    end

    style E stroke:#ef4444,stroke-width:3px
    style H stroke:#10b981,stroke-width:2px

La idea clave aquí es que RabbitMQ no solo deja caer mensajes cuando está bajo presión, sino que ralentiza la fuente. Este es un enfoque mucho más civilizado que descartar datos silenciosamente.

Agradecimientos a los consumidores

Los consumidores controlan el ritmo a través de acuse de recibo (ACKs y NACKs). Un mensaje no se elimina de la cola hasta que el consumidor lo reconozca explícitamente. Si un consumidor no hace mensajes ACK lo suficientemente rápido, la cola crece, lo que eventualmente activa el control de flujo aguas arriba.

También puede utilizar prefetch limits (QoS) para controlar cuántos mensajes no reconocidos puede tener un consumidor en vuelo a la vez, lo que impide que un único consumidor lento acumule mensajes.

// Set prefetch count to limit unacknowledged messages
channel.BasicQos(prefetchSize: 0, prefetchCount: 10, global: false);

Esto le dice a RabbitMQ: "Sólo envíeme 10 mensajes a la vez. Una vez que CK algunos, usted puede enviar más." Es el consumidor diciendo explícitamente la cantidad de presión que puede manejar. Documentación del cliente .NET cubre la API en detalle.

Cómo otros sistemas manejan la contrapresión

Los patrones son universales, pero las implementaciones difieren. He aquí cómo otros sistemas de mensajería populares abordan el mismo problema.

Kafka: Encuesta controlada por el consumidor

Kafka adopta un enfoque fundamentalmente diferente: los consumidores tire Esto hace que la contrapresión esté implícita: si un consumidor no hace una encuesta, no recibe mensajes. Al bróker no le importa; sólo mantiene los mensajes alrededor hasta que el consumidor esté listo.

// Kafka consumer with explicit backpressure control
using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");

while (!cancellationToken.IsCancellationRequested)
{
    // Only fetch what you can handle - this IS your backpressure
    var result = consumer.Consume(timeout: TimeSpan.FromSeconds(1));

    if (result != null)
    {
        await ProcessMessageAsync(result.Message.Value);

        // Manual commit = explicit acknowledgement
        consumer.Commit(result);
    }

    // If processing is slow, you simply poll less frequently
    // Kafka doesn't push more messages at you
}

El bit inteligente: los grupos de consumidores de Kafka reequilibran automáticamente las particiones. Si un consumidor se queda atrás, puede añadir más consumidores al grupo y las particiones se redistribuyen. La contrapresión se convierte en una decisión de escalado.

// Control batch size to manage memory pressure
var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "order-processors",
    AutoOffsetReset = AutoOffsetReset.Earliest,
    MaxPollIntervalMs = 300000,      // 5 mins max between polls
    MaxPartitionFetchBytes = 1048576, // 1MB max per partition fetch
    FetchMaxBytes = 52428800          // 50MB max total fetch
};

Autobús de servicio Azure: Control de mensajes concurrentes

Azure Service Bus utiliza un MaxConcurrentCalls configuración que es muy simple: controla cuántos mensajes maneja el procesador simultáneamente. La contrapresión es automática.

var processor = client.CreateProcessor("orders-queue", new ServiceBusProcessorOptions
{
    // This IS your backpressure - only process 10 at a time
    MaxConcurrentCalls = 10,
    AutoCompleteMessages = false,
    PrefetchCount = 20  // Buffer 20 messages locally
});

processor.ProcessMessageAsync += async args =>
{
    try
    {
        await ProcessOrderAsync(args.Message.Body.ToString());
        await args.CompleteMessageAsync(args.Message);
    }
    catch (Exception ex)
    {
        // Abandon returns message to queue for retry
        await args.AbandonMessageAsync(args.Message);
    }
};

processor.ProcessErrorAsync += args =>
{
    Console.WriteLine($"Error: {args.Exception.Message}");
    return Task.CompletedTask;
};

await processor.StartProcessingAsync();

Azure Service Bus también soporta períodos de sesiones para el tratamiento ordenado y colas de letras muertas para los mensajes que fallan repetidamente, ambos importantes para manejar la presión cuando las cosas salen mal.

AWS SQS: Danza de tiempo de espera de visibilidad

SQS utiliza los tiempos de visibilidad como su mecanismo de contrapresión. Cuando recibes un mensaje, se vuelve invisible para otros consumidores. Si no lo eliminas a tiempo, reaparece para que alguien más lo intente.

var sqsClient = new AmazonSQSClient();

// Receive with explicit backpressure control
var response = await sqsClient.ReceiveMessageAsync(new ReceiveMessageRequest
{
    QueueUrl = queueUrl,
    MaxNumberOfMessages = 10,           // Batch size = backpressure control
    WaitTimeSeconds = 20,               // Long polling
    VisibilityTimeout = 300             // 5 mins to process before retry
});

foreach (var message in response.Messages)
{
    try
    {
        await ProcessAsync(message.Body);

        // Only delete after successful processing
        await sqsClient.DeleteMessageAsync(queueUrl, message.ReceiptHandle);
    }
    catch
    {
        // Don't delete - message will become visible again after timeout
        // Optionally, change visibility timeout to retry sooner
        await sqsClient.ChangeMessageVisibilityAsync(queueUrl,
            message.ReceiptHandle, visibilityTimeout: 0);
    }
}

El truco inteligente SQS: utilizar ApproximateNumberOfMessages para controlar la profundidad de la cola y los consumidores a escala automática:

var attributes = await sqsClient.GetQueueAttributesAsync(new GetQueueAttributesRequest
{
    QueueUrl = queueUrl,
    AttributeNames = new List<string> { "ApproximateNumberOfMessages" }
});

var depth = int.Parse(attributes.Attributes["ApproximateNumberOfMessages"]);

if (depth > 1000)
{
    // Signal to scale up consumers
    await TriggerAutoScalingAsync();
}

NATS JetStream: Control de flujo incorporado

NATS JetStream tiene un control de flujo explícito con reconocimientos al consumidor y límites máximos de mensajes pendientes:

var js = connection.CreateJetStreamContext();

var subscription = js.PushSubscribeAsync("orders.>", (sender, args) =>
{
    try
    {
        ProcessMessage(args.Message.Data);
        args.Message.Ack();
    }
    catch
    {
        args.Message.Nak();  // Negative ack - redeliver
    }
}, new PushSubscribeOptions.Builder()
    .WithConfiguration(new ConsumerConfiguration.Builder()
        .WithMaxAckPending(100)     // Max unacked messages - THIS is backpressure
        .WithAckWait(30000)         // 30 seconds to ack
        .Build())
    .Build());

El hilo común

Observe lo que todos estos sistemas tienen en común:

  1. Límites explícitos de lote/moneda - controlas lo mucho que estás dispuesto a manejar
  2. Flujo basado en el reconocimiento - los mensajes permanecerán disponibles hasta que confirme el procesamiento
  3. Recuperación basada en el tiempo de espera - si fallas, los mensajes regresan para reintentar
  4. Monitorización de profundidad - siempre puedes preguntar "¿cómo estoy?".

La sintaxis difiere, pero la danza es la misma: "Aquí está lo mucho que puedo manejar. Dime cuando lo he manejado. Si no te lo digo a tiempo, asume que fallé".

C# Ejemplos de código

Bien, vamos a entrar en el código. Aquí hay ejemplos prácticos de implementar y responder a la contrapresión en sus aplicaciones C#.

Monitorización de la profundidad de la cola

Lo primero es lo primero: no se puede manejar lo que no se puede medir. Aquí está cómo comprobar cuántos mensajes están esperando en una cola:

var queue = channel.QueueDeclare(
    queue: "tasks",
    durable: true,
    exclusive: false,
    autoDelete: false);

Console.WriteLine($"Messages ready: {queue.MessageCount}");

// React to queue depth
if (queue.MessageCount > 1000)
{
    Console.WriteLine("Queue backing up - consider throttling producers");
}

Este fragmento comprueba cuántos mensajes están esperando. Si el conteo sube, esa es su señal para acelerar a los productores o aumentar el tamaño de los consumidores. No se preocupe por comprobar esto con demasiada frecuencia, un chequeo médico periódico suele ser suficiente.

Publisher Confirma

Publisher confirma hacerle saber cuando RabbitMQ ha recibido y procesado su mensaje con éxito. Si los reconocimientos se ralentizan, esa es una clara señal de contrapresión:

// Enable publisher confirms
channel.ConfirmSelect();

var body = Encoding.UTF8.GetBytes("Hello, Queue!");

channel.BasicPublish(
    exchange: "",
    routingKey: "tasks",
    basicProperties: null,
    body: body);

// Wait for confirmation - timeout indicates backpressure
bool confirmed = channel.WaitForConfirms(TimeSpan.FromSeconds(5));

if (!confirmed)
{
    Console.WriteLine("Message not confirmed - broker may be under pressure");
}

Si RabbitMQ está luchando, las confirmaciones tardan más tiempo. Su productor puede utilizar esta señal para retroceder en lugar de acumular más presión.

Para escenarios de alto rendimiento, querrá confirmar asíncrono:

channel.ConfirmSelect();

var outstandingConfirms = new ConcurrentDictionary<ulong, string>();

channel.BasicAcks += (sender, ea) =>
{
    if (ea.Multiple)
    {
        var confirmed = outstandingConfirms.Where(k => k.Key <= ea.DeliveryTag);
        foreach (var entry in confirmed)
        {
            outstandingConfirms.TryRemove(entry.Key, out _);
        }
    }
    else
    {
        outstandingConfirms.TryRemove(ea.DeliveryTag, out _);
    }
};

channel.BasicNacks += (sender, ea) =>
{
    // Message was rejected - implement retry logic
    Console.WriteLine($"Message {ea.DeliveryTag} was nacked - broker under pressure");
    // Back off before retrying
};

Reintentar con el retroceso exponencial

Cuando detectas contrapresión, lo peor que puedes hacer es volver a intentarlo inmediatamente a toda velocidad. Es como responder a un atasco de tráfico presionando el acelerador con más fuerza. En su lugar, implementa un retroceso exponencial:

public async Task PublishWithBackpressureAsync(
    IModel channel,
    byte[] body,
    int maxRetries = 5)
{
    int attempt = 0;

    while (attempt < maxRetries)
    {
        try
        {
            channel.ConfirmSelect();
            channel.BasicPublish(
                exchange: "",
                routingKey: "tasks",
                basicProperties: null,
                body: body);

            if (channel.WaitForConfirms(TimeSpan.FromSeconds(5)))
            {
                return; // Success
            }

            throw new Exception("Publish not confirmed");
        }
        catch (Exception ex)
        {
            attempt++;

            if (attempt >= maxRetries)
            {
                throw new Exception($"Failed to publish after {maxRetries} attempts", ex);
            }

            // Exponential backoff: 1s, 2s, 4s, 8s, 16s
            var delay = TimeSpan.FromSeconds(Math.Pow(2, attempt - 1));
            Console.WriteLine($"Backpressure detected - retry {attempt} after {delay}");

            await Task.Delay(delay);
        }
    }
}

Esto imita el patrón de HTTP 429 (Demasiadas peticiones). En lugar de martillar al bróker, hacemos una pausa antes de volver a intentarlo, dando al sistema tiempo para recuperarse.

Uso de canales para la contrapresión en proceso

Si está construyendo una tubería interna (productor → procesador → consumidor todo dentro de su aplicación), .NET's Channel<T> proporciona un elegante soporte de contrapresión:

// Create a bounded channel - backpressure is automatic
var channel = Channel.CreateBounded<WorkItem>(new BoundedChannelOptions(100)
{
    FullMode = BoundedChannelFullMode.Wait // Block producer when full
});

// Producer - will automatically wait when channel is full
async Task ProduceAsync(ChannelWriter<WorkItem> writer)
{
    for (int i = 0; i < 10000; i++)
    {
        var item = new WorkItem { Id = i };

        // This awaits if the channel is at capacity
        await writer.WriteAsync(item);

        Console.WriteLine($"Produced item {i}");
    }

    writer.Complete();
}

// Consumer - processes at its own pace
async Task ConsumeAsync(ChannelReader<WorkItem> reader)
{
    await foreach (var item in reader.ReadAllAsync())
    {
        // Simulate slow processing
        await Task.Delay(100);
        Console.WriteLine($"Processed item {item.Id}");
    }
}

// Run both concurrently
await Task.WhenAll(
    ProduceAsync(channel.Writer),
    ConsumeAsync(channel.Reader)
);

El canal delimitado aplica automáticamente la contrapresión: el productor bloquea cuando el canal está lleno, naturalmente desacelerando para igualar el ritmo del consumidor. No se requiere estrangulamiento manual.

Un editor completo de backpressure-Aware

He aquí un ejemplo más completo que reúne el monitoreo, confirma y retrocede:

public class BackpressureAwarePublisher : IDisposable
{
    private readonly IConnection _connection;
    private readonly IModel _channel;
    private readonly string _queueName;
    private readonly int _queueDepthThreshold;

    public BackpressureAwarePublisher(
        string hostName,
        string queueName,
        int queueDepthThreshold = 1000)
    {
        var factory = new ConnectionFactory { HostName = hostName };
        _connection = factory.CreateConnection();
        _channel = _connection.CreateModel();
        _queueName = queueName;
        _queueDepthThreshold = queueDepthThreshold;

        _channel.QueueDeclare(
            queue: queueName,
            durable: true,
            exclusive: false,
            autoDelete: false);

        _channel.ConfirmSelect();
    }

    public async Task<bool> PublishAsync(byte[] body, CancellationToken ct = default)
    {
        // Check queue depth first
        var queueInfo = _channel.QueueDeclarePassive(_queueName);

        if (queueInfo.MessageCount > _queueDepthThreshold)
        {
            Console.WriteLine($"Queue depth {queueInfo.MessageCount} exceeds threshold - applying backpressure");

            // Wait for queue to drain a bit
            while (queueInfo.MessageCount > _queueDepthThreshold * 0.8)
            {
                await Task.Delay(1000, ct);
                queueInfo = _channel.QueueDeclarePassive(_queueName);
            }
        }

        // Publish with retry
        for (int attempt = 1; attempt <= 3; attempt++)
        {
            try
            {
                var properties = _channel.CreateBasicProperties();
                properties.Persistent = true;

                _channel.BasicPublish(
                    exchange: "",
                    routingKey: _queueName,
                    basicProperties: properties,
                    body: body);

                if (_channel.WaitForConfirms(TimeSpan.FromSeconds(5)))
                {
                    return true;
                }
            }
            catch (Exception ex)
            {
                Console.WriteLine($"Publish attempt {attempt} failed: {ex.Message}");
            }

            if (attempt < 3)
            {
                await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, attempt)), ct);
            }
        }

        return false;
    }

    public void Dispose()
    {
        _channel?.Dispose();
        _connection?.Dispose();
    }
}

Mejores prácticas

Mantenga la calma

No entres en pánico cuando crezcan las colas. Un poco de profundidad es normal y saludable, significa que tu sistema está absorbiendo los picos de carga con gracia. El objetivo no es una cola vacía; es un estable Una cola que no crece sin límites.

Supervisar la profundidad de la cola con el tiempo. Busque tendencias, no instantáneas. Una cola consistentemente a 100 mensajes está bien. Una cola que ha crecido de 100 a 10.000 en la última hora necesita atención.

Sea pragmático

Aplicar patrones pragmáticamente. No todos los mensajes necesitan confirmación editorial. No todas las colas necesitan un manejo sofisticado de la contrapresión. Una cola que procesa 10 mensajes por hora probablemente no necesita la misma ingeniería de resiliencia que un procesamiento de 10.000 por segundo.

Pregúntese: "¿Cuál es el costo real si este mensaje se pierde o se retrasa?" Si la respuesta es "no mucho", no sobreingeniería. Si la respuesta es "impacto significativo de integridad financiera o de datos," invierta en el manejo adecuado de la contrapresión.

Escalar inteligente

Cuando las colas están retrocediendo, la respuesta no siempre es "añadir más productores". Eso es como tratar de arreglar un atasco de tráfico añadiendo más coches.

Considere:

  • Los consumidores de escala primero - ¿Puedes agregar más trabajadores para procesar el trabajo atrasado?
  • Comprobación de los cuellos de botella - es una dependencia lenta aguas abajo causando la copia de seguridad?
  • Lote cuando sea posible - pueden los consumidores procesar múltiples mensajes a la vez?
flowchart TD
    A[Queue Growing] --> B{Consumer<br/>Saturated?}
    B -->|Yes| C[Add Consumers]
    B -->|No| D{Downstream<br/>Bottleneck?}
    D -->|Yes| E[Fix/Scale Downstream]
    D -->|No| F{Can Batch<br/>Process?}
    F -->|Yes| G[Implement Batching]
    F -->|No| H[Accept Higher Latency<br/>or Reduce Load]

    style C stroke:#10b981,stroke-width:2px
    style E stroke:#f59e0b,stroke-width:2px
    style G stroke:#0ea5e9,stroke-width:2px

Monitor y alerta

Establecer alertas para:

  • Profundidad de la cola superior a los umbrales
  • Aumento del desfase de los consumidores
  • El editor confirma la hora de salida
  • Eventos bloqueados de conexión

¿Quieres saber sobre la contrapresión? antes se convierte en una crisis, no cuando su sistema ya ha caído.

Conclusión

La contrapresión no es solo una limitación de tasas, es una táctica de supervivencia. Al tratarla como una conversación entre productor y consumidor, se construyen sistemas que se mantienen resistentes bajo presión.

Los puntos de vista clave:

  1. La contrapresión es retroalimentación - productores y consumidores que colaboran para encontrar un rendimiento sostenible
  2. Monitorear la profundidad de la cola - no puedes manejar lo que no puedes medir
  3. El editor de uso confirma - saber cuando el corredor está luchando
  4. Implementar un retroceso exponencial - no martillees un sistema que ya está bajo presión
  5. Consumidores de escala, no sólo productores - arreglar el cuello de botella, no el síntoma

Cuando tu sistema dice "estoy lleno, dame un minuto", la respuesta correcta es "No hay problema, esperaré". Esa es la esencia de los sistemas distribuidos bien educados: sólidos, colaborativos y resistentes.

Mantenga la calma cuando las colas crecen un poco, y recuerde: un sistema bajo una elegante contrapresión es infinitamente mejor que uno que se ha caído por completo.

Finding related posts...
logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.