Competing Consumers

๐Ÿš€ Competing Consumers Pattern ุจุงุณุชุฎุฏุงู… Redis Streams ูˆ .NET Core

ู„ู…ุง Worker ูˆุงุญุฏ ูŠู‚ูˆู„ูƒ: «ูŠุง ุฌู…ุงุนุฉ ุงู„ุทู„ุจุงุช ูƒุชูŠุฑ ุนู„ูŠّุง!» ๐Ÿ˜…

ุชุฎูŠู„ ุฅู† ุนู†ุฏูƒ E-Commerce System، ูˆุงู„ุนู…ูŠู„ ุจูŠุนู…ู„ Order ุฌุฏูŠุฏ.

ุงู„ู€ API ุชุนู…ู„ ุดุบู„ู‡ุง ุงู„ุทุจูŠุนูŠ:

POST /orders

ุชุญูุธ ุงู„ู€ Order ููŠ ุงู„ู€ Database، ูˆุจุนุฏ ูƒุฏู‡ ู…ุญุชุงุฌ ุชุนู…ู„ ุดูˆูŠุฉ ุนู…ู„ูŠุงุช ุชู‚ูŠู„ุฉ ููŠ ุงู„ุฎู„ููŠุฉ:

  • ๐Ÿ’ณ ุชุฌู‡ูŠุฒ ุนู…ู„ูŠุฉ ุงู„ุฏูุน.
  • ๐Ÿ“ฆ ุชุฌู‡ูŠุฒ ุงู„ู€ Order ู„ู„ุดุญู†.
  • ๐Ÿงพ ุฅู†ุดุงุก Invoice.
  • ๐Ÿ“ง ุฅุฑุณุงู„ Confirmation Email.
  • ๐Ÿ“Š ุชุญุฏูŠุซ Analytics.

ุทุจุนًุง ู…ุด ู…ู†ุทู‚ูŠ ุชุฎู„ูŠ ุงู„ุนู…ูŠู„ ูˆุงู‚ู ู…ุณุชู†ูŠ ูƒู„ ุฏู‡ ูŠุญุตู„ ู‚ุจู„ ู…ุง ุงู„ู€ API ุชุฑุฌุน Response.

ูุฃูˆู„ ููƒุฑุฉ ู…ู…ูƒู† ุชูŠุฌูŠ ููŠ ุฏู…ุงุบู†ุง:

Task.Run(async () =>
{
   await ProcessOrderAsync(order);
});

ุฎู„ุตู†ุง ูŠุง ุจุงุดุง ๐Ÿ˜Ž

ู„ู„ุฃุณู… ู„ุฃ. ๐Ÿ˜„

ู„ุฃู† ุฃูˆู„ ู…ุง ุงู„ุณูŠุณุชู… ูŠูƒุจุฑ، ูˆุชุจุฏุฃ ุชุดุบู„ ุฃูƒุชุฑ ู…ู† Server، ูˆุชุญุตู„ Restart ุฃูˆ Deployment ุฃูˆ Crash، ู‡ุชูƒุชุดู ุฅู†ู†ุง ู…ุญุชุงุฌูŠู† ุญุงุฌุฉ ุฃู‚ูˆู‰ ุจูƒุชูŠุฑ.

ูˆู‡ู†ุง ุจูŠุฏุฎู„ ู…ุนุงู†ุง:

⚡ Competing Consumers Pattern

ูˆู…ุนุงู‡:

๐Ÿ”ด Redis Streams + Consumer Groups


๐ŸŽฏ ุงู„ุฃูˆู„: ูŠุนู†ูŠ ุฅูŠู‡ Competing Consumers؟

ุงู„ููƒุฑุฉ ุจุณูŠุทุฉ ุฌุฏًุง.

ุจุฏู„ ู…ุง ูŠูƒูˆู† ุนู†ุฏูƒ Worker ูˆุงุญุฏ ุจูŠู‚ุฑุฃ ุงู„ู€ Messages:

             ┌───────────────┐
Messages ───▶│   Worker #1   │
             └───────────────┘

ู…ู…ูƒู† ูŠูƒูˆู† ุนู†ุฏูƒ 3 ุฃูˆ 5 ุฃูˆ 20 Workers:

                       ┌──────────────┐
                  ┌──▶│ Consumer #1  │
                  │   └──────────────┘
                  │
Messages ──────────┼──▶┌──────────────┐
                  │   │ Consumer #2  │
                  │   └──────────────┘
                  │
                  └──▶┌──────────────┐
                      │ Consumer #3  │
                      └──────────────┘

ู„ูƒู† ุงู„ู…ู‡ู… ู‡ู†ุง:

ุงู„ุชู„ุงุชุฉ ู…ุด ุจูŠุนุงู„ุฌูˆุง ู†ูุณ ุงู„ู€ Message.

ู‡ู… ุจูŠุชู†ุงูุณูˆุง Compete ุนู„ู‰ ุงู„ุดุบู„.

ูŠุนู†ูŠ ู„ูˆ ุนู†ุฏู†ุง:

Order-101
Order-102
Order-103
Order-104
Order-105
Order-106

ู…ู…ูƒู† ุงู„ุชูˆุฒูŠุน ูŠุญุตู„ ูƒุฏู‡:

๐ŸŸข Consumer-A
 ├── Order-101
 └── Order-104

๐Ÿ”ต Consumer-B
 ├── Order-102
 └── Order-105

๐ŸŸฃ Consumer-C
 ├── Order-103
 └── Order-106

ูƒู„ Consumer ุฎุฏ ุฌุฒุก ู…ู† ุงู„ุดุบู„.

ูˆุฏู‡ ุจุงู„ุถุจุท ู…ูู‡ูˆู…:

Competing Consumers

Redis Consumer Groups ู…ุนู…ูˆู„ุฉ ุนุดุงู† ุงู„ุณูŠู†ุงุฑูŠูˆ ุฏู‡ ุชุญุฏูŠุฏًุง؛ ุฏุงุฎู„ ุงู„ู€ Group ุงู„ูˆุงุญุฏ، ุงู„ู€ consumers ุจูŠุชู‚ุงุณู…ูˆุง ุงู„ู€ entries ุจุฏู„ ู…ุง ูƒู„ ูˆุงุญุฏ ูŠุณุชู„ู… ู†ุณุฎุฉ ู…ู†ู‡ุง.


๐Ÿ”ด ุทูŠุจ ุฅูŠู‡ Redis Streams ุฃุตู„ًุง؟

Redis Stream ุชู‚ุฏุฑ ุชุชุฎูŠู„ู‡ ูƒุฃู†ู‡:

Append-Only Event Log

ูŠุนู†ูŠ ุจุฏู„ ู…ุง ุชุญุท ู‚ูŠู…ุฉ ุนุงุฏูŠุฉ:

key → value

ุนู†ุฏูƒ ุณู„ุณู„ุฉ Messages:

orders:stream

