Back to "Paineen alla: Miten jonotusjärjestelmät kestävät vastapainetta esimerkeillä C#"

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

Azure Service Bus C# Distributed Systems Kafka Messaging RabbitMQ

Paineen alla: Miten jonotusjärjestelmät kestävät vastapainetta esimerkeillä C#

Sunday, 23 November 2025

Vastapaine on hajautettujen järjestelmien loputon sankari. Se estää jonoja puhkeamasta saumoihin, kun tuottajat ampuvat viestejä nopeammin kuin kuluttajat voivat pureskella niitä läpi. Yksinkertaisesti sanottuna: se on järjestelmä, jossa sanotaan "odota hetki", kun asiat menevät liian kiireisiksi.

Juttu on näin: tämän artikkelin tekniikat pätevät käytännössä joka Viestijono ja palvelubussi – RabbitMQ, Kafka, Azure Service Bus, AWS SQS, NATS, voit sanoa sen. Yksityiskohdat vaihtelevat, mutta periaatteet ovat yleismaailmallisia. Käytän RabbitMQ:ta useimpiin esimerkkeihin, koska se on se, minkä tiedän parhaiten, mutta näytän, miten nämä kuviot muuttuvat eri alustoilla.

Tunnustus: Edes useimmat vanhemmat kehittäjät eivät toteuta asianmukaista vastapaineen käsittelyä. He rakentavat hyviä dev- ja lavastusjärjestelmiä, sitten ihmettelevät, miksi tuotanto kaatuu Black Fridayn aikana. Vastapaineen käsittely on yksi niistä tekniikoista, jotka erottavat "se toimii" ja "se vaa'at". Jos et ajattele sitä, rakennat järjestelmän, joka lopulta epäonnistuu kuormituksessa.

Mikä on vastapainetta

Pohjimmiltaan vastapaine on takaisinkytkentä, joka hidastaa tuottajia, kun kuluttajat ovat jäljessä. Ajattele sitä liikennevaloina liukutiellä – moottoritielle ei voi vain kasautua milloin haluaa. Valot hallitsevat virtaa, jolloin autot voivat fuusioitua turvallisesti aiheuttamatta kolaria.

Ilman vastapaineita nopea tuottaja yllättää hitaan kuluttajan. Viestit kasaantuvat jonoihin, muisti uupuu, ja lopulta järjestelmäsi kaatuu. Vastapaineen mukaan "vakaa" ennen kuin katastrofi iskee.

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

Vastapaineen kauneus on se, että se on keskustelu Kuluttaja viestittää: "Olen täynnä, anna minulle minuutti" ja tuottaja vastaa: "Ei ongelmaa, minä odotan." Se on kohteliasta, yhteistyöhaluista ja estää kaikkia kaatumasta.

Miten RabbitMQ käsittelee vastapainetta

RabbitMQ:lla on useita sisäänrakennettuja mekanismeja vastapaineen hallitsemiseksi, ja niiden ymmärtäminen on ratkaisevan tärkeää, jos rakennetaan järjestelmiä, joiden on pysyttävä pystyssä kuormituksen alla.

Virtauksen ohjaus

Kun RabbitMQ:n muistinkäyttö tai jonon syvyys ylittää määritetyt kynnykset, se aktivoituu Virtauksen ohjausTämä estää tilapäisesti kustantajan yhteydet – julkaisijat eivät voi lähettää uusia viestejä ennen kuin välittäjä on selvittänyt tarpeeksi ruuhkaa.

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

Avainymmärrys tässä on, että RabbitMQ ei vain pudota viestejä paineen alla – se hidastaa lähdettä. Tämä on paljon sivistyneempi lähestymistapa kuin tietojen hiljainen poisheittäminen.

Kuluttajan suosionosoitukset

Kuluttajat hallitsevat vauhtia tunnustamalla (asennukset ja NACKit). Viestiä ei poisteta jonosta, ennen kuin kuluttaja nimenomaisesti tunnustaa sen. Jos kuluttaja ei paukuta viestejä tarpeeksi nopeasti, jono kasvaa, mikä lopulta käynnistää virtauksen hallinnan ylävirtaan.

Voit myös käyttää prefektirajat Tämä estää yhden hitaan kuluttajan hamstraamasta viestejä.

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

Tämä kertoo RabbitMQ:lle: "Lähetä minulle vain 10 viestiä kerrallaan. Kun paukutan joitakin, voit lähettää lisää." Kuluttaja sanoo selvästi, kuinka paljon painetta se kestää.

Miten muut järjestelmät kestävät vastapainetta

