0Pricing
AI Agents · درس

قوائم انتظار الأحداث ووسطاء الرسائل

استخدام قوائم Redis وRabbitMQ وKafka لتواصل الوكلاء غير المترابط.

قوائم انتظار الأحداث ووسطاء الرسائل درس مجاني في AI Agents على CoddyKit. هذا هو الدرس 2 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في AI Agents، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة AI Agents 4 دروس في المجموع.

لماذا نستخدم قوائم انتظار الرسائل مع الوكلاء؟

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

قوائم Redis بوصفها قوائم انتظار

تعمل قوائم Redis بوصفها قوائم انتظار بسيطة: تضيف LPUSH العناصر إلى المقدمة (إدراج في قائمة الانتظار)، بينما تزيل BRPOP العناصر من المؤخرة وتحظر التنفيذ إذا كانت القائمة فارغة (إزالة من قائمة الانتظار). وهي قائمة انتظار موثوقة تعمل وفق مبدأ الوارد أولًا يصرف أولًا.

import redis

r = redis.Redis(host='localhost', port=6379, decode_responses=True)
QUEUE_NAME = 'agent:tasks'

# Producer: add a task to the queue
def enqueue_task(task: dict):
    import json
    r.lpush(QUEUE_NAME, json.dumps(task))
    print(f'Enqueued task: {task["id"]}')

# Consumer: blocking pop (waits up to 5 seconds for a message)
def dequeue_task(timeout: int = 5):
    import json
    result = r.brpop(QUEUE_NAME, timeout=timeout)
    if result:
        queue_name, raw_data = result
        return json.loads(raw_data)
    return None  # Timed out - no tasks

# Producer side
enqueue_task({'id': 'task-1', 'type': 'email', 'email_id': 'email-abc'})
enqueue_task({'id': 'task-2', 'type': 'email', 'email_id': 'email-def'})

# Check queue length
print(f'Queue length: {r.llen(QUEUE_NAME)}')

حلقة عامل الوكيل باستخدام Redis

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

import redis
import json
import time
import logging

logger = logging.getLogger('worker')
r = redis.Redis(host='localhost', port=6379, decode_responses=True)

def process_task(task: dict) -> bool:
    task_type = task.get('type')
    if task_type == 'email':
        logger.info(f'Processing email: {task.get("email_id")}')
        # Call email agent here
        return True
    elif task_type == 'file':
        logger.info(f'Processing file: {task.get("filepath")}')
        return True
    else:
        logger.warning(f'Unknown task type: {task_type}')
        return False

def run_worker(queue_name: str = 'agent:tasks'):
    print(f'Worker started, watching queue: {queue_name}')
    while True:
        try:
            task = dequeue_task(timeout=5)
            if task:
                success = process_task(task)
                if not success:
                    # Re-queue failed tasks for retry
                    r.lpush('agent:failed', json.dumps(task))
            else:
                logger.debug('No tasks, waiting...')
        except KeyboardInterrupt:
            print('Worker stopped')
            break
        except Exception as e:
            logger.error(f'Worker error: {e}')
            time.sleep(1)  # Brief pause on unexpected error

print('Worker function defined')

Redis Pub/Sub للأحداث

يختلف Redis Pub/Sub عن القوائم؛ فهو يذيع الرسائل إلى جميع المشتركين الحاليين. استخدموا Pub/Sub لإشعارات التوزيع المتشعب عندما تحتاج عدة وكلاء إلى الاستجابة للحدث نفسه.

import redis
import json
import threading

r = redis.Redis(host='localhost', port=6379, decode_responses=True)

# Publisher: send event to all subscribers
def publish_event(channel: str, event: dict):
    r.publish(channel, json.dumps(event))
    print(f'Published to {channel}: {event}')

# Subscriber: listen for events in a background thread
def subscribe_and_handle(channel: str, handler_fn):
    pubsub = r.pubsub()
    pubsub.subscribe(channel)
    
    def listener():
        for message in pubsub.listen():
            if message['type'] == 'message':
                event = json.loads(message['data'])
                handler_fn(event)
    
    thread = threading.Thread(target=listener, daemon=True)
    thread.start()
    return thread

def handle_agent_event(event):
    print(f'Agent received event: {event}')

# Start subscriber
subscribe_and_handle('agent:events', handle_agent_event)

# Publish an event
publish_event('agent:events', {'type': 'user_action', 'action': 'login', 'user_id': 42})

import time
time.sleep(0.1)  # Give subscriber time to receive

أساسيات RabbitMQ باستخدام pika

RabbitMQ وسيط رسائل متكامل الميزات، يوفر التوجيه، وتأكيدات الاستلام، وقوائم انتظار الرسائل الميتة. وتربط مكتبة pika لغة Python بـ RabbitMQ.

import pika
import json

# Connect to RabbitMQ
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='localhost',
        port=5672,
        credentials=pika.PlainCredentials('guest', 'guest')
    )
)
channel = connection.channel()