┌───────────────────────────────┐
│ 1750000010000-0  Order #101  │
├───────────────────────────────┤
│ 1750000010200-0  Order #102  │
├───────────────────────────────┤
│ 1750000010300-0  Order #103  │
├───────────────────────────────┤
│ 1750000010400-0  Order #104  │
└───────────────────────────────┘

ูˆูƒู„ Entry ู„ูŠู‡ุง ID ู…ุฑุชุจ ุฒู…ู†ูŠًุง.

ุงู„ู€ Producer ูŠุถูŠู ุจุงุณุชุฎุฏุงู…:

XADD

ูˆุงู„ู€ Consumers ุชู‚ุฑุฃ ุจุงุณุชุฎุฏุงู…:

XREADGROUP

ูˆุจุนุฏ ู…ุง ุงู„ู€ Consumer ูŠุฎู„ุต:

XACK

Redis Streams ุจุชุญุชูุธ ุจุงู„ู€ entries ูˆุชูˆูุฑ replay ูˆconsumer tracking ูˆacknowledgment، ูˆุฏู‡ ู…ุฎุชู„ู ุนู† Redis Pub/Sub ุงู„ู„ูŠ ุจุทุจูŠุนุชู‡ fire-and-forget ูˆู…ููŠู‡ูˆุด history ู„ู„ุฑุณุงุฆู„ ุงู„ู‚ุฏูŠู…ุฉ.


๐Ÿง  ุฃู‡ู… ุฌุฒุก: Consumer Group

ุฎู„ูŠู†ุง ู†ู‚ูˆู„ ุฅู† ุนู†ุฏู†ุง Stream ุงุณู…ู‡:

orders:stream

ูˆุนู†ุฏู†ุง Consumer Group:

order-processors

ูˆุฌูˆุงู‡:

order-worker-01
order-worker-02
order-worker-03

ุงู„ุตูˆุฑุฉ ุชุจู‚ู‰:

                    ๐Ÿ”ด Redis
             ┌──────────────────┐
             │  orders:stream   │
             │                  │
Producer ────▶│ Order-101        │
             │ Order-102        │
             │ Order-103        │
             │ Order-104        │
             │ Order-105        │
             └────────┬─────────┘
                      │
                      ▼
            ๐ŸŸ  Consumer Group
             order-processors
                      │
         ┌────────────┼────────────┐
         ▼            ▼            ▼
    ๐ŸŸข Worker-1   ๐Ÿ”ต Worker-2   ๐ŸŸฃ Worker-3
      101            102            103
      104            105            ...

Redis ู‡ูˆ ุงู„ู„ูŠ ุจูŠู…ุณูƒ:

  • ู…ูŠู† ุงุณุชู„ู… ุฅูŠู‡؟
  • ูˆู…ูŠู† ุฎู„ุต؟
  • ูˆู…ูŠู† ุงุณุชู„ู… Message ูˆู„ุณู‡ ู…ุงุนู…ู„ู‡ุงุด ACK؟

ูˆุฏู‡ ุจูŠุชู… ุนู† ุทุฑูŠู‚ ุญุงุฌุฉ ุงุณู…ู‡ุง:

๐Ÿ“‹ Pending Entries List — PEL


๐Ÿ“ฆ ุงู„ุณูŠู†ุงุฑูŠูˆ ุงู„ุนู…ู„ูŠ

ู‡ู†ุจู†ูŠ System ุจุณูŠุท:

Orders API
    │
    │ OrderCreated
    ▼
Redis Stream
"orders:created"
    │
    ▼
Consumer Group
"order-processing"
    │
 ┌───┼─────────────┐
 ▼   ▼             ▼
W1   W2            W3

ูƒู„ Worker ู…ู…ูƒู† ูŠูƒูˆู† ู…ูˆุฌูˆุฏ ุนู„ู‰ Server ู…ุฎุชู„ู.

ู…ุซู„ุงً:

VM-01
└── OrderWorker

VM-02
└── OrderWorker

VM-03
└── OrderWorker

ุฃูˆ Kubernetes:

order-worker Deployment

Replica #1
Replica #2
Replica #3
Replica #4
Replica #5

ูƒู„ู‡ู… ูŠุดุชุฑูƒูˆุง ููŠ ู†ูุณ:

Consumer Group = order-processing

ู„ูƒู† ูƒู„ Instance ู„ุงุฒู… ูŠูƒูˆู† ู„ูŠู‡:

Unique Consumer Name

๐Ÿ› ️ ุฅู†ุดุงุก ุงู„ู…ุดุฑูˆุน

ู‡ู†ุณุชุฎุฏู…:

.NET 8+
StackExchange.Redis
Redis Streams
BackgroundService

ู…ู…ูƒู† ุชุนู…ู„ Worker:

dotnet new worker -n OrderProcessor.Worker

ูˆุชุถูŠู Redis Client:

dotnet add package StackExchange.Redis

StackExchange.Redis ู‡ูˆ ุงู„ู€ Redis client ุงู„ุดุงุฆุน ููŠ .NET، ูˆRedis ู†ูุณู‡ ุนู†ุฏู‡ ุฃู…ุซู„ุฉ ุฑุณู…ูŠุฉ ู„ุงุณุชุฎุฏุงู… Streams ู…ู† C#.


๐Ÿ”Œ ุงู„ุงุชุตุงู„ ุจู€ Redis

ููŠ:

Program.cs

ู†ุนู…ู„:

using StackExchange.Redis;

var builder = Host.CreateApplicationBuilder(args);

builder.Services.AddSingleton<IConnectionMultiplexer>(_ =>
{
   return ConnectionMultiplexer.Connect(
       builder.Configuration.GetConnectionString("Redis")!
   );
});

builder.Services.AddSingleton(sp =>
{
   var redis = sp.GetRequiredService<IConnectionMultiplexer>();
   return redis.GetDatabase();
});

builder.Services.AddHostedService<OrderConsumerWorker>();

var host = builder.Build();

host.Run();

ูˆููŠ:

appsettings.json

ู†ุญุท:

{
 "ConnectionStrings": {
   "Redis": "localhost:6379"
 }
}

Microsoft ุจุชูˆูุฑ BackgroundService ูˆIHostedService ุชุญุฏูŠุฏًุง ู„ู„ู€ long-running background workloads، ูˆุจุชู‚ุฏุฑ ุชุณุฌู„ ุงู„ู€ worker ุจุงุณุชุฎุฏุงู… AddHostedService.


๐Ÿ“ค Producer: ุฅุถุงูุฉ Order ู„ู„ู€ Stream

ู†ูุชุฑุถ ุฅู† ุนู†ุฏู†ุง Event:

public record OrderCreatedEvent(
   Guid EventId,
   Guid OrderId,
   Guid CustomerId,
   decimal Total,
   DateTimeOffset CreatedAt
);

ู†ู‚ุฏุฑ ู†ู†ุดุฑู‡:

using System.Text.Json;
using StackExchange.Redis;

