0Pricing
AI Agents · Ders

Olay Kuyrukları ve Mesaj Aracıları

Ayrıştırılmış aracı iletişimi için Redis kuyrukları, RabbitMQ ve Kafka.

Olay Kuyrukları ve Mesaj Aracıları, CoddyKit'te ücretsiz bir AI Agents dersidir. Bu, 4 dersinin 2. dersidir. Aşağıdan dersin tamamını ücretsiz okuyabilir, sonra tarayıcıda yerleşik kod editörü ve 7/24 yapay zeka koçu ile uygulamalı olarak pratik yapabilirsin. Bu, AI Agents öğrenme yolunun bir parçasıdır ve ilerlemeniz web ve CoddyKit uygulaması arasında senkronize olur. AI Agents kursu toplamda 4 dersten oluşur.

Ajanlar İçin İleti Kuyrukları Neden Kullanılır?

İleti kuyrukları tetikleyiciyi (olay üreticisini) ajandan (olay tüketicisinden) ayırır. Üretici olayları herhangi bir hızda oluşturabilir; ajan bunları kendi hızında işler. Bu, aşırı yüklenmeyi önler ve yeniden denemeleri mümkün kılar.

Kuyruk Olarak Redis Listeleri

Redis listeleri basit kuyruklar olarak çalışır: LPUSH listenin başına ekler (kuyruğa alma), BRPOP listenin sonundan kaldırır ve liste boşsa engellenir (kuyruktan alma). Bu, güvenilir bir ilk giren ilk çıkar kuyruğudur.

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 ile Ajan İşçisi Döngüsü

Bir ajan işçisi, yeni görevler için kuyruğu sürekli yoklar, her görevi işler ve döngüyü sürdürür. Zaman aşımı süresi içinde hiçbir görev gelmezse kaynak tüketmeden yeniden döngüye girer.

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')

Olaylar İçin Redis Pub/Sub

Redis Pub/Sub listelerden farklıdır: iletileri o anda abone olan tüm alıcılara yayınlar. Aynı olaya birden çok ajanın tepki vermesi gereken, yayılma biçimli bildirimler için Pub/Sub kullanın.

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 ile RabbitMQ Temelleri

RabbitMQ; yönlendirme, onaylar ve ölü ileti kuyrukları sunan, kapsamlı özelliklere sahip bir ileti aracısıdır. pika kitaplığı Python'u RabbitMQ'ya bağlar.

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()

Onaylarla RabbitMQ Tüketicisi

Tüketiciler, iletileri işledikten sonra onaylamalıdır. Bir ajan onay vermeden önce çökerse RabbitMQ, iletinin başka bir tüketici tarafından işlenmesi için iletiyi kuyruğa yeniden ekler. Bu, hiçbir iletinin kaybolmamasını sağlar.

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')

Ölü İleti Kuyrukları

Çok fazla kez işlenemeyen iletiler, sonsuza kadar yeniden denenmek yerine manuel olarak incelenmek üzere bir ölü ileti kuyruğuna (DLQ) gönderilmelidir. Başarısız iletilerin DLQ'ya otomatik olarak yönlendirilmesi için RabbitMQ'yu yapılandırın.

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()

Görev Kuyruğu Deseni

Görev kuyruğu deseni, tetikleyiciyi yürütmeden ayırır: yalın bir üretici görev açıklamalarını kuyruğa iter; işçiler bunları çekip yürütür. Birden çok işçi görevleri paralel olarak işleyebilir.

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: Üst Düzey Görev Kuyruğu

Celery, en popüler Python görev kuyruğudur. Aracı olarak Redis ve RabbitMQ'yu, otomatik yeniden denemeleri, görev yönlendirmeyi ve Flower ile izlemeyi destekler. Üretim ortamındaki ajan iş yükleri için kullanın.

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')

Kuyruk Sağlığını İzleme

Bekleyen işleri saptamak için kuyruk derinliğini izleyin. Kuyruk derinliği artıyorsa daha fazla işçiye ihtiyacınız vardır veya ajan çok yavaştır. Derinlik bir eşiği aştığında uyarı oluşturun.

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 mi RabbitMQ mu?

Sadelik gerektiğinde, zaten Redis kullanıyorsanız veya görevler kısa ömürlüyse ve çökme sırasında birkaç iletinin kaybolması kabul edilebiliyorsa Redis kullanın. Garantili teslimat, karmaşık yönlendirme, ölü ileti kuyrukları veya farklı aboneliklere sahip birden çok tüketici gerektiğinde RabbitMQ kullanın.

# 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')

Bilgi Kontrolü: İleti Kuyrukları

Ajanlar için ileti kuyrukları ve aracıları ne kadar anladığınızı sınayın.

İleti Kuyrukları Özeti

İleti kuyrukları, ajan tetikleyicilerini yürütmeden ayırır: Redis listeleri LPUSH/BRPOP ile basit ilk giren ilk çıkar kuyrukları sağlar; Redis Pub/Sub yayın olaylarını mümkün kılar; RabbitMQ, onaylar ve ölü ileti kuyruklarıyla garantili teslimat sunar; Celery ise üst düzey zamanlama, yeniden denemeler ve izleme ekler. Kuyruk derinliğini izlemek, kesintiye dönüşmeden önce bekleyen işleri fark etmenizi sağlar.

Sıkça Sorulan Sorular

“Olay Kuyrukları ve Mesaj Aracıları” dersi ücretsiz mi?

Evet — “Olay Kuyrukları ve Mesaj Aracıları” dersin tüm metni burada web'de ücretsiz olarak okunabilir. Etkileşimli olarak pratik yapmak (yerleşik kod editörü ve 7/24 yapay zeka koçu) ve AI Agents kursunun geri kalanını açmak için CoddyKit PRO'ya yükselt. AI Agents kursu toplamda 4 dersten oluşur.

“Olay Kuyrukları ve Mesaj Aracıları” dersinde ne öğreneceğim?

Ayrıştırılmış aracı iletişimi için Redis kuyrukları, RabbitMQ ve Kafka. AI Agents ile uygulamalı kodu tarayıcıda doğrudan çalıştırarak pratik yaparsın ve 7/24 yapay zeka koçu dersi çalışırken sorularını yanıtlar.

AI Agents öğrenmeye başlamak için deneyim gerekli mi?

Önceden deneyim gerekmez. CoddyKit'te AI Agents, başlangıçtan ileri seviyeye kadar yapılandırıldığı için buradan başlayabilir veya başından başlayıp kendi hızında ilerleme yapabilirsin. Bu, 4 dersinin 2. dersidir.

“Olay Kuyrukları ve Mesaj Aracıları” dersi ne kadar sürer?

Çoğu CoddyKit dersi yaklaşık 5–10 dakika sürer. Her biri kısa ve etkileşimli olduğu için sabit ilerleme yaparsın ve web ile uygulama arasında tam olarak bıraktığın yerden devam edebilirsin.

Bu AI Agents dersinde kod yazıp çalıştırabilir miyim?

Evet. Her AI Agents dersi yerleşik bir kod editörü içerir, bu sayede tarayıcıda gerçek kod yazıp çalıştırabilir ve anlık yapay zeka geri bildirimi alırsın — yerel kurulum gerekli değildir.

Bu kursun tüm dersleri

  1. Aracı Geliştiricileri İçin Eşzamansız Python
  2. Olay Kuyrukları ve Mesaj Aracıları
  3. Engellemesiz Paralel Araç Çalıştırma
  4. Eşzamansız Aracı Çerçeveleri: LangChain ve Ötesi
← AI Agents Sayfasına Dön