# Declare a durable queue (survives broker restart)
channel.queue_declare(
    queue='agent_tasks',
    durable=True  # Queue survives RabbitMQ restart
)

# Publish a message
def publish_task(task: dict):
    channel.basic_publish(
        exchange='',
        routing_key='agent_tasks',
        body=json.dumps(task),
        properties=pika.BasicProperties(
            delivery_mode=2,  # Make message persistent
            content_type='application/json'
        )
    )
    print(f'Published task: {task["id"]}')

publish_task({'id': 'task-1', 'type': 'process_email', 'email_id': 'email-xyz'})
connection.close()

مستهلك RabbitMQ مع تأكيدات الاستلام

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

import pika
import json

def create_consumer():
    connection = pika.BlockingConnection(
        pika.ConnectionParameters(host='localhost')
    )
    channel = connection.channel()
    channel.queue_declare(queue='agent_tasks', durable=True)
    
    # Process one message at a time (fair dispatch)
    channel.basic_qos(prefetch_count=1)
    
    def process_message(ch, method, properties, body):
        task = json.loads(body)
        print(f'Processing: {task["id"]}')
        
        try:
            # Do agent work here
            print(f'Completed task: {task["id"]}')
            # Acknowledge: message is removed from queue
            ch.basic_ack(delivery_tag=method.delivery_tag)
        
        except Exception as e:
            print(f'Failed task {task["id"]}: {e}')
            # Negative acknowledge + requeue=True: puts message back
            ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
    
    channel.basic_consume(queue='agent_tasks', on_message_callback=process_message)
    
    print('Consumer ready. Waiting for messages...')
    channel.start_consuming()

print('RabbitMQ consumer defined')

قوائم انتظار الرسائل الميتة

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

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Create the dead letter queue first
channel.queue_declare(queue='agent_tasks_dlq', durable=True)

# Create main queue with dead-letter exchange config
channel.queue_declare(
    queue='agent_tasks',
    durable=True,
    arguments={
        'x-dead-letter-exchange': '',
        'x-dead-letter-routing-key': 'agent_tasks_dlq',
        'x-message-ttl': 3600000,  # Messages expire after 1 hour
        'x-max-delivery-count': 3  # Max 3 delivery attempts
    }
)

# Consumer: nack without requeue sends to DLQ
def careful_consumer(ch, method, properties, body):
    import json
    task = json.loads(body)
    
    try:
        # Process task
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception:
        # Do NOT requeue - send to DLQ
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

print('Dead letter queue configured')
connection.close()

نمط قائمة انتظار المهام

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

import redis
import json
from datetime import datetime

r = redis.Redis(host='localhost', port=6379, decode_responses=True)

# Task descriptor - what needs to be done
def create_agent_task(task_type: str, params: dict, priority: int = 5) -> dict:
    return {
        'id': f'{task_type}-{int(datetime.utcnow().timestamp() * 1000)}',
        'type': task_type,
        'params': params,
        'priority': priority,
        'created_at': datetime.utcnow().isoformat(),
        'retry_count': 0,
        'max_retries': 3
    }

# Enqueue with priority (multiple queues by priority)
def enqueue_with_priority(task: dict):
    priority = task.get('priority', 5)
    queue = f'agent:tasks:p{priority}'
    r.lpush(queue, json.dumps(task))
    print(f'Enqueued {task["id"]} to priority-{priority} queue')

# Worker dequeues from high-priority queue first
def dequeue_priority(timeout: int = 5):
    for p in [1, 2, 3, 4, 5]:  # Check priority 1 first
        result = r.brpop(f'agent:tasks:p{p}', timeout=0.1)
        if result:
            return json.loads(result[1])
    return None

enqueue_with_priority(create_agent_task('email_analysis', {'email_id': 'abc'}, priority=2))
enqueue_with_priority(create_agent_task('file_process', {'path': '/tmp/file.txt'}, priority=5))

Celery: قائمة انتظار مهام عالية المستوى

Celery هي قائمة انتظار المهام الأكثر شيوعًا في Python. وهي تدعم Redis وRabbitMQ بوصفهما وسيطين، وإعادة المحاولة التلقائية، وتوجيه المهام، والمراقبة باستخدام Flower. استخدموها لأحمال عمل الوكلاء في بيئات الإنتاج.

from celery import Celery
import os

# Create Celery app with Redis broker
app = Celery(
    'agent_tasks',
    broker=os.environ.get('REDIS_URL', 'redis://localhost:6379/0'),
    backend=os.environ.get('REDIS_URL', 'redis://localhost:6379/0')
)

# Configure retry behavior
app.conf.update(
    task_acks_late=True,
    task_reject_on_worker_lost=True,
    task_serializer='json',
    result_expires=3600
)