public class OrderEventPublisher
{
   private const string StreamName = "orders:created";

   private readonly IDatabase _redis;

   public OrderEventPublisher(IDatabase redis)
   {
       _redis = redis;
   }

   public async Task PublishAsync(
       OrderCreatedEvent order)
   {
       var json = JsonSerializer.Serialize(order);

       var values = new NameValueEntry[]
       {
           new("eventId", order.EventId.ToString()),
           new("orderId", order.OrderId.ToString()),
           new("payload", json),
           new(
               "createdAt",
               order.CreatedAt.ToUnixTimeMilliseconds()
           )
       };

       await _redis.StreamAddAsync(
           StreamName,
           values
       );
   }
}

ุฏูŠ ุชู‚ุฑูŠุจًุง ุจุชุชุฑุฌู… ููŠ Redis ุฅู„ู‰:

XADD orders:created *
   eventId ...
   orderId ...
   payload ...

Redis ูŠุนู…ู„ Entry ุฒูŠ:

orders:created

1750000001234-0
│
├── eventId   = "..."
├── orderId   = "..."
├── payload   = "{...}"
└── createdAt = "..."

๐Ÿ’ก ุทุจ ู„ูŠู‡ Stream ุจุฏู„ List؟

ู„ุฃู† Redis Stream ู…ุด ู…ุฌุฑุฏ:

Queue

ู‡ูˆ ุนู†ุฏู‡ ู…ูู‡ูˆู…:

Message History
Consumer Groups
Pending Messages
Acknowledgment
Replay
Claiming

ูŠุนู†ูŠ Redis ุนุงุฑู ุฅู†:

Order-101

ุงุชุณู„ู… ู„ู€:

worker-02

ู„ูƒู† ู„ุณู‡:

❌ Not Acknowledged

ูˆุฏูŠ ู…ุนู„ูˆู…ุฉ ุดุฏูŠุฏุฉ ุงู„ุฃู‡ู…ูŠุฉ ูˆู‚ุช ุงู„ุฃุนุทุงู„.


๐Ÿ‘ฅ ุฅู†ุดุงุก Consumer Group

ู‚ุจู„ ุงู„ู‚ุฑุงุกุฉ ู…ุญุชุงุฌูŠู† Group:

order-processing

ููŠ StackExchange.Redis:

private async Task EnsureGroupExistsAsync()
{
   try
   {
       await _redis.StreamCreateConsumerGroupAsync(
           "orders:created",
           "order-processing",
           "0-0",
           createStream: true
       );
   }
   catch (RedisServerException ex)
       when (ex.Message.Contains("BUSYGROUP"))
   {
       // Group already exists.
   }
}

0-0 ู…ุนู†ุงู‡ุง:

ู„ู…ุง ุงู„ู€ Group ูŠุชุนู…ู„، ุงุจุฏุฃ ู…ู† ุฃูˆู„ ุงู„ู€ Stream.

ู„ูˆ ุงุณุชุฎุฏู…ุช:

$

ู…ุนู†ุงู‡:

ุงุจุฏุฃ ู…ู† ุงู„ุฑุณุงุฆู„ ุงู„ุฌุฏูŠุฏุฉ ุงู„ู„ูŠ ุชูŠุฌูŠ ุจุนุฏ ุฅู†ุดุงุก ุงู„ู€ Group.


๐Ÿงต ู†ุนู…ู„ ุงู„ู€ Consumer Worker

public class OrderConsumerWorker : BackgroundService
{
   private const string StreamName =
       "orders:created";

   private const string GroupName =
       "order-processing";

   private readonly IDatabase _redis;

   private readonly ILogger<OrderConsumerWorker>
       _logger;

   private readonly string _consumerName;

   public OrderConsumerWorker(
       IDatabase redis,
       ILogger<OrderConsumerWorker> logger)
   {
       _redis = redis;
       _logger = logger;

       _consumerName =
           $"{Environment.MachineName}-" +
           $"{Guid.NewGuid():N}";
   }

   protected override async Task ExecuteAsync(
       CancellationToken stoppingToken)
   {
       await EnsureGroupExistsAsync();

       while (!stoppingToken.IsCancellationRequested)
       {
           var entries =
               await _redis.StreamReadGroupAsync(
                   StreamName,
                   GroupName,
                   _consumerName,
                   ">",
                   count: 10
               );

           if (entries.Length == 0)
           {
               await Task.Delay(
                   500,
                   stoppingToken
               );

               continue;
           }

           foreach (var entry in entries)
           {
               await ProcessMessageAsync(
                   entry,
                   stoppingToken
               );
           }
       }
   }
}

ุฑูƒุฒ ู…ุนุงูŠุง ููŠ:

">"

ุฏูŠ ู…ุนู†ุงู‡ุง:

ู‡ุงุช Messages ุฌุฏูŠุฏุฉ ู„ู… ูŠุชู… ุชุณู„ูŠู…ู‡ุง ู„ุฃูŠ Consumer ุฏุงุฎู„ ุงู„ู€ Group.

ูˆุฏู‡ ุจุงู„ุถุจุท ุฃุณุงุณ ุงู„ู€ competing behavior ููŠ Redis Consumer Groups.


๐Ÿ”ฅ ูˆู‡ู†ุง ุงู„ุณุญุฑ ุจูŠุญุตู„

ู„ูˆ ุนู†ุฏู†ุง 3 Instances ู…ู† ุงู„ุชุทุจูŠู‚:

Consumer Name:

SERVER-A-8fa23...
SERVER-B-24fd1...
SERVER-C-91ab4...

ูƒู„ู‡ู… ุจูŠู†ูุฐูˆุง:

XREADGROUP
GROUP order-processing

Redis ู…ู…ูƒู† ูŠูˆุฒุน:

                    orders:created
                        │
        ┌───────────────┼───────────────┐
        │               │               │
        ▼               ▼               ▼

๐ŸŸข SERVER-A        ๐Ÿ”ต SERVER-B        ๐ŸŸฃ SERVER-C

Order-101          Order-102          Order-103
Order-104          Order-106          Order-105
Order-109          Order-107          Order-108

ุฅุฐู† ู„ูˆ ุงู„ุถุบุท ุฒุงุฏ:

3 Workers

ุฎู„ูŠู‡ู…:

10 Workers

ู…ุด ู…ุญุชุงุฌ ุชุบูŠุฑ ุงู„ู€ Producer.

ูˆู…ุด ู…ุญุชุงุฌ ุชุบูŠุฑ ุงู„ู€ API.

ูˆู…ุด ู…ุญุชุงุฌ ุชุนู…ู„ Distribution Logic ุจู†ูุณูƒ.

ูˆุฏู‡:

๐Ÿš€ Horizontal Scaling


✅ ACK: ุงู„ู„ุญุธุฉ ุงู„ู…ู‡ู…ุฉ ุฌุฏًุง

ุจุนุฏ ู…ุง Worker ูŠุณุชู„ู… Message:

