# Onder druk: Hoe Wachtrijsystemen Backpressure hanteren met Voorbeelden in C#

<datetime class="hidden">2025-11-23T14:00</datetime>

<!-- category -- Distributed Systems, RabbitMQ, Kafka, Azure Service Bus, C#, Messaging -->
Backpressure is de unsung held van gedistribueerde systemen. Het is wat ervoor zorgt dat je wachtrijen niet barsten in de naden wanneer producenten sneller berichten afvuren dan consumenten er doorheen kunnen kauwen. Simpel gezegd: het is het systeem dat "wacht op een moment" zegt als dingen te druk worden.

**Het zit zo:** de technieken in dit artikel zijn vrijwel van toepassing op *elke* bericht wachtrij en service bus. RabbitMQ, Kafka, Azure Service Bus, AWS SQS, NATS, noem maar op. De details variëren, maar de principes zijn universeel. Ik zal RabbitMQ gebruiken voor de meeste voorbeelden omdat het is wat ik het beste weet, maar ik zal u laten zien hoe deze patronen vertalen over platforms.

> **Een bekentenis:** Zelfs de meeste senior ontwikkelaars implementeren niet de juiste backpressure handling. Ze bouwen happy-path systemen die prima werken in dev en enscenering, dan vraag je je af waarom productie valt tijdens Black Friday. Backpressure handling is een van die technieken die "het werkt" scheidt van "het schalen." Als je er niet over nadenkt, bouw je een systeem dat uiteindelijk zal falen onder belasting.

[TOC]

## Wat Backpressure is

In de kern, backpressure is een terugkoppelingslus die producenten vertraagt wanneer consumenten achterblijven achter. Denk aan het als verkeerslichten op een slip weg.Je kunt niet gewoon stapelen op de snelweg wanneer je wilt. De lichten controleren de stroom, waardoor auto's veilig samensmelten zonder een stapeling te veroorzaken.

Zonder tegendruk zal een snelle producent een trage consument overweldigen. Berichten stapelen zich op in rijen, geheugen raakt uitgeput, en uiteindelijk valt je systeem om. Backpressure zegt "steady on" voordat de ramp toeslaat.

```mermaid
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
```

Het mooie van tegendruk is dat het een *gesprek* De consument geeft aan: "Ik ben vol, geef me een minuut" en de producent antwoordt: "Geen probleem, ik wacht." Het is beleefd, collaboratief, en zorgt ervoor dat iedereen niet omvalt.

## Hoe RabbitMQ Backpressure behandelt

