تصميم الخلفية

بنية قوائم الانتظار: دليل تنفيذ Redis و RabbitMQ و SQS

بناء أنظمة قوائم انتظار موثوقة مع Redis و RabbitMQ و AWS SQS. يغطي dead letter queues والتكافؤ وأنماط المعالجة الحقيقية.

Khalid Aboubakr
19 دقيقة قراءة
Message QueueRabbitmqRedisAsync ProcessingReliabilitySqsDead Letter QueueIdempotency

مقدمة

الطلب-الاستجابة المتزامن ليس دائماً النمط الصحيح. عندما تكون العمليات بطيئة أو غير موثوقة أو تحتاج للمعالجة بالترتيب، توفر طوابير الرسائل موثوقية وقابلية توسع لا يمكن لاستدعاءات API المباشرة مطابقتها.

متى تستخدم الطوابير

المرشحون الجيدون للمعالجة غير المتزامنة

// 1. العمليات طويلة التشغيل async function generateReport(reportId: string): Promise<void> { await queue.add('reports', { reportId, requestedBy: currentUser.id, }); return { jobId: reportId, status: 'queued' }; } // 2. استدعاءات الخدمات الخارجية التي قد تفشل async function sendNotifications(orderId: string): Promise<void> { await queue.add('notifications', { orderId, channels: ['email', 'sms', 'push'], }, { attempts: 5, backoff: { type: 'exponential', delay: 1000 }, }); }

تنفيذ معالجة طوابير موثوقة

إعداد BullMQ (قائم على Redis)

import { Queue, Worker, QueueScheduler } from 'bullmq'; const emailQueue = new Queue('emails', { connection }); const worker = new Worker('emails', async (job) => { const { to, subject, template, data } = job.data; try { await emailService.send({ to, subject, template, data }); return { sent: true, sentAt: new Date() }; } catch (error) { throw error; // رمي الخطأ يؤدي إلى إعادة المحاولة } }, { connection, concurrency: 10, limiter: { max: 100, duration: 1000 }, });

استراتيجيات إعادة المحاولة

// التراجع الأسي مع الاهتزاز const jobOptions = { attempts: 5, backoff: { type: 'exponential', delay: 1000 }, }; // استراتيجية تراجع مخصصة const customBackoff = (attemptsMade: number, error: Error): number => { if (error instanceof RateLimitError) { return error.retryAfter * 1000; } return Math.min(1000 * Math.pow(2, attemptsMade), 30000); };

المعالجة مرة واحدة بالضبط

تحقيق دلالات مرة واحدة بالضبط يتطلب الاستقرار:

class IdempotentProcessor { async process<T>( jobId: string, processor: () => Promise<T> ): Promise<{ result: T; cached: boolean }> { const resultKey = `result:${jobId}`; // التحقق من المعالجة السابقة const cached = await this.redis.get(resultKey); if (cached) { return { result: JSON.parse(cached), cached: true }; } // محاولة الحصول على القفل const acquired = await this.redis.set(lockKey, 'processing', 'EX', 300, 'NX'); if (!acquired) { await this.waitForResult(resultKey); const result = await this.redis.get(resultKey); return { result: JSON.parse(result!), cached: true }; } try { const result = await processor(); await this.redis.setex(resultKey, this.ttl, JSON.stringify(result)); return { result, cached: false }; } finally { await this.redis.del(lockKey); } } }

طوابير الرسائل الميتة

const deadLetterQueue = new Queue('dead-letters', { connection }); const dlqWorker = new Worker('dead-letters', async (job) => { await alertService.notify({ channel: 'slack', message: `فشل المهمة بشكل دائم: ${job.name}`, data: { originalQueue: job.data.queue, error: job.data.error, }, }); await failedJobsRepository.save({ jobId: job.id, ...job.data, receivedAt: new Date(), }); }, { connection });

الخلاصة

البنية القائمة على الطوابير توفر:

  1. الموثوقية: المهام تستمر حتى لو تعطل العمال
  2. قابلية التوسع: أضف عمالاً لزيادة الإنتاجية
  3. المرونة: إعادات المحاولة التلقائية تتعامل مع الإخفاقات العابرة
  4. الترتيب: معالجة المهام بالتسلسل عند الحاجة
  5. الرؤية: مراقبة أعماق الطوابير وأوقات المعالجة

مقالات ذات صلة