Redis
│
│ Order-101
▼
Worker-2

Redis ุจูŠุนุชุจุฑู‡ุง:

Pending

ู…ุด Finished.

ูŠุนู†ูŠ:

orders:created

Order-101
   │
   ▼
Consumer: worker-2
Status: PENDING ⚠️

ุจุนุฏ ู…ุง ู†ุฎู„ุต ุงู„ู€ Business Logic:

await ProcessOrderAsync(order);

await _redis.StreamAcknowledgeAsync(
   StreamName,
   GroupName,
   entry.Id
);

ุณุงุนุชู‡ุง ูู‚ุท:

PENDING
  │
  │ XACK
  ▼
✅ Processed

Redis ุจูŠูˆูุฑ XACK ุนุดุงู† ูŠุดูŠู„ ุงู„ุฑุณุงู„ุฉ ู…ู† Pending Entries List ุจุนุฏ ู†ุฌุงุญ ุงู„ู…ุนุงู„ุฌุฉ.


๐Ÿงฉ ProcessMessageAsync

ู…ุซู„ุงً:

private async Task ProcessMessageAsync(
   StreamEntry entry,
   CancellationToken cancellationToken)
{
   try
   {
       var payload =
           entry.Values
               .First(x => x.Name == "payload")
               .Value
               .ToString();

       var order =
           JsonSerializer.Deserialize<OrderCreatedEvent>(
               payload
           );

       if (order is null)
           throw new InvalidOperationException(
               "Invalid order message."
           );

       _logger.LogInformation(
           "Consumer {Consumer} processing Order {OrderId}",
           _consumerName,
           order.OrderId
       );

       await ProcessOrderAsync(
           order,
           cancellationToken
       );

       await _redis.StreamAcknowledgeAsync(
           StreamName,
           GroupName,
           entry.Id
       );

       _logger.LogInformation(
           "Order {OrderId} completed.",
           order.OrderId
       );
   }
   catch (Exception ex)
   {
       _logger.LogError(
           ex,
           "Failed processing message {MessageId}",
           entry.Id
       );

       // ู…ู‡ู…:
       // ู…ุง ู†ุนู…ู„ุด ACK.
   }
}

ูˆู‡ู†ุง ุฎุฏ ุจุงู„ูƒ ู…ู† ุงู„ู‚ุฑุงุฑ ุงู„ู…ู‡ู…:

Success
  ↓
XACK ✅

Failure
  ↓
NO ACK ❌

๐Ÿ’ฅ ุทุจ ู„ูˆ Worker ูˆู‚ุน ููŠ ู†ุต ุงู„ุชู†ููŠุฐ؟

ู‡ู†ุง ุจู‚ู‰ Redis Streams ุจุชุธู‡ุฑ ู‚ูˆุชู‡ุง.

ุชุฎูŠู„:

Redis
│
│ Order-501
▼
Worker-A
│
├── Received ✅
│
├── Started processing...
│
๐Ÿ’ฅ CRASH

Worker-A ู…ุงุนู…ู„ุด:

XACK

ุฅุฐู† Redis ุดุงูŠู:

Order-501

Owner:
Worker-A

State:
PENDING

ุงู„ุฑุณุงู„ุฉ ู…ุด ุงุฎุชูุช.

ูˆุฏู‡ ุฌุฒุก ุฃุณุงุณูŠ ู…ู† reliability model ุจุชุงุน Consumer Groups؛ Redis ูŠู‚ุฏุฑ ูŠุนุฑุถ ุงู„ู€ pending entries، ูˆุงู„ุฑุณุงุฆู„ ุงู„ู„ูŠ ูุถู„ุช ุจุฏูˆู† ACK ู…ู…ูƒู† ุชุชุนู…ู„ ู„ู‡ุง claim ุจูˆุงุณุทุฉ Consumer ุณู„ูŠู… ุจุงุณุชุฎุฏุงู… XCLAIM ุฃูˆ XAUTOCLAIM.


๐Ÿ›Ÿ XAUTOCLAIM

ู…ุซู„ุงً:

Worker-A
๐Ÿ’€ Dead

Order-501
⏳ Pending for 60 seconds

Worker-B ูŠู‚ุฏุฑ ูŠู‚ูˆู„:

ุฃู†ุง ู‡ุงุฎุฏ ุงู„ุฑุณุงุฆู„ ุงู„ู„ูŠ ุฃุตุญุงุจู‡ุง ูˆู‚ุนูˆุง.

ุจุงุณุชุฎุฏุงู…:

XAUTOCLAIM

ุงู„ุตูˆุฑุฉ:

              Pending Entries
                   │
            Order-501
                   │
                   ▼
            ๐Ÿ’€ Worker-A
                   │
            idle > 60 sec
                   │
                   ▼
             XAUTOCLAIM
                   │
                   ▼
            ๐ŸŸข Worker-B
                   │
                   ▼
             Retry Order-501

StackExchange.Redis ุจูŠุฏุนู… StreamAutoClaimAsync، ูˆุจุงู„ุชุงู„ูŠ ุชู‚ุฏุฑ ุชุนู…ู„ Recovery Worker ุฃูˆ periodic recovery loop ู„ู„ู€ abandoned messages.


๐Ÿ” ุฅุฐู† Redis Streams ุจูŠุฏูŠู†ูŠ At-Least-Once Delivery

ูˆุฏูŠ ู†ู‚ุทุฉ ู…ู‡ู…ุฉ ุฌุฏًุง.

ู…ุด:

Exactly Once

ูˆู„ูƒู† ุชู‚ุฑูŠุจًุง:

Message delivered
    │
    ▼
Process
    │
    ├──── Success ────▶ ACK
    │
    └──── Crash ──────▶ Retry

ูˆุฏู‡ ู…ุนู†ุงู‡ ุฅู† ุงู„ุฑุณุงู„ุฉ ู…ู…ูƒู† ุชุชุนุงู„ุฌ:

ู…ุฑุฉ

ุฃูˆ ููŠ ุจุนุถ ุญุงู„ุงุช ุงู„ุฃุนุทุงู„:

ุฃูƒุชุฑ ู…ู† ู…ุฑุฉ

ูˆู…ู† ู‡ู†ุง ุชูŠุฌูŠ ุฃู‡ู…ูŠุฉ:

๐Ÿ›ก️ Idempotency

ู…ุซู„ุงً ู„ูˆ Order:

OrderId = 500

ุงุชุนุงู„ุฌ ู…ุฑุชูŠู†، ู…ุงูŠู†ูุนุด ุชุนู…ู„:

Charge Customer
Charge Customer again ๐Ÿ˜ฑ

ู„ุงุฒู… ุงู„ู€ processing layer ุชุจู‚ู‰ ู…ุตู…ู…ุฉ ุจุญูŠุซ ุฅุนุงุฏุฉ ู†ูุณ ุงู„ู€ message ู…ุง ุชุนู…ู„ุด side effects ู…ูƒุฑุฑุฉ.