RabbitMQ heeft verschillende ingebouwde mechanismen voor het hanteren van tegendruk, en het begrijpen ervan is cruciaal als je systemen bouwt die rechtop onder belasting moeten blijven. [KonijnMQ documentatie](https://www.rabbitmq.com/docs) Ik zal link naar specifieke pagina's als we gaan.

### Stroomregeling

Wanneer RabbitMQ's geheugengebruik of wachtrijdiepte de geconfigureerde drempels overschrijdt, activeert het [stroomregeling](https://www.rabbitmq.com/docs/flow-control). Dit blokkeert tijdelijk publisher connecties . Uitgevers kunnen geen nieuwe berichten verzenden totdat de makelaar voldoende achterstand heeft opgeruimd . Zie ook de documenten op [geheugenalarmen](https://www.rabbitmq.com/docs/memory) en [schijfalarmen](https://www.rabbitmq.com/docs/disk-alarms) welke trigger flow control.

```mermaid
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
```

Het belangrijkste inzicht hier is dat RabbitMQ niet alleen berichten laat vallen wanneer onder druk het vertraagt de bron. Dit is een veel beschaafder aanpak dan stilletjes gegevens weggooien.

### Erkenningen van de consument

Consumenten controleren het tempo door [bevestigingen](https://www.rabbitmq.com/docs/confirms#consumer-acknowledgements) (ACKs en NACKs). Een bericht wordt niet uit de wachtrij verwijderd totdat de consument het expliciet erkent. Als een consument niet snel genoeg berichten ACK, groeit de wachtrij die uiteindelijk flow control stroomopwaarts activeert.

U kunt ook [prefetch-limieten](https://www.rabbitmq.com/docs/consumer-prefetch) (QoS) om te controleren hoeveel niet-geannexeerde berichten een consument in een keer kan hebben in de vlucht. Dit voorkomt dat een enkele trage consument berichten verzamelt.

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

Dit vertelt RabbitMQ: "Stuur me slechts 10 berichten per keer. Zodra ik ACK sommige, kunt u meer sturen." Het is de consument expliciet zeggen hoeveel druk het aankan. [.NET client documentatie](https://www.rabbitmq.com/client-libraries/dotnet-api-guide) behandelt de API in detail.

## Hoe andere systemen omgaan met Backpressure

De patronen zijn universeel, maar de implementaties verschillen. Hier is hoe sommige andere populaire messaging systemen benaderen hetzelfde probleem.

### Kafka: Consumer Controlled Polling

Kafka kiest een fundamenteel andere aanpak. *trekken* Dit maakt tegendruk impliciet: als een consument niet peilt, ontvangt hij geen berichten. De makelaar kan het niet schelen; het houdt de boodschappen gewoon rond tot de consument klaar is.

```csharp
// 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
}
```

Het slimme stukje: Kafka's consumentengroepen herbalanceren automatisch partities. Als één consument achterop raakt, kunt u meer consumenten toevoegen aan de groep en partities worden herverdeeld. Backpressure wordt een schalende beslissing.

```csharp
// 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
};
```

### Azure Service Bus: Gelijktijdige Berichtcontrole

Azure Service Bus maakt gebruik van een `MaxConcurrentCalls` instelling dat is prachtig eenvoudig . Controleert hoeveel berichten uw processor tegelijkertijd behandelt. Backpressure is automatisch.

```csharp
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 ondersteunt ook *sessies* voor bestelde verwerking en *wachtrijen met dode letters* voor berichten die herhaaldelijk falen ...beide belangrijk voor het beheer van druk wanneer dingen fout gaan.

### AWS SQS: Zichtbaarheid Timeout Dance

SQS maakt gebruik van zicht timeouts als zijn backpressure mechanisme. Wanneer u een bericht ontvangt, wordt het onzichtbaar voor andere consumenten. Als u het niet op tijd verwijdert, komt het opnieuw voor iemand anders om het te proberen.

```csharp
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);
    }
}
```

De slimme SQS truc: gebruik `ApproximateNumberOfMessages` om de wachtrijdiepte en de automatische schaal van de consument te controleren:

```csharp
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: Flow Control Ingebouwd

NATS JetStream heeft expliciete stroomcontrole met consumentenbevestigingen en maximale berichtlimieten:

```csharp
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());
```

### De gemeenschappelijke draad

Merk op wat al deze systemen gemeen hebben:

1. **Expliciete batch/concurrency limieten** - je bepaalt hoeveel je bereid bent om te gaan
2. **Op erkenning gebaseerde stroom** - berichten blijven beschikbaar totdat u de verwerking bevestigt
3. **Timeout-gebaseerd herstel** - als u faalt, keren berichten terug om opnieuw te proberen
4. **Dieptebewaking** Je kunt altijd vragen hoe het met me zit.

De syntaxis verschilt, maar de dans is hetzelfde: *"Hier is hoeveel ik aankan. Zeg me wanneer ik het heb afgehandeld. Als ik het je niet op tijd vertel, neem ik aan dat ik gefaald heb."*

## C# Code Voorbeelden

Hier zijn praktische voorbeelden van het implementeren en reageren op tegendruk in je C#-toepassingen.

### Wachtrij Dieptebewaking

Eerst kunt u niet beheren wat u niet kunt meten. Hier is hoe u kunt controleren hoeveel berichten wachten in een wachtrij:

```csharp
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");
}
```

Dit knipsel controleert hoeveel berichten er wachten. Als de telling klimt, dat is uw cue te gastelen producenten of opschalen consumenten. Maak je geen zorgen over het controleren van deze te vaak een periodieke gezondheidscontrole is meestal voldoende.

### Uitgever bevestigt

Uitgever bevestigt laat u weten wanneer RabbitMQ uw bericht succesvol heeft ontvangen en verwerkt. Als bevestigingen vertragen, is dat een duidelijk signaal van tegendruk:

```csharp
// 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");
}
```

Als RabbitMQ worstelt, duurt de bevestiging langer of de tijd volledig. Uw producer kan dit signaal gebruiken om zich terug te trekken in plaats van te stapelen op meer druk.

Voor high-throughput scenario's wil je asynchrone bevestigingen:

```csharp
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
};
```

### Met exponentieel Backoff opnieuw proberen

Als je tegendruk ontdekt, is het ergste wat je kunt doen direct opnieuw proberen op volle snelheid. Dat is als reageren op een file door het gaspedaal harder in te drukken. In plaats daarvan, exponentieel backoff implementeren:

```csharp
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);
        }
    }
}
```

Dit bootst HTTP's 429 (Too Many Requests) patroon na. In plaats van de makelaar te hameren, pauzeren we voordat we het opnieuw proberen, waardoor het systeem tijd krijgt om te herstellen.

### Kanalen gebruiken voor In-Process Backpressure

Als je een interne pijpleiding bouwt (producent → processor → consument alle binnen uw toepassing), .NET's `Channel<T>` biedt elegante backpressure ondersteuning:

```csharp
// 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)
);
```

Het begrensde kanaal voert automatisch tegendruk uit de producent blokkeert wanneer het kanaal vol is, natuurlijk vertragend om het tempo van de consument aan te passen. Geen handmatige throttling vereist.

### Een complete Backpressure-Aware Publisher

Hier is een completer voorbeeld dat monitoring, bevestigingen en back-off samenbrengt:

```csharp
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();
    }
}
```

## Beste praktijken

### Rustig blijven.

Raak niet in paniek wanneer wachtrijen groeien. Een beetje diepte is normaal en gezond dus uw systeem is het absorberen van belasting pieken sierlijk. Het doel is niet een lege wachtrij; het is een *stabiel* wachtrij die niet ongebonden groeit.

Monitor wachtrijdiepte na verloop van tijd. Zoek naar trends, geen snapshots. Een wachtrij die constant 100 berichten bevat is prima. Een wachtrij die is gegroeid van 100 naar 10.000 het afgelopen uur heeft aandacht nodig.

### Pragmatisch zijn

Pas patronen pragmatisch toe. Niet elk bericht heeft een uitgever nodig. Niet elke wachtrij heeft een verfijnde backpressure handling nodig. Een wachtrij die 10 berichten per uur verwerkt heeft waarschijnlijk niet dezelfde veerkrachtstechniek nodig als één verwerking van 10.000 per seconde.

Vraag jezelf af: "Wat is de werkelijke kosten als dit bericht is verloren of vertraagd?" Als het antwoord is "niet veel," niet over-engineer. Als het antwoord is "significante financiële of gegevensintegriteit impact," investeren in de juiste backpressure handling.

### Slimme schalen

Wanneer wachtrijen back-up, het antwoord is niet altijd "toevoegen meer producenten." Dat is als proberen om een file te repareren door het toevoegen van meer auto's.

Overweeg:

- **Eerst de consument schalen** - kunt u meer werknemers toevoegen om de achterstand te verwerken?
- **Controle op knelpunten** - Is een langzame stroomafwaartse afhankelijkheid de oorzaak van de back-up?
- **Lot waar mogelijk** - kunnen consumenten meerdere berichten tegelijk verwerken?

```mermaid
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 en alarm

Instellen van signaleringen voor:

- Wachtrijdiepte boven drempels
- De vertraging van de consument neemt toe
- Uitgever bevestigt timing uit
- Verbinding geblokkeerde gebeurtenissen

U wilt weten over tegendruk *voor* Het wordt een crisis, niet als je systeem al is omgevallen.

## Conclusie

Backpressure is niet alleen het beperken van de snelheid is een overlevingstactiek. Door het te behandelen als een gesprek tussen producent en consument, bouw je systemen die veerkrachtig blijven onder druk.

De belangrijkste inzichten:

1. **Backpressure is feedback** - producenten en consumenten die samenwerken om duurzame doorvoer te vinden
2. **Monitor-wachtrijdiepte** - je kunt niet beheren wat je niet kunt meten
3. **Uitgever gebruiken bevestigt** - weet wanneer de makelaar worstelt
4. **Exponentiële back-off uitvoeren** - niet een systeem hameren dat al onder druk staat
5. **Schaal de consumenten, niet alleen de producenten** - repareer het bottleneck, niet het symptoom

Wanneer uw systeem zegt "Ik ben vol, geef me een minuut," de juiste reactie is "Geen probleem, ik zal wachten." Dat is de essentie van goed gedragen gedistribueerde systemen beschaafd, samenwerkend en veerkrachtig.

Blijf kalm als wachtrijen een beetje groeien, en onthoud: een systeem onder sierlijke tegendruk is oneindig beter dan een systeem dat volledig is omgevallen.