הגדרת תורי הודעות ב-Redis: שימוש ב-Pub/Sub וב-Streams
דמיינו שחנות המסחר האלקטרוני שלכם שולחת 10,000 מיילים בשעה. אם זה נעשה באופן סינכרוני, השרת נתקע לדקה, והמשתמש ממתין לתגובה. תור Redis אסינכרוני פותר זאת: המיילים נשלחים ברקע, והבקשה מעובדת מיידית. ב-5 שנות עבודה, יישמנו תורי Redis ביותר מ-50 פרויקטים, הפחתנו את עומס השרת בעד 70% וחסכנו ללקוחות בממוצע 2,000 דולר לחודש בעלויות תשתית.
Redis Pub/Sub ו-Streams הם שני פתרונות פופולריים לעיבוד משימות אסינכרוני. אנו עוזרים להקים תשתית זו במפתח מלא החל מ-500 דולר, עם אחריות לאפס אובדן נתונים עם Streams. נבחן את הפרויקט שלכם.
איך לבחור בין Pub/Sub ל-Streams?
Redis מספק שני מנגנונים להודעות אסינכרוניות: Pub/Sub — פשוט, ללא שמירה, ו-Streams — תור מתמשך עם קבוצות צרכנים, דומה ל-Kafka קל. הבחירה תלויה במשימה: התראות בזמן אמת (Pub/Sub) או תור משימות אמין (Streams).
| תכונה | Pub/Sub | Streams | רשימות (LPUSH/BRPOP) |
|---|---|---|---|
| שמירת נתונים | לא | כן | כן |
| קבוצות צרכנים | לא | כן | לא |
| השמעה חוזרת של היסטוריה | לא | כן | לא |
| מורכבות | מינימלית | בינונית | מינימלית |
| ביצועים ב-10K הודעות/שנייה | 2.1 אלפיות שנייה | 3.4 אלפיות שנייה | 1.8 אלפיות שנייה |
| מקרה שימוש | אירועים בזמן אמת | תור משימות | תור פשוט |
צריכים תור אמין? הלקוחות שלנו חוסכים בדרך כלל 30% מזמן הפיתוח על ידי שימוש ב-Streams במקום לוגיקת ניסיון חוזר מותאמת אישית. צרו קשר — נעזור לכם לבחור את האפשרות האופטימלית.
למה Redis Streams עדיפים על Pub/Sub למשימות קריטיות
Streams נוחים פי 2–3 מ-Pub/Sub בקנה מידה: הם תומכים בקבוצות צרכנים, מאפשרים אישור עיבוד וקריאה חוזרת של הודעות שנכשלו. Pub/Sub הוא פתרון פשוט לאירועים בזמן אמת, אבל ככל שהעומס גדל או שנדרשות ערבויות מסירה, בחרו ב-Streams. זה מוביל לחיסכון של 30% בשעות מהנדס: אין צורך לכתוב לוגיקת ניסיון חוזר מותאמת אישית. בנוסף, עלויות התשתית יורדות בעד 50% בשל פחות אירועי השבתה.
הגדרת Redis Pub/Sub
מתאים להתראות בזמן אמת בתוך האפליקציה. הודעות אינן נשמרות — אם מנוי מנותק, ההודעה אובדת.
// Laravel: публикация через Redis Pub/Sub
use Illuminate\Support\Facades\Redis;
// Publisher
Redis::publish('user-notifications', json_encode([
'user_id' => $userId,
'type' => 'order.shipped',
'message' => 'Ваш заказ отправлен',
]));
// Subscriber (console command)
class RedisSubscribeCommand extends Command
{
protected $signature = 'redis:subscribe';
public function handle(): void
{
Redis::subscribe(['user-notifications'], function (string $message) {
$data = json_decode($message, true);
broadcast(new UserNotificationEvent($data)); // → WebSocket
});
}
} הגדרת Redis Streams
Streams הם הבחירה הנכונה לתור משימות על Redis. הודעות נשמרות בזרם, קבוצות צרכנים עוקבות אחר התקדמות, ורשומות ממתינות עוקבות אחר הודעות שלא עובדו. ערבות מסירה: הודעה נמחקת רק לאחר XACK.
# Создать поток и добавить сообщение
XADD emails * user_id 123 email [email protected] template welcome
# Создать consumer group
XGROUP CREATE emails email-workers $ MKSTREAM
# Читать новые сообщения (воркер 1)
XREADGROUP GROUP email-workers worker-1 COUNT 10 BLOCK 5000 STREAMS emails >
# Подтвердить обработку
XACK emails email-workers <message-id> איך להגדיר קבוצות צרכנים ב-Redis Streams?
קבוצות צרכנים מאפשרות לחלק הודעות בין עובדים. כל עובד מקבל הודעות ייחודיות, ורשומות ממתינות עוקבות אחר אלו שלא עובדו. זהו הבסיס לסובלנות תקלות.
דוגמת עובד ב-PHP
use Illuminate\Support\Facades\Redis;
class RedisStreamWorker
{
private string $stream = 'emails';
private string $group = 'email-workers';
private string $consumer;
public function __construct()
{
$this->consumer = gethostname() . ':' . getmypid();
$this->ensureGroup();
}
private function ensureGroup(): void
{
try {
Redis::xgroup('CREATE', $this->stream, $this->group, '$', true);
} catch (\Throwable) {
// Группа уже существует
}
}
public function run(): void
{
while (true) {
// Сначала обработать pending (не подтверждённые с прошлого запуска)
$pending = Redis::xreadgroup(
$this->group,
$this->consumer,
[$this->stream => '0'], // '0' = pending messages
10
);
$this->processMessages($pending);
// Затем новые сообщения
$messages = Redis::xreadgroup(
$this->group,
$this->consumer,
[$this->stream => '>'], // '>' = only new
10,
5000 // блокировка 5 секунд
);
$this->processMessages($messages);
}
}
private function processMessages(?array $streams): void
{
if (!$streams) return;
foreach ($streams[$this->stream] ?? [] as [$id, $fields]) {
try {
$this->handleEmail($fields);
Redis::xack($this->stream, $this->group, $id);
} catch (\Throwable $e) {
Log::error('Stream message failed', ['id' => $id, 'error' => $e->getMessage()]);
// Сообщение остаётся в pending — будет перечитано при следующем запуске
}
}
}
private function handleEmail(array $fields): void
{
Mail::to($fields['email'])->send(new TemplateMail($fields['template'], $fields));
}
} דוגמת עובד ב-Node.js
import Redis from 'ioredis';
const redis = new Redis({ host: 'redis', port: 6379 });
const STREAM = 'emails';
const GROUP = 'email-workers';
const CONSUMER = `worker-${process.pid}`;
async function startWorker(): Promise<void> {
// Создать группу если не существует
try {
await redis.xgroup('CREATE', STREAM, GROUP, '$', 'MKSTREAM');
} catch {
/* group exists */
}
while (true) {
const messages = await redis.xreadgroup(
'GROUP', GROUP, CONSUMER,
'COUNT', '10',
'BLOCK', '5000',
'STREAMS', STREAM, '>'
) as [string, [string, string[]][]][] | null;
if (!messages) continue;
for (const [, entries] of messages) {
for (const [id, fields] of entries) {
const data = Object.fromEntries(
fields.reduce((acc, val, i) => (i % 2 === 0 ? acc.push([val, fields[i+1]]) : acc, acc), [] as [string,string][])
);
try {
await sendEmail(data);
await redis.xack(STREAM, GROUP, id);
} catch (err) {
console.error('Email failed:', id, err);
}
}
}
}
} השוואה בין PHP ל-Node.js ליישום עובד
| תכונה | PHP (Laravel) | Node.js (ioredis) |
|---|---|---|
| מקביליות | תהליכים (supervisor) | לולאת אירועים |
| טיפול בממתינות | מובנה (Laravel Horizon) | ידני |
| פופולריות | נפוץ מאוד | ביצועים גבוהים |
| מורכבות הגדרה | בינונית | נמוכה |
ניהול זרמים וקיצוץ
# Обрезать поток до 10000 последних сообщений
XTRIM emails MAXLEN ~ 10000
# Автоматически при добавлении
XADD emails MAXLEN ~ 100000 * user_id 123 template welcome ניטור וניפוי שגיאות
עקבו אחר // Laravel: публикация через Redis Pub/Sub use Illuminate\Support\Facades\Redis; // Publisher Redis::publish('user-notifications', json_encode([ 'user_id' => $userId, 'type' => 'order.shipped', 'message' => 'Ваш заказ отправлен', ])); // Subscriber (console command) class RedisSubscribeCommand extends Command { protected $signature = 'redis:subscribe'; public function handle(): void { Redis::subscribe(['user-notifications'], function (string $message) { $data = json_decode($message, true); broadcast(new UserNotificationEvent($data)); // → WebSocket }); } } — הודעות שלא עובדו. אם מספרן עולה, העובד מפגר. השתמשו ב-# Создать поток и добавить сообщение XADD emails * user_id 123 email [email protected] template welcome # Создать consumer group XGROUP CREATE emails email-workers $ MKSTREAM # Читать новые сообщения (воркер 1) XREADGROUP GROUP email-workers worker-1 COUNT 10 BLOCK 5000 STREAMS emails > # Подтвердить обработку XACK emails email-workers <message-id> לבדיקת סטטוס. הגדירו התראות על אורך הממתינות. לדוגמה, כשהממתינות עולות על 1000, אנו שולחים התראה ל-Slack/Telegram.
Redis Streams — תור מתמשך עם קבוצות צרכנים. (תיעוד Redis)[https://redis.io/docs/latest/develop/data-types/streams/]
פרטי ניטור:
- התראות ב-Telegram/Slack כשהממתינות חורגות מהסף (לדוגמה, > 1000).
- רישום שגיאות עם מזהה הודעה לעיבוד ידני חוזר.
- שימוש ב-
use Illuminate\Support\Facades\Redis; class RedisStreamWorker { private string $stream = 'emails'; private string $group = 'email-workers'; private string $consumer; public function __construct() { $this->consumer = gethostname() . ':' . getmypid(); $this->ensureGroup(); } private function ensureGroup(): void { try { Redis::xgroup('CREATE', $this->stream, $this->group, '$', true); } catch (\Throwable) { // Группа уже существует } } public function run(): void { while (true) { // Сначала обработать pending (не подтверждённые с прошлого запуска) $pending = Redis::xreadgroup( $this->group, $this->consumer, [$this->stream => '0'], // '0' = pending messages 10 ); $this->processMessages($pending); // Затем новые сообщения $messages = Redis::xreadgroup( $this->group, $this->consumer, [$this->stream => '>'], // '>' = only new 10, 5000 // блокировка 5 секунд ); $this->processMessages($messages); } } private function processMessages(?array $streams): void { if (!$streams) return; foreach ($streams[$this->stream] ?? [] as [$id, $fields]) { try { $this->handleEmail($fields); Redis::xack($this->stream, $this->group, $id); } catch (\Throwable $e) { Log::error('Stream message failed', ['id' => $id, 'error' => $e->getMessage()]); // Сообщение остаётся в pending — будет перечитано при следующем запуске } } } private function handleEmail(array $fields): void { Mail::to($fields['email'])->send(new TemplateMail($fields['template'], $fields)); } }להקצאת הודעות תקועות לעובד אחר.
טעויות נפוצות ופתרונות
- אין טיפול בממתינות: העובד קרס, הודעות תקועות. פתרון — תמיד לעבד ממתינות בעת ההפעלה.
- אין ערבות אידמפוטנטיות: מייל כפול. השתמשו במפתח אידמפוטנטיות.
- קיצוץ אגרסיבי מדי: אובדן הודעות שלא עובדו. השתמשו ב-~ (tilde) לקיצוץ מקורב.
מה כלול ולוחות זמנים
- ניתוח דרישות ועיצוב סכמת זרמים.
- יישום עובדים ב-PHP או Node.js עם טיפול בשגיאות.
- הגדרת קבוצות צרכנים, רשומות ממתינות וניטור.
- תיעוד תפעולי (קיצוץ, התראות).
- הדרכת צוות (שעה אחת).
- אחריות על קוד — 3 חודשים.
יישום בסיסי של עובד Streams (מיילים, התראות) — החל מ-500 דולר, 1–2 ימים. עם ניטור, התראות ותיעוד — 1,200 דולר, 2–3 ימים. העלות מחושבת באופן אישי לפי מורכבות הפרויקט.
צרו קשר לייעוץ — נבחן את הפרויקט שלכם תוך יום אחד. חסכו עד 2,000 דולר לחודש במשאבי שרת וזמן מפתחים עם הגדרת תורי ההודעות שלנו ב-Redis.