๐Ÿง  ุฏู„ูˆู‚ุชูŠ ุงู„ุณุคุงู„ ุงู„ู…ู‡ู…:

ุฅูŠู‡ ุงู„ูุฑู‚ ุนู† Task.Run؟

ู†ูุชุฑุถ ุฅู† ุงู„ู€ Controller ููŠู‡:

[HttpPost]
public async Task<IActionResult> CreateOrder()
{
   var order = await CreateOrderAsync();

   _ = Task.Run(async () =>
   {
       await ProcessOrderAsync(order);
   });

   return Ok();
}

ุธุงู‡ุฑูŠًุง ุงู„ู…ูˆุถูˆุน ุฌู…ูŠู„:

HTTP Request
  │
  ├── Create Order
  │
  └── Task.Run()
          │
          ▼
    Background Work

ู„ูƒู† ุนู†ุฏูƒ ู…ุดุงูƒู„ ู…ุนู…ุงุฑูŠุฉ ูƒุจูŠุฑุฉ.


๐Ÿ’ฃ ุงู„ู…ุดูƒู„ุฉ ุงู„ุฃูˆู„ู‰: Process Crash

ู…ุน:

Task.Run

ุงู„ุดุบู„ ู…ูˆุฌูˆุฏ ููŠ Memory ุจุชุงุนุฉ:

ASP.NET Process

ูŠุนู†ูŠ:

Order
 │
 ▼
Task.Run()
 │
 ▼
RAM

ู„ูˆ ุญุตู„:

Application Restart
IIS Recycle
Pod Restart
Deployment
VM Crash
Process Kill

ูุงู„ู€ Task ู…ู…ูƒู† ุจุจุณุงุทุฉ:

๐Ÿ’จ ุชุฎุชููŠ.

ู…ููŠุด Broker ุฎุงุฑุฌูŠ ูุงูƒุฑ:

ูƒุงู† ููŠู‡ Order ู„ุณู‡ ู…ุญุชุงุฌ ูŠุชุนุงู„ุฌ.

Microsoft ู†ูุณู‡ุง ุจุชูุฑู‚ ุจูŠู† arbitrary background threads ูˆุจูŠู† hosted services ุงู„ู…ุฑุชุจุทุฉ ุจุนู…ุฑ ุงู„ุชุทุจูŠู‚، ูˆุจุชูˆุตูŠ ุจุงุณุชุฎุฏุงู… Hosted Services ู„ู„ู€ long-running background workloads ุจุฏู„ ุงู„ุงุนุชู…ุงุฏ ุนู„ู‰ thread ุนุดูˆุงุฆูŠ ุบูŠุฑ ู…ُุฏุงุฑ.


๐Ÿ”ด Redis Streams

ู‡ู†ุง ุงู„ูˆุถุน ู…ุฎุชู„ู:

Application
  │
  │ XADD
  ▼
┌────────────────────┐
│      Redis         │
│                    │
│ Order-101          │
│ Order-102          │
│ Order-103          │
└────────────────────┘

ุงู„ู€ Application ู†ูุณู‡ุง ู…ู…ูƒู† ุชู…ูˆุช:

Application ๐Ÿ’€

ู„ูƒู† ุงู„ู€ Stream ู…ูˆุฌูˆุฏ ุฎุงุฑุฌ ุงู„ู€ process.

Worker ุฌุฏูŠุฏ ูŠุทู„ุน:

Worker-B ๐ŸŸข

ูˆูŠุฑุฌุน ูŠูƒู…ู„.

ูˆุฏู‡ ุงู„ู„ูŠ ู…ู…ูƒู† ู†ุณู…ูŠู‡ ู‡ู†ุง:

Persistence / Recoverability / Rebindability

ูŠุนู†ูŠ ุงู„ู€ Work ู…ุด ู…ุฑุจูˆุท ุจุนู…ุฑ ุงู„ู€ Process ุงู„ู„ูŠ ุฃู†ุดุฃู‡.


๐Ÿงต ุทูŠุจ ูˆุงู„ู€ Channel<T>؟

.NET ุนู†ุฏู‡:

Channel<T>

ูˆุฏูŠ ุญุงุฌุฉ ู…ู…ุชุงุฒุฉ ุฌุฏًุง.

ู…ุซู„ุงً:

var channel =
  Channel.CreateBounded<Order>(1000);

Producer:

await channel.Writer.WriteAsync(order);

Consumer:

await foreach (
  var order in channel.Reader.ReadAllAsync())
{
  await ProcessAsync(order);
}

ุฏู‡ ู…ู…ุชุงุฒ ู„ูˆ ูƒู„ ุงู„ุดุบู„:

ุฏุงุฎู„ ู†ูุณ ุงู„ู€ Process

๐ŸŸข Architecture ู…ุน Channel

┌─────────────────────────────┐
│       ASP.NET Process       │
│                             │
│ API                         │
│  │                          │
│  ▼                          │
│ Channel<Order>              │
│  │                          │
│  ├── Worker 1               │
│  ├── Worker 2               │
│  └── Worker 3               │
│                             │
└─────────────────────────────┘

ุฌู…ูŠู„ ุฌุฏًุง.

ู„ูƒู† ูƒู„ู‡ ู…ูˆุฌูˆุฏ ุฏุงุฎู„:

Process ูˆุงุญุฏ.

๐Ÿ”ด Architecture ู…ุน Redis Streams

┌──────────────┐
│ Orders API   │
│ Server A     │
└──────┬───────┘
      │
      │ XADD
      ▼
╔════════════════════════╗
║      REDIS STREAM      ║
║                        ║
║ orders:created         ║
╚═══════════╤════════════╝
           │
           │ Consumer Group
           │
     ┌─────┼─────┐
     │     │     │
     ▼     ▼     ▼

┌────────┐ ┌────────┐ ┌────────┐
│Worker A│ │Worker B│ │Worker C│
│Server 1│ │Server 2│ │Server 3│
└────────┘ └────────┘ └────────┘

ูˆู‡ู†ุง ุฅุญู†ุง ุฏุฎู„ู†ุง ุนุงู„ู…:

๐ŸŒ Distributed Processing


⚖️ Redis Streams vs Channel vs Task.Run

ุงู„ุฎุงุตูŠุฉ Task.Run Channel<T> Redis Streams
Background Processing
ุฏุงุฎู„ ู†ูุณ Process❌ ู…ุด ู„ุงุฒู…
Multi-Server
Message History
Replay
Consumer GroupsIn-Process ูู‚ุท
Explicit ACK
Pending Messages
Recovery ุจุนุฏ Worker CrashุถุนูŠูุถุนูŠู
Horizontal Scalingู…ุญุฏูˆุฏุฏุงุฎู„ Process
Distributed Consumers
Retry InfrastructureManualManual✅ ู‚ุงุจู„ ู„ู„ุจู†ุงุก ููˆู‚ PEL
Survives App Restart
ู…ู†ุงุณุจ ู„ู€ durable jobs✅ ู…ุน ุฅุนุฏุงุฏ Redis ู…ู†ุงุณุจ

