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 Groups | ❌ | In-Process ููุท | ✅ |
| Explicit ACK | ❌ | ❌ | ✅ |
| Pending Messages | ❌ | ❌ | ✅ |
| Recovery ุจุนุฏ Worker Crash | ุถุนูู | ุถุนูู | ✅ |
| Horizontal Scaling | ู ุญุฏูุฏ | ุฏุงุฎู Process | ✅ |
| Distributed Consumers | ❌ | ❌ | ✅ |
| Retry Infrastructure | Manual | Manual | ✅ ูุงุจู ููุจูุงุก ููู 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 ูุงุฒู ูุญุตู ูู:
- Payment
- Notification
- 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
Post a Comment