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
الضغط على الظهر هو البطل غير المُعَلَّق للنظم الموزعة. هذا ما يمنع طوابيرك من الإنفجار عند الطوابير عندما يقوم المنتجون بإطلاق رسائل أسرع مما يستطيع المستهلكون مضغها من خلالها. ضع ببساطة: إنه النظام الذي يقول "انتظر لحظة" عندما تصبح الأمور مشغولة جداً.
هذا هو الشيء: تنطبق التقنيات الواردة في هذه المادة على كل كل قائمة الرسائل والحافلة الخدمية - RabbitMQ, Kafka, Azure Service Service Bus, AWS SQS, NATS, سمها. التفاصيل تختلف، لكن المبادئ عالمية. سأستخدم AnterMQ لمعظم الأمثلة لأنها ما أعرفه أفضل، لكنني سأريكم كيف تترجم هذه الأنماط عبر المنصات.
(أ) الاعتراف: حتى معظم كبار المطورين لا ينفذون مناولة ملائمة للضغط على الظهر. إنهم يبنون أنظمة سعيدة تعمل بشكل جيد في التنضيم والتنشيط، ثم يتساءلون لماذا ينهار الإنتاج أثناء يوم الجمعة الأسود. معالجة الضغط على الظهر هي واحدة من التقنيات التي تفصل بين "تعمل" و"تعمل" و"تدرجات". إذا كنت لا تفكر في ذلك، فأنت تبني نظاماً سيفشل في نهاية المطاف تحت الحمل.
في جوهره، الضغط على الظهر هو حلقة تغذية مرتدة تبطئ المنتجين عندما يتخلف المستهلكون عن الركب. فكّر فيها مثل إشارات المرور على طريق زلة - لا يمكنك فقط أن تتراكم على الطريق السريع متى ما رغبت. الأضواء تتحكم في التدفق، تسمح للسيارات بالاندماج بأمان دون التسبب في تكدس.
من دون الضغط على الظهر، فإن المنتج السريع سوف يتغلب على المستهلك البطيء. الرسائل تتراكم في الطابور، الذاكرة تُستنفد، وفي النهاية ينهار نظامك. ضغط الظهر يقول "مستقر" قبل وقوع الكارثة.
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 لديه عدة آليات مدمجة لمعالجة الضغط الظهري، وفهمها أمر حاسم إذا كنت تبني أنظمة تحتاج إلى البقاء مستقيمة تحت الحمل. وثائق منسِفQ ممتاز - سأربط صفحات محددة بينما نحن نسير.
عندما يقوم استخدام الذاكرة أو عمق الطابير الذي يزيد عن عمقه على استخدام الذاكرة أو العتبات المُعَدّة، فإنه يُشغّل (الإنفاق)هذا مؤقتًا يمنع روابط الناشرين - لا يمكن للناشرين إرسال رسائل جديدة حتى يقوم السمسار بتصفية ما يكفي من الأعمال المتراكمة. وقد عقد مؤتمراً بشأن الذي يُؤدّي إلى التحكّم في التصرّف.
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
البصيرة الرئيسية هنا هي أن AnbMQ لا يُسقط الرسائل فقط عندما تحت الضغط - بل يبطئ المصدر. هذا نهج حضاري أكثر بكثير من التخلص من البيانات بصمت.
المستهلكون يتحكمون في الوتيرة من خلال الإقرارات الرسالة لا تُزال من الصف حتى يعترف المستهلك بها صراحة. إذا لم يقم المستهلك بتوصيل الرسائل بسرعة كافية، فإن الطابور ينمو مما يؤدي في نهاية المطاف إلى السيطرة على التدفق إلى أعلى.
يمكنك أيضاً استخدام مُكَفِف حدود (QOS) للتحكم في عدد الرسائل غير المعترف بها التي يمكن للمستهلك الحصول عليها أثناء الطيران في وقت واحد. وهذا يمنع المستهلك البطيء الوحيد من تخزين الرسائل.
// Set prefetch count to limit unacknowledged messages
channel.BasicQos(prefetchSize: 0, prefetchCount: 10, global: false);
يقول هذا AnterMQ: "فقط أرسل لي 10 رسائل في كل مرة. بمجرد أن أقوم بطرح بعضها، يمكنك إرسال المزيد." هو المستهلك الذي يقول صراحة كم من الضغط يمكنه التعامل معه. الـ وثائق العملاء ويغطي هذا المؤشر بالتفصيل.
الأنماط عالمية، لكن التنفيذ يختلف. وهنا كيف أن بعض أنظمة الرسائل الشعبية الأخرى تتعامل مع نفس المشكلة.
يأخذ كافكا نهجاً مختلفاً جوهرياً - نهجاً - مستهلكاً & حرّر بدلاً من أن يدفعوها. هذا يجعل الضغط على الظهر ضمنياً: إن لم يصوت المستهلك، فإنه لا يتلقى رسائل. السمسار لا يهتم. إنه يحفظ الرسائل حتى يكون المستهلك مستعداً.
// 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();
كما يقدم الدعم لحافز خدمات الأززززز ما يُجرِّه من دورات طلب من الأمين العام ** للرسائل التي تفشل مراراً وتكراراً- كلاهما مهم لإدارة الضغط عندما تسوء الأمور.
SQS يستخدم وقت الرؤية كآلية ضغط خلفية. عندما تستقبل رسالة، تصبح غير مرئية للمستهلكين الآخرين. إذا لم تحذفها في الوقت، فإنها تظهر لشخص آخر ليجربها.
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);
}
}
الخدعة الذكية SQS SQS: استخدام: 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 JetStream لديه مراقبة واضحة للتدفق مع إقرارات المستهلك وأقصى حدود معلّقة للرسائل:
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");
}
هذه القصاصة تتحقق من عدد الرسائل التي تنتظر. إذا صعد العد، هذه هي إشارتك لخنق المنتجين أو رفع مستوى المستهلكين. لا تقلق بشأن التحقق من هذا بشكل متكرر جداً - الفحص الصحي الدوري عادة ما يكون كافياً.
الناشر أكّد إلى علم متى AnbMQ نجح في إستلام و معالجة الرسالة. إذا تباطأت الإقرارات، هذه إشارة واضحة من ضغط الظهر:
// 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");
}
إذا كان AnbMQ يكافح، فإن التأكيدات تستغرق وقتاً أطول أو وقتاً بعيداً تماماً. يمكن لمنتجك استخدام هذه الإشارة للتراجع بدلاً من التثبيت على المزيد من الضغط.
بالنسبة للسيناريوهات عالية الفعالية، سوف تريد تأكيدات شبه متزامنة:
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);
}
}
}
هذا يحاكي نمط HTTTP 429 (الكثير من الطلبات). بدلاً من ضرب السمسار، نتوقف قبل إعادة المحاولة، نعطي النظام وقت للاستعادة.
إذا كنت تبني خط أنابيب داخلي (المنتج o المعالج 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 إلى 10000 على مدار الساعة الماضية يحتاج إلى اهتمام.
تطبيق أنماط عمليّة. ليس كل رسالة تحتاج إلى تأكيدات من الناشر. ليس كل طابور بحاجة إلى معالجة متطورة للضغط الظهري. a طابور يقوم بتجهيز 10 رسائل في الساعة على الأرجح لا يحتاج إلى نفس هندسة القدرة على التحمل مثل معالجة 10,000 في الثانية.
اسأل نفسك: "ما هي التكلفة الفعلية إذا كانت هذه الرسالة مفقودة أو متأخرة؟" إذا كانت الإجابة "ليس الكثير"، لا تفرط في الهندسة. إذا كانت الإجابة هي "تأثير مالي ذي شأن أو تأثير سلامة البيانات"، استثمر في معالجة الضغط الخلفي السليم.
عندما تتراجع الطوابير، الإجابة ليست دائماً "تزيد من المنتجين". هذا أشبه بمحاولة إصلاح ازدحام مروري بإضافة المزيد من السيارات.
ما يلي:
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
إنشاء تنبيهات لـ:
تريد أن تعرف عن الضغط على الظهر ثالثاً يُصبحُ a أزمة، لَيسَ عندما نظامكَ يَسْقطُ.
الضغط على الظهر ليس فقط الحد من المعدل، بل هو وسيلة للبقاء على قيد الحياة. من خلال معاملتها على أنها محادثة بين المنتج والمستهلك، يمكنك بناء الأنظمة التي تبقى مرنة تحت الضغط.
الرؤى الرئيسية:
عندما يقول نظامك "أنا كامل، أعطني دقيقة،" الرد الصحيح هو "لا مشكلة، سأنتظر." هذا هو جوهر الأنظمة الموزعة المهذبة جيداً - سياسة، تعاونية، ومرنة.
إبقوا هادئين عندما تنمو الطوابير قليلاً، وتذكروا: النظام تحت ضغط الظهر الرشيق أفضل إلى الأبد من النظام الذي سقط تماماً.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.