תורי הודעות ב-Redis: הגדרת Pub/Sub ו-Streams

כאשר שרת מבזבז זמן על משימות רקע, משתמשים ממתינים לתגובות והעסק מאבד מהירות. אנו מגדירים תורי הודעות על Redis באמצעות Pub/Sub ו-Streams כך שעיבוד משימות אסינכרוני פועל ללא תקלות. הצוות שלנו מספק את הפרויקט במפתח מלא—מבחירת המנגנון ועד היישום והתמיכה השוטפת.

פיתוח ותחזוקה של כל סוגי האתרים:

אתרי מידע או יישומי אינטרנט
אתרי תדמית, דפי נחיתה, אתרי חברה, קטלוגים מקוונים, חידונים, אתרי קידום, בלוגים, מקורות חדשות, פורטלי מידע, פורומים, אגרגטורים
אתרי מסחר אלקטרוני או יישומי אינטרנט
חנויות מקוונות, פורטלי B2B, שווקים, בורסות מקוונות, אתרי קאשבק, בורסות, פלטפורמות דרופשיפינג, מנתחי מוצרים
יישומי אינטרנט לניהול תהליכים עסקיים
מערכות CRM, מערכות ERP, פורטלים ארגוניים, מערכות ניהול ייצור, מנתחי מידע
אתרי שירות אלקטרוני או יישומי אינטרנט
פלטפורמות מודעות, בתי ספר מקוונים, בתי קולנוע מקוונים, בוני אתרים, פורטלים לשירותים אלקטרוניים, פלטפורמות אירוח וידאו, פורטלים נושאיים

אלה רק חלק מהסוגים הטכניים של אתרים שאנו עובדים איתם, ולכל אחד מהם יכולים להיות מאפיינים ופונקציונליות ספציפיים משלו, וכן ניתן להתאים אותם לצרכים ולמטרות הספציפיים של הלקוח.

השירותים שאנו מציעים
מציג 1 מתוך 1כל 2062 השירותים
תורי הודעות ב-Redis: הגדרת Pub/Sub ו-Streams
בינוני
~2-3 ימים

הכישורים שלנו:

שאלות נפוצות

העבודות האחרונות

  • פיתוח אתר חברה B2B ADVANCE
    פיתוח אתר חברה B2B ADVANCE
    1504
  • פיתוח אפליקציית ווב עבור FEEDME
    פיתוח אפליקציית ווב עבור FEEDME
    1344
  • פיתוח אתר עבור BELFINGROUP
    פיתוח אתר עבור BELFINGROUP
    1052
  • פיתוח חנות מקוונת לחברת FURNORO
    פיתוח חנות מקוונת לחברת FURNORO
    1307
  • פיתוח אפליקציית ווב עבור Enviok
    פיתוח אפליקציית ווב עבור Enviok
    1050
  • פיתוח אתר לחברת FIXPER
    פיתוח אתר לחברת FIXPER
    1033

הגדרת תורי הודעות ב-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) לקיצוץ מקורב.

מה כלול ולוחות זמנים

  1. ניתוח דרישות ועיצוב סכמת זרמים.
  2. יישום עובדים ב-PHP או Node.js עם טיפול בשגיאות.
  3. הגדרת קבוצות צרכנים, רשומות ממתינות וניטור.
  4. תיעוד תפעולי (קיצוץ, התראות).
  5. הדרכת צוות (שעה אחת).
  6. אחריות על קוד — 3 חודשים.

יישום בסיסי של עובד Streams (מיילים, התראות) — החל מ-500 דולר, 1–2 ימים. עם ניטור, התראות ותיעוד — 1,200 דולר, 2–3 ימים. העלות מחושבת באופן אישי לפי מורכבות הפרויקט.

צרו קשר לייעוץ — נבחן את הפרויקט שלכם תוך יום אחד. חסכו עד 2,000 דולר לחודש במשאבי שרת וזמן מפתחים עם הגדרת תורי ההודעות שלנו ב-Redis.