@app.task(bind=True, max_retries=3, default_retry_delay=60)
def run_email_agent(self, email_id: str):
    try:
        # Agent logic here
        print(f'Processing email: {email_id}')
        return {'status': 'success', 'email_id': email_id}
    except Exception as exc:
        raise self.retry(exc=exc)

# Trigger a task (from any Python code)
# run_email_agent.delay('email-abc')  # Fire and forget
# result = run_email_agent.apply_async(args=['email-abc'], countdown=60)  # Delayed
print('Celery task defined')

مراقبة حالة قائمة الانتظار

راقبوا عمق قائمة الانتظار لاكتشاف التراكمات. فإذا ازداد العمق، فأنتم بحاجة إلى مزيد من العمال أو أن الوكيل بطيء جدًا. أطلقوا تنبيهًا عندما يتجاوز العمق حدًا معينًا.

import redis

r = redis.Redis(host='localhost', port=6379, decode_responses=True)

def check_queue_health(queue_name: str, max_depth: int = 100) -> dict:
    depth = r.llen(queue_name)
    oldest_raw = r.lindex(queue_name, -1)  # Get last item (oldest)
    
    oldest_age_seconds = None
    if oldest_raw:
        import json
        from datetime import datetime
        oldest = json.loads(oldest_raw)
        if 'created_at' in oldest:
            created = datetime.fromisoformat(oldest['created_at'])
            oldest_age_seconds = (datetime.utcnow() - created).total_seconds()
    
    health = {
        'queue': queue_name,
        'depth': depth,
        'max_depth': max_depth,
        'overloaded': depth > max_depth,
        'oldest_item_age_seconds': oldest_age_seconds
    }
    
    if health['overloaded']:
        print(f'ALERT: Queue {queue_name} depth={depth} exceeds max={max_depth}')
    
    return health

print('Queue health monitor defined')
print('Usage: check_queue_health("agent:tasks")')

اختيار Redis أو RabbitMQ

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

# Decision guide as code comments

# Use Redis lists when:
# - Simple FIFO queue is enough
# - Redis already in stack
# - Can tolerate rare message loss on Redis crash
# - Low operational overhead matters

# Use RabbitMQ when:
# - Message acknowledgment is critical (no data loss)
# - Need complex routing (topic exchanges, fanout)
# - Need dead-letter queues for failed messages
# - Multiple consumer types need different message subsets
# - Need message TTL and expiry

# Use Celery (on top of either) when:
# - Need task scheduling and delays
# - Need automatic retries with exponential backoff
# - Need task result storage
# - Need monitoring dashboard (Flower)

print('Redis: simple, fast, ok with rare loss')
print('RabbitMQ: guaranteed delivery, complex routing')
print('Celery: high-level abstraction over both')

اختبار المعرفة: قوائم انتظار الرسائل

اختبروا مدى فهمكم لقوائم انتظار الرسائل ووسطاء الرسائل المستخدمين مع الوكلاء.

ملخص قوائم انتظار الرسائل

تفصل قوائم انتظار الرسائل بين مشغلات الوكلاء والتنفيذ: توفر قوائم Redis قوائم انتظار بسيطة وفق مبدأ الوارد أولًا يصرف أولًا باستخدام LPUSH/BRPOP؛ ويتيح Redis Pub/Sub بث الأحداث؛ ويوفر RabbitMQ تسليمًا مضمونًا مع تأكيدات الاستلام وقوائم انتظار الرسائل الميتة؛ وتضيف Celery الجدولة العالية المستوى، وإعادة المحاولة، والمراقبة. وتنبهكم مراقبة عمق قائمة الانتظار إلى التراكمات قبل أن تتحول إلى أعطال.

الأسئلة الشائعة

هل درس «قوائم انتظار الأحداث ووسطاء الرسائل» مجاني؟

نعم — نص درس «قوائم انتظار الأحداث ووسطاء الرسائل» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة AI Agents، انتقل إلى CoddyKit PRO. تتضمن دورة AI Agents 4 دروس في المجموع.

ماذا ستتعلم في «قوائم انتظار الأحداث ووسطاء الرسائل»؟

استخدام قوائم Redis وRabbitMQ وKafka لتواصل الوكلاء غير المترابط. تتمرن على AI Agents مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.

هل أحتاج إلى خبرة سابقة لأبدأ AI Agents؟

لا تُشترط خبرة سابقة. AI Agents على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 2 من أصل 4.

كم من الوقت يستغرق درس «قوائم انتظار الأحداث ووسطاء الرسائل»؟

معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.

هل يمكنني كتابة وتشغيل أكواد في درس AI Agents هذا؟

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

جميع الدروس في هذه الدورة

  1. Python غير المتزامن لمطوّري الوكلاء
  2. قوائم انتظار الأحداث ووسطاء الرسائل
  3. تنفيذ الأدوات بالتوازي دون حجب
  4. أطر الوكلاء غير المتزامنة: LangChain وما بعده
← العودة إلى AI Agents