AI एजेंट · पाठ

इवेंट कतारें और संदेश ब्रोकर

अलग-अलग रखे गए एजेंट संचार के लिए Redis कतारें, RabbitMQ और Kafka।

पाठ 2, कुल 4 में से13 चरण

इवेंट कतारें और संदेश ब्रोकर, CoddyKit पर AI एजेंट का एक निःशुल्क पाठ है। यह 4 में से 2वाँ पाठ है। आप नीचे पूरा पाठ निःशुल्क पढ़ सकते हैं—फिर अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर के साथ ब्राउज़र में इसका व्यावहारिक अभ्यास कर सकते हैं। यह AI एजेंट सीखने के मार्ग का हिस्सा है और आपकी प्रगति वेब तथा CoddyKit ऐप पर सिंक होती रहती है। AI एजेंट पाठ्यक्रम में कुल 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

pika के साथ RabbitMQ की मूल बातें

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) में भेजना चाहिए। विफल संदेशों को स्वचालित रूप से DLQ में भेजने के लिए RabbitMQ को कॉन्फ़िगर करें।

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 उच्च-स्तरीय शेड्यूलिंग, पुनः प्रयास और निगरानी जोड़ता है। कतार की गहराई की निगरानी सेवा-विच्छेद बनने से पहले लंबित कार्यों के बारे में चेतावनी देती है।

शुरुआत निःशुल्क

एआई शिक्षक के साथ AI एजेंट सीखें — निःशुल्क

अपने ब्राउज़र में वास्तविक कोड लिखें और चलाएँ, चौबीसों घंटे एआई शिक्षक से तुरंत सहायता पाएँ, और वेब या ऐप पर वहीं से शुरू करें जहाँ आपने छोड़ा था।

पाठ्यक्रम
60
पाठ
239

अक्सर पूछे जाने वाले प्रश्न

क्या “इवेंट कतारें और संदेश ब्रोकर” पाठ निःशुल्क है?

हाँ—“इवेंट कतारें और संदेश ब्रोकर” का पूरा पाठ यहाँ वेब पर निःशुल्क पढ़ा जा सकता है। इंटरैक्टिव अभ्यास (अंतर्निहित कोड संपादक और 24/7 एआई ट्यूटर) करने और AI एजेंट पाठ्यक्रम का बाकी हिस्सा अनलॉक करने के लिए CoddyKit PRO लें। AI एजेंट पाठ्यक्रम में कुल 4 पाठ शामिल हैं।

“इवेंट कतारें और संदेश ब्रोकर” में मैं क्या सीखूँगा?

अलग-अलग रखे गए एजेंट संचार के लिए Redis कतारें, RabbitMQ और Kafka। आप ब्राउज़र में सीधे चलाए जाने वाले व्यावहारिक कोड के साथ AI एजेंट का अभ्यास करते हैं, और पाठ पूरा करते समय 24/7 एआई ट्यूटर आपके प्रश्नों के उत्तर देता है।

क्या AI एजेंट शुरू करने के लिए मुझे किसी अनुभव की आवश्यकता है?

पहले के अनुभव की आवश्यकता नहीं है। CoddyKit पर AI एजेंट शुरुआती से लेकर उन्नत शिक्षार्थियों तक सभी के लिए व्यवस्थित किया गया है, इसलिए आप यहीं से या शुरुआत से सीखना शुरू कर सकते हैं और अपनी गति से आगे बढ़ सकते हैं। यह 4 में से 2वाँ पाठ है।

“इवेंट कतारें और संदेश ब्रोकर” पाठ पूरा करने में कितना समय लगता है?

CoddyKit का अधिकांश पाठ लगभग 5–10 मिनट में पूरा हो जाता है। हर पाठ छोटा और संवादात्मक है, इसलिए आप लगातार प्रगति करते हैं और वेब या ऐप पर वहीं से सीखना जारी रख सकते हैं जहाँ आपने छोड़ा था।

क्या मैं इस AI एजेंट पाठ में कोड लिख और चला सकता हूँ?

हाँ। हर AI एजेंट पाठ में एक अंतर्निर्मित कोड संपादक शामिल है, जिससे आप सीधे अपने ब्राउज़र में वास्तविक कोड लिख और चला सकते हैं और तुरंत एआई प्रतिक्रिया पा सकते हैं—स्थानीय सेटअप की आवश्यकता नहीं है।

इस पाठ्यक्रम के सभी पाठ

  1. एजेंट डेवलपर के लिए Async Python
  2. इवेंट कतारें और संदेश ब्रोकर
  3. गैर-अवरोधक समानांतर टूल निष्पादन
  4. Async एजेंट फ़्रेमवर्क: LangChain और उससे आगे
← AI एजेंट पर वापस जाएँ