نظام التنفيذ المؤجل على RabbitMQ



مرحبا!



اسمي نيكيتا ، وأشرف على تطوير العديد من المشاريع في DomClick. اليوم أريد أن أكمل موضوع "صور مضحكة" في عالم RabbitMQ. في مقالته ، اعتبر أليكسي كازاكوف أداة قوية مثل قوائم الانتظار المؤجلة والتطبيقات المختلفة لاستراتيجية إعادة المحاولة. سنتحدث اليوم عن كيفية استخدام RabbitMQ لجدولة المهام الدورية.



لماذا احتجنا إلى إنشاء دراجتنا الخاصة ولماذا تخلينا عن الكرفس وأدوات إدارة المهام الأخرى؟ الحقيقة هي أنها لم تتناسب مع مهامنا ومتطلباتنا للتسامح مع الأخطاء ، والتي تعتبر صارمة للغاية في شركتنا.



عند التبديل إلى Docker و Kubernetes ، يواجه العديد من المطورين مشكلات تنظيم المهام الدورية ، ويتم إطلاق التيجان بدف ، ويترك التحكم في العملية الكثير مما هو مرغوب فيه. ثم هناك مشاكل في ذروة الأحمال خلال النهار.



كانت مهمتي أن أنفذ في المشروع نظامًا موثوقًا به لمعالجة المهام الدورية ، مع قابلية التوسع والتحمل للأخطاء بسهولة. مشروعنا بلغة Python ، لذلك كان من المنطقي أن نرى كيف يناسبنا الكرفس. هذه أداة جيدة ، لكننا واجهنا معها غالبًا الموثوقية وقابلية التوسع ومشكلات الإصدار السلس. جراب واحد - مجموعة معالجة واحدة. عند قياس الكرفس ، يجب عليك زيادة موارد حجرة واحدة ، لأنه لا توجد مزامنة بين الحاضنات ، مما يعني إيقاف معالجة المهام ، وإن كان ذلك مؤقتًا. وإذا كانت المهام طويلة الأجل أيضًا ، فقد خمنت بالفعل مدى صعوبة إدارتها. العيب الثاني الواضح: من خارج الصندوق ، لا يوجد دعم لعدم التزامن ، وهذا مهم بالنسبة لنا ، لأن المهام تحتوي بشكل أساسي على عمليات الإدخال / الإخراج ، ويعمل الكرفس على سلاسل العمليات.



في ذلك الوقت (2018) ، لم نجد أداة مناسبة جاهزة ، وبدأنا في تطوير أداة خاصة بنا. مع الأخذ في الاعتبار وظيفة التنفيذ المؤجل للمهام وتبادل الرسائل الميتة ، قررنا إنشاء نظام لمعالجة المهام الدورية. بدا المفهوم شيئًا كالتالي:



صورة



سأحاول شرح ماذا.



  1. يتم إرسال المهام في شكل رسالة إلى تبادل المجدول.
  2. و routing_keyالبرنامج يحصل في التفريخ طابور المطلوبة، الذي يحتوي على معلمة message_ttl، وكذلك الاتصال مع تبادل المعالج كما تبادل إلكتروني صفقة. لا ترتبط قائمة انتظار "النضج" بنوع المهام ، فهي تلعب دور "المؤقت" فقط ، أي أنه يمكنك إنشاء العديد من قوائم الانتظار التي تحتاج إليها وإدارتها routing_key.
  3. نظرًا لأن قائمة الانتظار لا تحتوي على مستمعين ، فإن الرسائل بعد "النضج" في قائمة الانتظار تنتقل إلى تبادل المعالج.
  4. ثم يلتقط المستهلك المجاني (مستهلك المعالجة) الرسالة وينفذها. بعد التنفيذ ، تتكرر الدورة إذا لزم الأمر.


ما هي ميزة مثل هذا المخطط؟



  1. التنفيذ المرحلي ، أي لن تتم معالجة مهمة جديدة إذا لم تكتمل المهمة السابقة.
  2. مستمع واحد (مستهلك) ، أي يمكنك إنشاء عاملين عالميين ومتخصصين. تم التوسع عن طريق زيادة عدد القرون المطلوبة.
  3. انشر المهام الجديدة دون تعطيل عمل المهام الحالية. يكفي تحديث قرون المستمع برفق وإرسال الرسالة المناسبة إلى قائمة الانتظار. أي أنه يمكنك رفع الكبسولات برمز جديد ، والذي سيتعامل مع الرسائل الجديدة ، وستستمر العمليات الحالية في الكبسولات القديمة. هذا يعطينا تحديث سلس.
  4. يمكنك استخدام التعليمات البرمجية غير المتزامنة وأي بنية أساسية ، بينما تكون مكدسًا بشكل مستقل.
  5. يمكنك التحكم في تنفيذ المهام على المستوى الأصلي ack/ المستوى reject، وكذلك الحصول على قائمة انتظار اختيارية إضافية (قائمة انتظار التحكم) يمكنها تتبع دورة حياة المهام.


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