Mallit ovat universaaleja, mutta toteutustavat eroavat toisistaan. Näin jotkut muut suositut viestijärjestelmät lähestyvät samaa ongelmaa.

Kafka: Kuluttajaohjattu mielipidemittaus

Kafka omaksuu perinpohjaisesti erilaisen lähestymistavan – kuluttajat vedä Tämä tekee vastapaineesta epäsuoran: jos kuluttaja ei tee kyselyä, hän ei saa viestejä. Välittäjä ei välitä, se vain pitää viestit ympärillään, kunnes kuluttaja on valmis.

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

Näppärää on se, että Kafkan kuluttajaryhmät tasapainottavat osioita automaattisesti uudelleen. Jos yksi kuluttaja jää jälkeen, voit lisätä ryhmään lisää kuluttajia ja jakoja jaetaan uudelleen. Vastapaineesta tulee skaalaava päätös.

// 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-palvelubussi: Rinnakkaisviestin ohjaus

Azure Service Bus käyttää MaxConcurrentCalls Asettaminen on kauniisti yksinkertaista – se säätelee kuinka monta viestiä prosessori käsittelee samanaikaisesti. Vastapaine on automaattinen.

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 tukee myös istunnot tilattavalle käsittelylle ja dead-letter -jonoja Viestit, jotka epäonnistuvat toistuvasti – molemmat ovat tärkeitä paineen hallitsemiseksi, kun asiat menevät pieleen.

AWS SQS: Näkyvyysaikalisätanssi

SQS käyttää vastapainemekanisminaan näkyvyyden aikakatkaisuja. Kun saat viestin, se muuttuu näkymättömäksi muille kuluttajille. Jos et poista sitä ajoissa, se ilmestyy uudelleen jonkun muun yrittämään.

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

Näppärä SQS-temppu: käytä ApproximateNumberOfMessages Seurataan jonojen syvyyttä ja automaattisia mittauksia käyttäviä kuluttajia:

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 Sisäänrakennettu

NATS JetStreamilla on selkeä virtauksen ohjaus, jossa on kuluttajatunnuksia ja max odottavien viestien raja-arvot:

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

Tavanomaiset säikeet

Huomaa, mitä yhteistä kaikilla näillä järjestelmillä on:

  1. Eksplisiittiset erä- ja valuuttarajat - hallitset, kuinka paljon olet valmis käsittelemään
  2. Rekrytointiin perustuva virta - viestit pysyvät saatavilla, kunnes varmistat käsittelyn
  3. Aikalisäpohjainen toipuminen - jos epäonnistut, viestit palaavat yritettäväksi uudelleen
  4. Syvyyden seuranta Voit aina kysyä: "Kuinka tukossa olen?"

Syntaksi on erilainen, mutta tanssi on sama: "Näin paljon kestän. Kerro, kun olen hoitanut asian. Jos en kerro ajoissa, oletan epäonnistuneeni."

C#-koodiesimerkkejä

Tässä käytännön esimerkkejä C#-sovellusten vastapaineiden toteuttamisesta ja niihin reagoimisesta.

Kytkentäsyvyyden seuranta

Tärkeimmät asiat ensin – ei voi hallita sitä, mitä ei pysty mittaamaan. Näin voit tarkistaa, kuinka monta viestiä jonossa odottaa:

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

Tämä nappula tarkistaa, kuinka monta viestiä odottaa. Jos luku nousee, se on merkkisi tuottajien kuristamiseen tai kuluttajien skaalautumiseen. Älä huolehdi siitä, että tarkistat tämän liian usein – säännöllinen terveystarkastus on yleensä riittävä.

Julkaisija vahvistaa

Julkaisija vahvistaa, että ilmoita, kun RabbitMQ on onnistuneesti saanut ja käsitellyt viestisi. Jos vahvistukset hidastuvat, se on selkeä vastapaineen merkki:

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

Jos RabbitMQ kamppailee, vahvistukset vievät aikaa tai aikaa kokonaan. Tuottajasi voi tämän signaalin avulla perääntyä eikä kasata lisää paineita.

Korkeatehoisten skenaarioiden osalta haluat, että asynkroninen vahvistus on:

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

Yritä uudelleen eksponentiaalisella backoffilla

Kun havaitset vastapaineen, pahinta on yrittää heti uudelleen täydellä teholla. Se on kuin vastaisi liikenneruuhkiin painamalla kiihdytintä kovemmin. Sen sijaan toteuta eksponentiaalinen peräänajo:

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