๐ŸŽฏ ุฅู…ุชู‰ ุฃุณุชุฎุฏู… Channel؟

ุฃู†ุง ู…ุด ุจู‚ูˆู„:

Redis Streams replaces Channel

ู„ุฃ ุฎุงู„ุต.

Channel<T> ู…ู…ุชุงุฒ ู„ู…ุง ุงู„ู…ุดูƒู„ุฉ:

In-Process Concurrency

ู…ุซู„ุงً:

Image Processing Pipeline
ุฏุงุฎู„ Service ูˆุงุญุฏ

ุฃูˆ:

Buffer ุจูŠู† Producer ูˆConsumer
ุฏุงุฎู„ ู†ูุณ Application Instance

ูˆู‚ุชู‡ุง ุฅุฏุฎุงู„ Redis ู…ู…ูƒู† ูŠูƒูˆู† Overengineering.


๐ŸŽฏ ุฅู…ุชู‰ Redis Streams؟

Redis Streams ุชุจุฏุฃ ุชุจู‚ู‰ ุฌุฐุงุจุฉ ู„ู…ุง ุชู‚ูˆู„:

ุงู„ุดุบู„ ู…ู‡ู… ูˆู…ูŠู†ูุนุด ูŠุถูŠุน.

ุฃูˆ:

ุฃู†ุง ู…ุญุชุงุฌ ุฃูƒุชุฑ ู…ู† Server ูŠุนุงู„ุฌ.

ุฃูˆ:

ู„ูˆ Worker ูˆู‚ุน ุญุฏ ุชุงู†ูŠ ูŠูƒู…ู„.

ุฃูˆ:

ุนุงูŠุฒ ุฃุฒูˆุฏ Workers ูˆู‚ุช ุงู„ุถุบุท.

ุฃูˆ:

ุนุงูŠุฒ ุฃุนุฑู ู…ูŠู† ุงุณุชู„ู… ุงู„ุฑุณุงู„ุฉ ูˆู„ุณู‡ ู…ุงุฎู„ุตู‡ุงุด.

ุณุงุนุชู‡ุง ุฅุญู†ุง ุฎุฑุฌู†ุง ู…ู†:

Background Programming

ูˆุฏุฎู„ู†ุง ููŠ:

Distributed Messaging Architecture


๐Ÿง  ู†ู‚ุทุฉ ู…ุนู…ุงุฑูŠุฉ ู…ู‡ู…ุฉ ุฌุฏًุง:

Consumer ≠ Consumer Group

ุฏูŠ ู…ู† ุฃูƒุชุฑ ุงู„ุญุงุฌุงุช ุงู„ู„ูŠ ุจุชุชู„ุฎุจุท.

ู†ูุชุฑุถ ุฅู† ุงู„ู€ Order ู„ุงุฒู… ูŠุญุตู„ ู„ู‡:

  1. Payment
  2. Notification
  3. Analytics

ู„ูˆ ุญุทูŠุช ุงู„ุชู„ุงุชุฉ Consumers ููŠ:

Consumer Group ูˆุงุญุฏ

ู‡ูŠุญุตู„ ุฅูŠู‡؟

Order-100
  │
  ▼
Consumer Group
  │
  ├──── Payment Worker
  ├──── Notification Worker
  └──── Analytics Worker

Redis ู‡ูŠุฎุชุงุฑ ูˆุงุญุฏ ูู‚ุท ู…ู†ู‡ู… ูŠุณุชู„ู… ุงู„ู€ Message.

ูˆุฏู‡ ุบู„ุท.

ู„ุฃู†ู†ุง ู…ุญุชุงุฌูŠู† ุงู„ุชู„ุงุช ุฎุฏู…ุงุช ุชุดูˆู ุงู„ู€ Order.


✅ ุงู„ุญู„

ู†ุนู…ู„:

orders:created

ู„ูƒู† 3 Groups:

payment-processors

notification-processors

analytics-processors

ูุชุจู‚ู‰ ุงู„ุตูˆุฑุฉ:

                    ๐Ÿ”ด orders:created
                          │
         ┌────────────────┼────────────────┐
         │                │                │
         ▼                ▼                ▼

๐ŸŸข payment-group   ๐Ÿ”ต notification-group   ๐ŸŸฃ analytics-group

     │                    │                    │
┌────┼────┐          ┌────┼────┐          ┌────┼────┐
▼    ▼    ▼          ▼    ▼    ▼          ▼    ▼    ▼
P1   P2   P3         N1   N2   N3         A1   A2   A3

ูƒู„ Group ูŠุดูˆู ุงู„ู€ Stream ุจุดูƒู„ ู…ุณุชู‚ู„، ู„ูƒู† ุฌูˆّู‡ ูƒู„ Group ุงู„ู€ Workers ูŠู†ุงูุณูˆุง ุจุนุถ.

ุฏูŠ ูˆุงุญุฏุฉ ู…ู† ุฃู‚ูˆู‰ ุงู„ุฃููƒุงุฑ ููŠ Redis Streams. Redis ุจูŠุณู…ุญ ุจุฃูƒุชุฑ ู…ู† consumer group ู„ู†ูุณ ุงู„ู€ stream، ูˆูƒู„ Group ุนู†ุฏู‡ cursor ูˆุชุชุจุน ู…ุณุชู‚ู„.


⚡ ู…ุซุงู„ Scaling ุญู‚ูŠู‚ูŠ

ุงู„ุณุงุนุฉ 11 ุงู„ุตุจุญ:

Orders/sec = 50

ุนู†ุฏูƒ:

2 Workers

ูƒูˆูŠุณูŠู† ุฌุฏًุง.

ุงู„ุณุงุนุฉ 8 ู…ุณุงุกً Black Friday:

Orders/sec = 5000

ุชู‚ุฏุฑ ุชุนู…ู„:

kubectl scale deployment order-worker
--replicas=20

ูุชุชุญูˆู„:

Redis Stream
   │
┌───┼───┐
▼   ▼   ▼
W1  W2  W3

ุฅู„ู‰:

Redis Stream
│
├── W01
├── W02
├── W03
├── W04
├── W05
├── ...
└── W20

ู…ู† ุบูŠุฑ ู…ุง ุงู„ู€ Producer ูŠุนุฑู ุฃุตู„ุงً ุฅู†ูƒ ุฒูˆุฏุช Workers.

ูˆุฏูŠ ู†ู‚ุทุฉ ุดุฏูŠุฏุฉ ุงู„ุฃู‡ู…ูŠุฉ:

Producer ูˆ Consumers Decoupled


๐Ÿ”ฅ Backpressure ุจุดูƒู„ ุทุจูŠุนูŠ

ุชุฎูŠู„:

