This is a viewer only at the moment see the article on how this works.
To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk
This is a preview from the server running through my markdig pipeline
Sunday, 23 November 2025
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.
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.
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.
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 Ik zal link naar specifieke pagina's als we gaan.
Wanneer RabbitMQ's geheugengebruik of wachtrijdiepte de geconfigureerde drempels overschrijdt, activeert het stroomregeling. Dit blokkeert tijdelijk publisher connecties . Uitgevers kunnen geen nieuwe berichten verzenden totdat de makelaar voldoende achterstand heeft opgeruimd . Zie ook de documenten op geheugenalarmen en schijfalarmen welke trigger flow control.
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.
Consumenten controleren het tempo door bevestigingen (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 (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.
// 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 behandelt de API in detail.
De patronen zijn universeel, maar de implementaties verschillen. Hier is hoe sommige andere populaire messaging systemen benaderen hetzelfde probleem.
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.
// 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.
// 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 maakt gebruik van een MaxConcurrentCalls instelling dat is prachtig eenvoudig . Controleert hoeveel berichten uw processor tegelijkertijd behandelt. Backpressure is automatisch.
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.
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.
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:
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 heeft expliciete stroomcontrole met consumentenbevestigingen en maximale berichtlimieten:
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());
Merk op wat al deze systemen gemeen hebben:
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."
Hier zijn praktische voorbeelden van het implementeren en reageren op tegendruk in je C#-toepassingen.
Eerst kunt u niet beheren wat u niet kunt meten. Hier is hoe u kunt controleren hoeveel berichten wachten in een wachtrij:
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 laat u weten wanneer RabbitMQ uw bericht succesvol heeft ontvangen en verwerkt. Als bevestigingen vertragen, is dat een duidelijk signaal van tegendruk:
// 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:
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
};
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:
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.
Als je een interne pijpleiding bouwt (producent → processor → consument alle binnen uw toepassing), .NET's Channel<T> biedt elegante backpressure ondersteuning:
// 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.
Hier is een completer voorbeeld dat monitoring, bevestigingen en back-off samenbrengt:
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();
}
}
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.
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.
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:
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
Instellen van signaleringen voor:
U wilt weten over tegendruk voor Het wordt een crisis, niet als je systeem al is omgevallen.
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:
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.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.