Tämä jäljittelee HTTP:n 429 (Liian monta pyyntöä) -mallia. Sen sijaan, että moukaroisimme välittäjää, keskeytämme ennen kuin yritämme uudelleen, jolloin järjestelmä ehtii toipua.

Kanavien käyttö prosessorien vastapaineessa

Jos rakennat sisäistä putkistoa (tuottaja → prosessori → kuluttaja kaikki sovelluksessasi), .NET's Channel<T> tarjoaa tyylikkään vastapainetuen:

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

Rajattu kanava käyttää automaattisesti vastapainetta – tuottaja lohkoa, kun kanava on täynnä, luonnollisesti hidastumassa kuluttajan tahdin mukaan. Käsikäyttöistä kuristamista ei tarvita.

Täydellinen vastapaine-tietoinen julkaisija

Tässä on täydellisempi esimerkki, joka yhdistää seurannan, vahvistaa ja peruuttaa:

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

Parhaita käytäntöjä

Pysy rauhallisena

Älä joudu paniikkiin jonojen kasvaessa. Syvyys on normaalia ja terveellistä – se tarkoittaa, että järjestelmäsi imee kuormapiikit sulavasti. Tavoite ei ole tyhjä jono, se on vakaa Jono, joka ei kasva rajattomasti.

Seuraa jonon syvyyttä ajan mittaan. Etsi trendejä, älä kuvakuvia. 100 viestin jono on hyvä, ja jono, joka on kasvanut 100:sta 10 000:een viimeisen tunnin aikana, vaatii huomiota.

Ole pragmaattinen

Käytä kuvioita käytännönläheisesti. Jokainen viesti ei tarvitse kustantajan vahvistusta. Jokainen jono ei tarvitse pitkälle kehitettyä vastapaineen käsittelyä. 10 viestiä tunnissa prosessoiva jono ei todennäköisesti tarvitse samaa sietokykyä kuin yksi käsittely 10 000 sekunnissa.

Kysy itseltäsi: "Mitä maksaa, jos tämä viesti katoaa tai viivästyy?" Jos vastaus on "ei paljon", älä yli-insinööri. Jos vastaus on "merkittävä taloudellinen tai tietojen eheysvaikutus", investoi asianmukaiseen vastapaineen käsittelyyn.

Scale Smart

Kun jonot perääntyvät, vastaus ei aina ole "lisää tuottajia". Se on kuin yrittäisi korjata ruuhkaa lisäämällä autoja.

Harkitsehan:

  • Mittaa kuluttajat ensin - Voitko lisätä työntekijöitä ruuhkan käsittelyyn?
  • Pullonkaulojen tarkastaminen Aiheuttaako yksi hidas riippuvuus tukijoukkoja?
  • Erä mahdollisuuksien mukaan - voivatko kuluttajat käsitellä useita viestejä kerralla?
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

Seuraa ja hälytä

Asetetaan kuulutukset:

  • Kynnyksen syvyys ylittää kynnykset
  • Kuluttajaviive kasvaa
  • Julkaisija vahvistaa ajoituksen
  • Yhteys esti tapahtumat

Haluat tietää vastapaineesta ennen siitä tulee kriisi, ei silloin, kun järjestelmä on jo kaatunut.

Päätelmät

Vastapaine ei ole vain kurssirajoitus, vaan se on selviytymistaktiikka. Kohtelemalla sitä tuottajan ja kuluttajan välisenä keskusteluna rakennat järjestelmiä, jotka pysyvät kestävinä paineen alla.

Tärkeimmät oivallukset:

  1. Vastapaine on palautetta - tuottajat ja kuluttajat tekevät yhteistyötä kestävän läpimenon löytämiseksi
  2. Seuraa jonon syvyyttä - et pysty siihen, mitä et pysty mittaamaan
  3. Käytä julkaisijaa vahvistaaksesi - tietää, milloin meklari kamppailee
  4. Toteuta eksponentiaalinen perääntyminen - Älä moukaroi järjestelmää, joka on jo paineen alla
  5. Mittakaavakuluttajat, eivät vain tuottajat - korjaa pullonkaula, ei oiretta

Kun järjestelmäsi sanoo "olen täynnä, anna minulle minuutti", oikea vastaus on "Ei ongelmaa, minä odotan". Se on hyvin käyttäytyvien hajautettujen järjestelmien – poliittisten, yhteistyökykyisten ja kestävien – ydin.

Pysy rauhallisena, kun jonot hieman kasvavat, ja muista: sulavan vastapaineen alla oleva järjestelmä on äärettömän parempi kuin kokonaan kaatunut järjestelmä.

logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.