Producer Rate
1000 msg/sec

ู„ูƒู† Workers ูŠู‚ุฏุฑูˆุง ูŠุนุงู„ุฌูˆุง:

600 msg/sec

ู…ุน Task.Run ู…ู…ูƒู† ุชุจุฏุฃ ุชูˆู„ุน Thread Pool ๐Ÿ˜…:

1000
2000
5000
10000 Tasks...

ู„ูƒู† ููŠ Stream:

Producer
1000 msg/sec
    │
    ▼
Redis Stream
████████████████
    │
    ▼
Consumers
600 msg/sec

ุงู„ูุฑู‚:

400 Message/sec

ูŠุชุฑุงูƒู… ููŠ ุงู„ู€ Stream ุจุฏู„ ู…ุง ูŠุจู‚ู‰ ุขู„ุงู ุงู„ู€ Tasks ุงู„ุนุดูˆุงุฆูŠุฉ ุฏุงุฎู„ Memory ุจุชุงุนุฉ ุงู„ู€ API.

ูˆุฏู‡ ูŠุณู…ุญ ู„ูƒ ุชุฑุงู‚ุจ:

  • Consumer Lag
  • Pending Messages
  • Processing Rate

ูˆุชุงุฎุฏ ู‚ุฑุงุฑ:

Scale Out ๐Ÿš€

Redis ุจูŠูˆูุฑ group state ูˆpending information ูˆุฃุฏูˆุงุช ุฒูŠ XPENDING ุชุณุงุนุฏูƒ ุชุฑุงู‚ุจ ุงู„ุฑุณุงุฆู„ ุงู„ู„ูŠ ู„ุณู‡ ุชุญุช ุงู„ู…ุนุงู„ุฌุฉ.


⚠️ ACK ู‚ุจู„ ูˆู„ุง ุจุนุฏ ุงู„ู€ Processing؟

ู…ู…ู†ูˆุน ุชุนู…ู„:

Receive
↓
ACK
↓
Process

ู„ูŠู‡؟

ู„ุฃู† ุงู„ุณูŠู†ุงุฑูŠูˆ ุฏู‡ ู…ู…ูƒู† ูŠุญุตู„:

Receive Order
    │
    ▼
ACK ✅
    │
    ▼
Process
    │
    ▼
๐Ÿ’ฅ CRASH

Redis ูุงูƒุฑ ุฅู†ูƒ ุฎู„ุตุช.

ู„ูƒู† ุงู„ุญู‚ูŠู‚ุฉ:

Order ุถุงุน.

ุงู„ุตุญ:

Receive
 │
 ▼
Process
 │
 ▼
Business Transaction
 │
 ▼
Success
 │
 ▼
XACK ✅

๐Ÿ›ก️ Transactional Outbox ู…ู‡ู… ุจุฑุถู‡

Redis Streams ุจุชุญู„ ู„ูƒ ู…ุดูƒู„ุฉ:

Consumer Reliability

ู„ูƒู† ููŠ Producer ููŠู‡ ู…ุดูƒู„ุฉ ู…ุฎุชู„ูุฉ.

ู†ูุชุฑุถ ุฅู† Orders API ุชุนู…ู„:

INSERT Order INTO SQL

ูˆุจุนุฏู‡ุง:

XADD orders:created

ู…ู…ูƒู† ูŠุญุตู„:

SQL INSERT ✅

      ↓

Application Crash ๐Ÿ’ฅ

      ↓

XADD ❌

ุงู„ู€ Order ุงุชุญูุธ.

ู„ูƒู† ุงู„ู€ Event ู…ุง ุงุชุจุนุชุด.

ูˆุฏูŠ ุงุณู…ู‡ุง:

Dual Write Problem

ู„ูˆ ุงู„ู€ Event Critical، ุงู„ุญู„ ุงู„ุฃุดู‡ุฑ ู…ุนู…ุงุฑูŠًุง ูŠุจู‚ู‰:

SQL Transaction
   │
   ├── Insert Order
   └── Insert OutboxMessage
            │
            ▼
      Outbox Publisher
            │
            ▼
        Redis Stream

ูŠุนู†ูŠ Redis Streams ู…ุด ุจุฏูŠู„ ุนู† Transactional Outbox؛ ุงู„ุงุชู†ูŠู† ุจูŠุญู„ูˆุง ู…ุดูƒู„ุชูŠู† ู…ุฎุชู„ููŠู†.


๐Ÿงฌ Architecture ุงู„ู†ู‡ุงุฆูŠุฉ

                      ๐Ÿ‘ค Customer
                          │
                          ▼
                 ┌────────────────┐
                 │   Orders API   │
                 └────────┬───────┘
                          │
                          ▼
                   SQL Transaction
                          │
               ┌─────────┴────────┐
               ▼                  ▼
           ๐Ÿ›’ Orders          ๐Ÿ“ค Outbox
                                  │
                                  ▼
                        Outbox Publisher
                                  │
                                  │ XADD
                                  ▼

                 ╔══════════════════════╗
                 ║     REDIS STREAM     ║
                 ║                      ║
                 ║    orders:created    ║
                 ╚══════════╤═══════════╝
                            │
                Consumer Group:
                order-processing
                            │
             ┌──────────────┼──────────────┐
             │              │              │
             ▼              ▼              ▼
       ๐ŸŸข Worker-01    ๐Ÿ”ต Worker-02    ๐ŸŸฃ Worker-03
             │              │              │
             ▼              ▼              ▼
         Process          Process         Process
             │              │              │
             └──────┬───────┴───────┬─────┘
                    │               │
                Success          Failure
                    │               │
                    ▼               ▼
                XACK ✅        Pending ⚠️
                                    │
                                    ▼
                                Retry /
                              XAUTOCLAIM

ุฏูŠ Architecture ู‚ุฑูŠุจุฉ ุฌุฏًุง ู…ู† ุงู„ู„ูŠ ู…ู…ูƒู† ุชุณุชุฎุฏู…ู‡ ูุนู„ุงً ููŠ Production.


๐Ÿง  ุงู„ุฎู„ุงุตุฉ ุงู„ู…ุนู…ุงุฑูŠุฉ

ุงู„ู€ Competing Consumers Pattern ู…ุด ู…ุฌุฑุฏ:

ู†ุดุบู„ ุดูˆูŠุฉ Threads.

ุงู„ููƒุฑุฉ ุงู„ุญู‚ูŠู‚ูŠุฉ:

                 Workload
                   │
                   ▼
         Shared Message Source
                   │
                   ▼
           Consumer Group
                   │
      ┌────────────┼────────────┐
      ▼            ▼            ▼
   Worker A     Worker B     Worker C

ูˆูƒู„ Worker:

Stateless ู‚ุฏุฑ ุงู„ุฅู…ูƒุงู†

ูˆุจุงู„ุชุงู„ูŠ ุชู‚ุฏุฑ:

Scale Out
Scale In
Restart
Deploy
Recover
Retry

