पीछेवाले नायकों को वितरित तंत्रों के अनंग हीरो है. यह है कि जब निर्माता उन लोगों के माध्यम से प्राप्त करने से तेजी से संदेशों को तेजी से उड़ा रहे हैं तो अपने कतारों को जारी रखता है. बस कहना है: यह प्रणाली "एक पल पर पकड़" कहते हैं जब स्थिति बहुत व्यस्त हो.
यहाँ बात है: इस लेख की तकनीकों पर लगभग लागू होता है प्रत्येक संदेश कतार और सेवा बस
एक पाप - स्वीकृति: यहाँ तक कि अधिकांश वयस्क विकासकर्ता भी उचित गति से काम नहीं कर रहे हैं। वे खुश पथ प्रणाली बनाते हैं कि अच्छा काम करते हैं, तो आश्चर्य है कि क्यों उत्पादन ब्लैक शुक्रवार के दौरान गिर जाता है। बैकर नियंत्रण उन तकनीकों में से एक है कि "यह काम करता है" से अलग है। यदि आप इसके बारे में सोच रहे हैं, आप एक इमारत है कि अंततः लोड हो जाएगा।
अपने कोर में, पीछे वाले एक फ़ीडबैक लूप है कि जब उपभोक्ता पीछे रह जाता है. इसके बारे में सोचो एक फिसली सड़क पर यातायात रोशनी - आप सिर्फ कार मार्ग पर ही नहीं मिल सकता जब आप कल्पना कर सकते हैं. प्रकाश प्रवाह को नियंत्रित करते हैं, कारों को सुरक्षित रूप से एक ढेर बनाने के बिना.
बिना तेजी से, एक तेजी से निर्माता एक धीमा उपभोक्ता को कुचल देगा. संदेश कतार में ढेर, स्मृति समाप्त हो जाती है, और बाद में आपका सिस्टम गिर जाता है. वापस दबाव का कहना है कि विपत्ति आने से पहले "पर" कहते हैं.
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
पीछे हटने की खूबसूरती यह है कि यह एक है वार्तालाप "मैं पूरा हूँ, मुझे एक मिनट दो" और निर्माता जवाब देता है "कोई समस्या नहीं, मैं इंतजार करेंगे." यह शांत, सहयोग से, और हर किसी के गिरने से दूर रहता है.
खरगोशQ के पास बहुत से प्रक्रिया के लिए बनाया गया है... ... पीठ थर्मेशन के लिए, और उन्हें समझने के लिए महत्वपूर्ण है उन्हें समझने के लिए है अगर आप निर्माण व्यवस्थाओं को तैयार कर रहे हैं जो बोझ तले सीधा रहने की जरूरत है. खरगोशएम दस्तावेज़Comment बढ़िया है, मैं विशेष पृष्ठों को लिंक के रूप में हम जाने के रूप में।
जब रब्बीटएम स्मृति प्रयोग या कतार गहराई अधिक से अधिक कॉन्फ़िगर किए गए सीमाों को सक्षम करता है, यह सक्रिय करता है फ्लो कंट्रोल. यह अस्थायी ब्लॉक कनेक्शन - सेब्स नए संदेश नहीं भेज सकता जब तक कि दलाल पर्याप्त बैकलॉग को मंजूर नहीं कर देता. स्मरण अलार्म@ info: whatsthis और डिस्क अलार्म@ info: whatsthis जो (कुफ्फ़ार की रूह) डूब कर सख्ती से खींच लेते हैं
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
महत्वपूर्ण अंतर्दृष्टि यहाँ है कि खरगोश MQ सिर्फ जब दबाव के तहत संदेश छोड़ नहीं करता है - यह स्रोत को धीमा कर देता है. यह चुपचाप डेटा बंद करने से एक बहुत अधिक नागरिक तरीका है.
समाचार - पत्र गति को गति से नियंत्रित करते हैं अनुमोदन (csss और Nccs). एक संदेश कतार से हटा नहीं है जब तक उपभोक्ता स्पष्ट रूप से यह स्वीकार नहीं कर लेते. यदि एक उपभोक्ता पर्याप्त तेज नहीं करता है, कतार बढ़ता है, जो अंततः sons प्रवाहित करता है.
आप भी उपयोग कर सकते हैं प्रीफेच सीमा (QoS) नियंत्रण करने के लिए कि कितने ज्ञानहीन संदेशों को एक बार में उड़ सकते हैं. यह एक एकल धीमा उपभोक्ता संदेश इकट्ठा करने से रोकता है.
// Set prefetch count to limit unacknowledged messages
channel.BasicQos(prefetchSize: 0, prefetchCount: 10, global: false);
यह खरगोशQ से कहता है: "सिर्फ मुझे 10 संदेश भेजने के लिए एक समय पर. एक बार जब मैं ऊपर कुछ उठा, आप और भेज सकते हैं." यह उपभोक्ता स्पष्ट रूप से कह सकते हैं कि यह कितना दबाव कर सकते हैं. .नेटटी क्लाएंट दस्तावेज़ और उस वक्त (क़ी हालत) पड़े देखा करते हैं
पैटर्न मौजूद हैं, लेकिन कार्यान्वयन अलग है. यहाँ कैसे कुछ लोकप्रिय संदेश व्यवस्थाओं एक ही समस्या के पास आते हैं.
काफ़्का एक मूलभूत दृष्टिकोण लेता है, दबायें संदेशों को धक्का देने के बजाय। यह वापस supdmundssssypedssssssssss. अगर एक उपभोक्ता को नहीं मिलता है, यह संदेश प्राप्त नहीं करता है. विक्रेता परवाह नहीं करता है, यह सिर्फ चारों ओर रहते हैं जब तक उपभोक्ता तैयार नहीं है.
// 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
}
चतुर बिट: काफ़काकाका के उपभोक्ता समूह स्वचालित रूप से पुनः लोड होता है. यदि एक उपभोक्ता पीछे गिर जाता है, तो आप समूह के लिए और उपभोक्ताओं को और अधिक उपभोक्ताओं को वितरित कर सकते हैं. डायटिंग एक स्केलिंग निर्णय बन जाता है.
// 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
};
एमिल सेवा बस एक प्रयोग करता है MaxConcurrentCalls यह बहुत सरल है, यह नियंत्रण करता है कि कितनी सारे संदेश आपके प्रक्रिया को एक साथ संभालता है. वापस आने वाले स्वचालित है.
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();
एमिल सेवा बस भी समर्थन करता है सत्र निर्देशित किया गया मृत-लेट कतार संदेशों के लिए जो बारंबार असफल हो जाते हैं - जब मामले गलत हो जाते हैं तब दबाव का प्रबंधन करने के लिए महत्वपूर्ण हैं.
SQQS का उपयोग करता है समय की गति का उपयोग अपनी पीठ कमजोर प्रणाली के रूप में. जब आप एक संदेश मिलता है, यह अन्य उपभोक्ताओं के लिए अदृश्य हो जाता है. अगर आप समय में इसे मिटा नहीं देते हैं, यह किसी और की कोशिश करने के लिए फिर से शुरू करता है.
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);
}
}
चतुर SQECKICKICKS चाल: इस्तेमाल करें ApproximateNumberOfMessages कतार गहराई को मॉनीटर करने के लिए तथा स्वचलित सफेद उपभोक्ताओं को निगरानी में रखें:
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 Jetetiz ने उपभोक्ताों को स्वीकार करने और अधिकतम स्थगित संदेश सीमा के साथ सुस्पष्ट प्रवाह नियंत्रण किया है:
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());
ध्यान दीजिए कि इन सभी व्यवस्थाओं में क्या समानता है:
वाक्य भिन्न है, लेकिन नृत्य एक समान है: "यहाँ मैं कितना संभाल सकते हैं. मुझे बताओ जब मैं इसे संभाल लिया है. अगर मैं आपको समय में नहीं बताता, मान लीजिए कि मैं असफल हो गया. "
ठीक है, चलो कोड में मिलता है. यहाँ पर लागू करने के व्यावहारिक उदाहरण हैं अपने C# अनुप्रयोगों में वापस उठने के लिए।
पहली चीजें हैं - आप क्या माप नहीं दे सकते. यहाँ कैसे पता लगाने के लिए कि कितने संदेश कतार में इंतजार कर रहे हैं:
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");
}
यह जाँच कर रहे हैं कि कितने संदेश इंतजार कर रहे हैं. यदि गिनती आती है, तो यह आपके क्यूरों के लिए है जो fretting निर्माता या पैमाने पर उपभोक्ताओं को. इस तरह की जाँच करने के बारे में चिंता मत करो, एक आवधिक स्वास्थ्य जाँच आम तौर पर काफी है.
प्रकाशक पुष्टि करते हैं कि जब रब्बीएएमQ सफलता पूर्वक प्राप्त किया गया है और आपके संदेश को प्रस्तुत किया है. यदि इनकार करना धीमा हो, तो यह वापसी वापसी का एक स्पष्ट संकेत है:
// 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");
}
यदि रब्बीतQ संघर्ष कर रहा है, पुष्टि पूरी तरह से या समय बाहर ले जाने के लिए पूरी तरह से समय लगता है ।
उच्च रूप से उदाहरणों के लिए, आप अतुल्यकालिक रूप से पुष्टि करेगा:
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
};
जब आप वापस दबाव का पता लगाएँ, तो सबसे बुरी बात जो आप कर सकते हैं पूर्ण गति पर तुरंत फिर से कोशिश. यह है कि एक ट्रैफिक जाम को दबाकर कठोर. इसके बजाय, सफल वापस बंद करें:
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);
}
}
}
यह नरम HTTP9 के 429 (o कई निवेदनों) पैटर्न। दलाल को बनाने के बजाय, हम प्रयास करने से पहले, सिस्टम को पुनः प्राप्त करने के लिए समय दे।
यदि आप एक आंतरिक स्पिट्यूट बना रहे हैं (अपने अनुप्रयोग के भीतर सभी प्रकार के उपभोक्ता, .नेटटी' Channel<T> साथ ही, वे दूसरों की मदद भी करते हैं:
// 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)
);
बाध्य्ड चैनल स्वतः ही वापस लागू होता है - निर्माता ब्लॉक जब चैनल पूरा होता है, स्वाभाविक रूप से उपभोक्ता की गति से मैच करने में धीमी होती है. कोई वैध विकल्प नहीं.
यहाँ एक और पूर्ण उदाहरण है जो एक साथ जांच लाता है, पुष्टि करता है, और बैकऑफ:
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();
}
}
कतार के बढ़ने पर मत घबरा. गहराई की एक बिट सामान्य है और स्वस्थ है - इसका मतलब है कि आपके तंत्र को दिलचस्प लोड है. लक्ष्य एक खाली कतार नहीं है, यह एक है. स्थिर कतार जो तेजी से विकसित नहीं होता है.
निगरानी कतार गहराई समय पर है. खेलों के लिए देखो, नहीं. एक कतार जो लगातार 100 संदेश में ठीक है. एक कतार है कि पिछले घंटे में 100 से 10,000 से अधिक घंटे ध्यान की आवश्यकता है.
पैटर्न प्रति सेकंड के रूप में लागू करें। हर संदेश प्रचारक की पुष्टि करने की जरूरत नहीं है। हर कतार में गंभीर रूप से वापस लाने की जरूरत होती है। एक कतार है कि एक घंटे में 10 संदेश की प्रक्रिया की आवश्यकता नहीं है। शायद एक ही mapultion की आवश्यकता नहीं है प्रति सेकंड के रूप में एक सेकंड के रूप में दस सेकंड के लिए।
अपने आप से पूछिए: "अगर यह संदेश खो गया है या देर हो गई है?" यदि उत्तर "बहुत अधिक नहीं," नहीं के ऊपर नहीं है. यदि जवाब है, तो उत्तर "यदि जवाब है "प्रिक वित्तीय वित्तीय या डाटा अखंडता प्रभाव" सही वापसी नियंत्रण में निवेश करता है.
जब कतारें समर्थन कर रही हैं, जवाब हमेशा "और अधिक निर्माताओं" नहीं है कि अधिक कारों को जोड़ने के द्वारा एक ट्रैफिक जाम को ठीक करने की कोशिश करना चाहते हैं.
विचार कीजिए:
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
इसके लिए चेतावनी सेट करें:
आप वापस बैठने के बारे में जानना चाहते हैं पहले यह एक संकट हो जाता है, नहीं जब अपने तंत्र के पहले से गिर चुका है.
बैकर केवल दर सीमित नहीं है - यह एक जीवित चाल है. इसे निर्माता और उपभोक्ताओं के बीच एक बातचीत के रूप में व्यवहार करने के द्वारा, आप उन सिस्टमों को निर्माण करते हैं जो दबाव के तहत निष्क्रिय रहते हैं.
मुख्य अन्तर्दृष्टि:
जब आपका तंत्र कहता है कि "मैं पूरा हूँ, मुझे एक मिनट दे" सही जवाब है "कोई समस्या नहीं, मैं इंतजार करूँगा." यह अच्छी तरह से वितरित सिस्टमों का सार है - विदेशी, सहयोग, और मजबूत.
शांत रहो जब कतार बन जाता है, और याद रखें: एक प्रणाली के नीचे सुंदर पीठ के नीचे एक प्रणाली बहुत ही बेहतर है कि पूरी तरह से गिर गया है.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.