def rmq_scheduler(routing_key_for_delay_queue, routing_key_for_processing_queue):
    def decorator(func):
        @wraps(func)
        async def wrapper(channel, body, envelope, properties):
            try:
                res = await func(channel, body, envelope, properties)
                await channel.publish(
                    payload=body,
                    exchange_name='',
                    routing_key=routing_key_for_delay_queue,
                )
                await channel.basic_client_ack(envelope.delivery_tag)
                return res
            except Exception as e:
                log_error(e)
                redelivered_count = get_count_of_redelivery_attempts(properties)
                if redelivered_count <= 3:
                    await resend_msg(
                        channel=channel,
                        body=body,
                        properties=properties,
                        routing_key=routing_key_for_processing_queue)
                else:
                    async with app.natalya_db_engine.acquire() as conn:
                        async with conn.begin():
                            await channel.publish(
                                payload=body,
                                exchange_name='',
                                routing_key=routing_key_for_delay_queue,
                            )
                await channel.basic_client_ack(envelope.delivery_tag)

        return wrapper

    return decorator


نستخدم الآن هذا المخطط لأداء المهام المتسلسلة الدورية فقط ، ولكن يمكن استخدامه أيضًا عندما يكون من المهم البدء في تنفيذ مهمة في وقت محدد ، دون تحويل الوقت إلى التنفيذ نفسه. للقيام بذلك ، يكفي إعادة جدولة المهمة بعد وصول الرسالة إلى المشرف.



صحيح أن هذا النهج له تكاليف عامة إضافية. يجب أن تفهم أنه في حالة حدوث خطأ ، ستعود الرسالة إلى قائمة الانتظار ، وسيقوم عامل آخر باستلامها والبدء في تنفيذها على الفور. لذلك ، تحتاج إلى فصل معالجة الأخطاء وفقًا لدرجة الأهمية والتفكير مسبقًا فيما يجب فعله بالرسالة في حالة حدوث هذا الخطأ أو ذاك.



الخيارات الممكنة:



  1. سيصلح الخطأ نفسه (على سبيل المثال ، خطأ في النظام): إرسال noackوتكرار معالجة الأخطاء.
  2. خطأ منطق العمل: تحتاج إلى مقاطعة الدورة - إرسال ack.
  3. يتكرر الخطأ من النقطة 1 كثيرًا: نحن نسمم rejectالمطورين ونشير إليهم . هناك خيارات هنا. يمكنك إنشاء قائمة انتظار للرسائل المراد إيداعها لإرجاع الرسالة بعد التحليل ، أو يمكنك استخدام تقنية إعادة المحاولة (حدد message_ttl).


مثال مصمم:



def auto_ack_or_nack(log_message):
   def decorator(func):
       @wraps(func)
       async def wrapper(channel, body, envelope, properties):
           try:
               res = await func(channel, body, envelope, properties)
               await channel.basic_client_ack(envelope.delivery_tag)
               return res
           except Exception as e:
               await channel.basic_client_nack(envelope.delivery_tag, requeue=False)
               log_error(log_message, exception=e)
 
       return wrapper
 
   return decorator


يعمل هذا المخطط معنا لمدة نصف عام ، وهو موثوق تمامًا ولا يتطلب عمليا الاهتمام. لا يؤدي تعطل التطبيق إلى تعطيل المجدول ويؤخر تنفيذ المهام بشكل طفيف.



لا توجد إيجابيات بدون سلبيات. يحتوي هذا المخطط أيضًا على نقطة ضعف خطيرة. إذا حدث شيء ما لـ RabbitMQ واختفت الرسائل ، فأنت بحاجة إلى إلقاء نظرة يدويًا على ما تم فقده وبدء الحلقة مرة أخرى. لكن هذا موقف غير محتمل للغاية حيث سيتعين عليك التفكير في هذه الخدمة أخيرًا :)



ملاحظة : إذا كان موضوع جدولة المهام الدورية مثيرًا للاهتمام بالنسبة لك ، ففي المقالة التالية ، سأخبرك بمزيد من التفصيل عن كيفية أتمتة إنشاء قائمة الانتظار ، وكذلك حول المشرف.



الروابط:






All Articles