ู…ู† ุบูŠุฑ ู…ุง ุงู„ู€ Producer ูŠูƒูˆู† ู…ุฑุชุจุท ุจุนุฏุฏ ุงู„ู€ Workers.


๐ŸฅŠ Task.Run vs Channel vs Redis Streams ููŠ ุฌู…ู„ุฉ ูˆุงุญุฏุฉ

Task.Run

"ุงุนู…ู„ ุงู„ุญุงุฌุฉ ุฏูŠ ููŠ ุงู„ุฎู„ููŠุฉ ุฌูˆู‡ ู†ูุณ ุงู„ู€ Process."

Channel<T>

"ุงุนู…ู„ ู„ูŠ Producer/Consumer Queue ู…ุญุชุฑู…ุฉ،
ุจุณ ุฌูˆู‡ ู†ูุณ ุงู„ู€ Process."

Redis Streams + Consumer Groups

"ุงุนู…ู„ ู„ูŠ Distributed Work Queue / Event Stream
ุชู‚ุฏุฑ Workers ุนู„ู‰ Servers ู…ุฎุชู„ูุฉ ุชุชู†ุงูุณ ุนู„ูŠู‡ุง،
ู…ุน Tracking ูˆACK ูˆRecovery."

ูˆุฏูŠ ู‡ูŠ ุงู„ู†ู‚ู„ุฉ ุงู„ู…ุนู…ุงุฑูŠุฉ ุงู„ุญู‚ูŠู‚ูŠุฉ. ๐Ÿš€


๐ŸŽฏ ู‚ุงุนุฏุฉ ุณู‡ู„ุฉ ุชุญูุธู‡ุง

ู„ูˆ ุงู„ุดุบู„:

Cheap + Local + Loss Acceptable

ููƒุฑ ููŠ:

Task / Channel

ู„ูƒู† ู„ูˆ ุงู„ุดุบู„:

Important
+
Distributed
+
Must Survive App Restarts
+
Needs Retry
+
Needs Horizontal Scaling

ุงุจุฏุฃ ุชููƒุฑ ููŠ:

Redis Streams
RabbitMQ
Azure Service Bus
Kafka

ุญุณุจ ุทุจูŠุนุฉ ุงู„ู…ุดูƒู„ุฉ.

Redis Streams ุงุฎุชูŠุงุฑ ู„ุทูŠู ุฌุฏًุง ุฎุตูˆุตًุง ู„ู…ุง ูŠูƒูˆู† Redis ู…ูˆุฌูˆุฏ ุจุงู„ูุนู„ ููŠ ุงู„ู€ infrastructure ูˆู…ุญุชุงุฌ Streaming/queue-like capabilities ุจุฏูˆู† ุฅุฏุฎุงู„ ู…ู†ุตุฉ ุฃูƒุจุฑ ู…ู† ุงุญุชูŠุงุฌูƒ. Redis ู†ูุณู‡ุง ุจุชู‚ุฏู… Streams ูƒุญู„ ู„ู„ู€ ordered event processing ูˆุงู„consumer groups ูˆุงู„replay ูˆุงู„short/moderate retention workloads.


⚠️ Challenges ู„ุงุฒู… ุชุจู‚ู‰ ููŠ ุญุณุงุจูƒ

ูˆุฃู†ุช ุจุชุญูˆู„ ุงู„ุชุตู…ูŠู… ุฏู‡ ู„ู€ Production، ุญุท ููŠ Architecture Checklist:

  • ☠️ Poison Messages
  • ๐Ÿ” Retry Strategy
  • ๐Ÿ›ก️ Idempotency
  • ๐Ÿ”ข Message Ordering
  • ๐Ÿ’€ Dead Letter Queue / Dead Letter Stream
  • ๐Ÿ“‹ Pending Entries Management
  • ๐Ÿ›Ÿ XAUTOCLAIM / Abandoned Consumers
  • ๐Ÿ“ˆ Consumer Lag Monitoring
  • ๐Ÿ“ฆ Stream Retention / Trimming
  • ๐Ÿ’พ Redis Persistence Configuration
  • ๐Ÿ” Security & Authentication
  • ๐Ÿ”„ Graceful Shutdown
  • ๐Ÿงฌ Message Schema Versioning
  • ๐Ÿ†” Unique Consumer Naming
  • ⚖️ Load Distribution
  • ๐Ÿงพ Transactional Outbox
  • ๐Ÿ”ญ Observability & Distributed Tracing

๐Ÿš€ ุงู„ุฎู„ุงุตุฉ

ุฃู‡ู… ุญุงุฌุฉ ุชุทู„ุน ุจูŠู‡ุง ู…ู† ุงู„ู…ู‚ุงู„ ุฅู†:

Competing Consumers

ู…ุด ู…ุฌุฑุฏ Optimization ุนุดุงู† ู†ุฎู„ูŠ ุงู„ู€ processing ุฃุณุฑุน.

ู‡ูˆ Pattern ุจูŠุบูŠุฑ ุทุฑูŠู‚ุฉ ุชููƒูŠุฑูƒ ู…ู†:

Application
 │
 └── Background Task

ุฅู„ู‰:

                   Distributed System

Producer
 │
 ▼
Durable / Recoverable Work
 │
 ▼
Consumer Group
 │
 ├── Consumer #1
 ├── Consumer #2
 ├── Consumer #3
 └── Consumer #N

ูุจุฏู„ ู…ุง ุงู„ู€ API ุชู‚ูˆู„:

"ุฃู†ุง ู‡ุนู…ู„ ุงู„ู€ Order Processing ุจู†ูุณูŠ ููŠ ุงู„ุฎู„ููŠุฉ."

ุจุชู‚ูˆู„:

"ุฃู†ุง ู‡ุณุฌู„ ุฅู† ููŠู‡ Work ู…ุญุชุงุฌ ูŠุชุนู…ู„، ูˆุฃุณูŠุจ Pool ู…ู† ุงู„ู€ Workers ูŠุชู†ุงูุณ ุนู„ูŠู‡."

ูˆุฏู‡ ุจูŠุฏูŠูƒ:

๐Ÿš€ Scalability

๐Ÿ›ก️ Reliability

๐Ÿ” Retries

๐Ÿ’พ Persistence / Recoverability

๐ŸŒ Distributed Processing

⚡ Loose Coupling

ูˆุฏู‡ ุจุงู„ุถุจุท ุงู„ุณุจุจ ุงู„ู„ูŠ ุจูŠุฎู„ูŠ Redis Streams + Consumer Groups ุงุฎุชูŠุงุฑ ู‚ูˆูŠ ุฌุฏًุง ู„ุชู†ููŠุฐ Competing Consumers Pattern ููŠ .NET.

Comments

Popular posts from this blog

Best Practices for Storing and Loading JSON Objects from a Large SQL Server Table Using .NET Core

Maxpooling vs minpooling vs average pooling

Understand the Softmax Function in Minutes