Mottryck är den osung hjälte för distribuerade system. Det är vad som hindrar dina köer från att brista i sömmarna när producenter skjuter meddelanden snabbare än konsumenterna kan tugga igenom dem. Enkelt uttryckt: det är systemet som säger "vänta på ett ögonblick" när saker och ting blir för upptagna.
Så här är det: teknikerna i denna artikel gäller praktiskt taget var och en av dessa meddelande kö och service buss-RabbitMQ, Kafka, Azure Service Bus, AWS SQS, NATS, du namnger det. Detaljerna varierar, men principerna är universella. Jag kommer att använda RabbitMQ för de flesta exempel eftersom det är vad jag vet bäst, men jag ska visa dig hur dessa mönster översätter över plattformar.
En bekännelse: Även de flesta seniora utvecklare inte implementera korrekt backpressure hantering. De bygger happy-path system som fungerar bra i dev och iscensättning, sedan undrar varför produktionen faller över under Black Friday. Mottryck hantering är en av de tekniker som skiljer "det fungerar" från "det skalor." Om du inte tänker på det, du bygger ett system som så småningom kommer att misslyckas under belastning.
I dess kärna är backpressure en återkoppling slinga som saktar producenterna när konsumenterna släpar efter. Tänk på det som trafikljus på en halkväg – du kan inte bara stapla på motorvägen när du vill. Ljusen styr flödet, låta bilar gå samman säkert utan att orsaka en hög-up.
Utan mottryck, en snabb producent kommer att överväldiga en långsam konsument. Meddelanden staplas upp i köer, minnet blir utmattad, och så småningom ditt system faller över. Backpressure säger "fast på" innan katastrofen slår till.
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
Skönheten med mottryck är att det är en samtal Mellan producent och konsument. Konsumenten signalerar "Jag är full, ge mig en minut" och producenten svarar "Inga problem, jag väntar." Det är artigt, samarbete, och hindrar alla från att falla över.
RabbitMQ har flera inbyggda mekanismer för hantering av mottryck, och förståelse för dem är avgörande om du bygger system som behöver hålla sig upprätt under belastning.
När RabbitMQ:s minnesanvändning eller ködjup överskrider inställda tröskelvärden aktiveras flödesreglering. Detta blockerar tillfälligt förlagsanslutningar – publicister kan inte skicka nya meddelanden förrän mäklaren har rensat tillräckligt med backlog.
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
Nyckelinsikten här är att RabbitMQ inte bara släpper meddelanden under tryck – det saktar ner källan. Detta är ett mycket mer civiliserat tillvägagångssätt än att tyst kasta data.
Konsumenter kontrollerar hastigheten genom erkännanden (ACKs och NACKs). Ett meddelande tas inte bort från kön förrän konsumenten uttryckligen erkänner det. Om en konsument inte ACK meddelanden snabbt nog, växer kön – som så småningom utlöser flödeskontroll uppströms.
Du kan också använda Gränser för förhämtning Att kontrollera hur många meddelanden som en konsument kan ha under flygning på en gång, vilket hindrar en enda långsam konsument från att lagra meddelanden.
// Set prefetch count to limit unacknowledged messages
channel.BasicQos(prefetchSize: 0, prefetchCount: 10, global: false);
Detta säger RabbitMQ: "Skicka mig bara 10 meddelanden åt gången. När jag väl ACK några, kan du skicka mer." Det är konsumenten som uttryckligen säger hur mycket tryck det kan hantera.
Mönstren är universella, men implementationerna skiljer sig åt. Här är hur några andra populära meddelandesystem närmar sig samma problem.
Kafka intar en helt annan hållning – konsumenterna dra Meddelanden snarare än att få dem knuffade. Detta gör mottryck implicit: om en konsument inte gör en enkät, tar den inte emot meddelanden. mäklaren bryr sig inte, det bara håller meddelanden runt tills konsumenten är redo.
// 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
}
Den smarta biten: Kafkas konsumentgrupper återbalanserar automatiskt partitioner. Om en konsument hamnar efter, kan du lägga till fler konsumenter till gruppen och partitioner omfördelas. Backpressure blir ett skalningsbeslut.
// 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 Buss använder en MaxConcurrentCalls inställning som är vackert enkel – den styr hur många meddelanden din processor hanterar samtidigt. Backpressure är automatiskt.
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 stöder också sessioner för beställd bearbetning och köer med dödbokstäver för meddelanden som misslyckas upprepade gånger – båda viktiga för att hantera påtryckningar när saker går fel.
SQS använder sikt timeouts som sin mottrycksmekanism. När du får ett meddelande blir det osynligt för andra konsumenter. Om du inte tar bort det i tid, dyker det upp igen för någon annan att prova.
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);
}
}
Det smarta SQS tricket: använd ApproximateNumberOfMessages För övervakning av ködjup och konsumenter i egen skala:
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 har explicit flödeskontroll med konsumentbekräftelser och max väntande meddelandegränser:
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());
Lägg märke till vad alla dessa system har gemensamt:
Syntaxen skiljer sig åt, men dansen är densamma: "Om jag inte berättar i tid, anta att jag misslyckades."
Här är praktiska exempel på hur du implementerar och reagerar på mottryck i dina C#-applikationer.
Först och främst – du kan inte hantera vad du inte kan mäta. Så här kontrollerar du hur många meddelanden som väntar i en kö:
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");
}
Den här snippan kontrollerar hur många meddelanden som väntar. Om räkningen klättrar, är det din signal att strypa producenter eller skala upp konsumenter. Oroa dig inte för att kontrollera detta alltför ofta – en periodisk hälsokontroll är vanligtvis tillräckligt.
Utgivare bekräftar att du vet när RabbitMQ framgångsrikt har tagit emot och behandlat ditt meddelande. Om bekräftelsen saktar ner, är det en tydlig signal om mottryck:
// 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");
}
Om RabbitMQ kämpar tar det längre tid eller tid helt och hållet. Din producent kan använda denna signal för att backa istället för att stapla på mer tryck.
För höggenomströmningsscenarier vill du ha asynkront bekräftande:
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
};
När du upptäcker mottryck, det värsta du kan göra är omedelbart försök i full fart. Det är som att svara på en trafikstockning genom att trycka på acceleratorn hårdare. Istället, genomföra exponentiellt backoff:
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);
}
}
}
Detta efterliknar HTTP: s 429 (för många förfrågningar) mönster. I stället för att hamra mäklaren, pausar vi innan försök, vilket ger systemet tid att återhämta sig.
Om du bygger en intern pipeline (producent → processor → konsument alla inom din ansökan), .NET: s Channel<T> ger elegant mottrycksstöd:
// 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)
);
Den avgränsade kanalen tillämpar automatiskt mottryck – producenten blockerar när kanalen är full, naturligtvis saktar ner för att matcha konsumentens takt. Ingen manuell strypning krävs.
Här är ett mer komplett exempel som sammanför övervakning, bekräftar, och backoff:
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();
}
}
Få inte panik när köerna växer. Lite djup är normalt och hälsosamt – det betyder att ditt system absorberar belastningstoppar graciöst. Målet är inte en tom kö; det är en stabil kö som inte växer oöverskådligt.
Övervaka ködjup över tid. Leta efter trender, inte ögonblicksbilder. En kö som är konsekvent på 100 meddelanden är bra. En kö som har vuxit från 100 till 10 000 under den senaste timmen behöver uppmärksamhet.
Tillämpa mönster pragmatiskt. Inte varje meddelande behöver utgivare bekräftar. Inte varje kö behöver sofistikerad mottryck hantering. En kö som behandlar 10 meddelanden per timme behöver förmodligen inte samma motståndskraft teknik som en bearbetning 10 000 per sekund.
Fråga dig själv: "Vad är den faktiska kostnaden om detta meddelande förloras eller försenas?" Om svaret är "inte mycket", inte överingenjör. Om svaret är "betydande ekonomisk eller data integritet inverkan," investera i korrekt backpressure hantering.
När köer backar upp, svaret är inte alltid "lägga till fler producenter." Det är som att försöka fixa en trafikstockning genom att lägga till fler bilar.
Tänk på följande:
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
Inrätta registreringar för
Du vill veta om mottryck före Det blir en kris, inte när ditt system redan är över.
Mottrycket är inte bara hastighetsbegränsning – det är en överlevnadstaktik. Genom att behandla det som ett samtal mellan producent och konsument bygger du system som håller sig motståndskraftiga under tryck.
De viktigaste insikterna:
När ditt system säger "Jag är full, ge mig en minut", är rätt svar "Inga problem, jag kommer att vänta." Det är kärnan i väluppfostrade distribuerade system-polit, samarbete, och motståndskraftig.
Håll dig lugn när köerna växer lite, och kom ihåg: ett system under graciöst mottryck är oändligt mycket bättre än ett system som fallit över helt och hållet.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.