AI-agenter · leksjon

Hendelseskøer og meldingsmeglere

Redis-køer, RabbitMQ og Kafka for frakoblet kommunikasjon mellom agenter.

Leksjon 2 av 413 trinn

Hendelseskøer og meldingsmeglere er en gratis leksjon i AI-agenter på CoddyKit. Dette er leksjon 2 av 4. Du kan lese hele leksjonen gratis nedenfor – og deretter øve praktisk i nettleseren med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt. Den er en del av læringsløpet i AI-agenter, og fremdriften din synkroniseres mellom nettet og CoddyKit-appen. Kurset i AI-agenter inneholder totalt 4 leksjoner.

Hvorfor bruke meldingskøer for agenter?

Meldingskøer avkobler utløseren (hendelsesprodusenten) fra agenten (hendelseskonsumenten). Produsenten sender hendelser i valgfritt tempo, mens agenten behandler dem i sitt eget tempo. Dette forhindrer overbelastning og muliggjør nye forsøk.

Redis-lister som køer

Redis-lister fungerer som enkle køer: LPUSH legger til fremst (enqueue), mens BRPOP fjerner fra bakerst og blokkerer hvis listen er tom (dequeue). Dette er en pålitelig først inn, først ut-kø.

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

Agentens arbeidersløyfe med Redis

En agentarbeider spør kontinuerlig køen etter nye oppgaver, behandler hver oppgave og gjentar dette. Hvis ingen oppgaver kommer innen tidsavbruddet, gjentas løkken uten å bruke ressurser.

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 for hendelser

Redis Pub/Sub fungerer annerledes enn lister: den kringkaster meldinger til alle aktive abonnenter. Bruk Pub/Sub for kringkastingsvarsler der flere agenter må reagere på den samme hendelsen.

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

Grunnleggende om RabbitMQ med pika

RabbitMQ er en meldingsmegler med omfattende funksjonalitet, blant annet ruting, kvitteringer og køer for meldinger som er sendt til side. Biblioteket pika kobler Python til 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-konsument med kvitteringer

Konsumenter må kvittere for meldinger etter behandlingen. Hvis en agent krasjer før den kvitterer, legger RabbitMQ meldingen tilbake i køen for en annen konsument. Dette sikrer at ingen meldinger går tapt.

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

Køer for meldinger som er sendt til side

Meldinger som ikke kan behandles etter for mange forsøk, bør sendes til en kø for meldinger som er sendt til side (DLQ) for manuell kontroll i stedet for å prøves på nytt i det uendelige. Konfigurer RabbitMQ til å rute mislykkede meldinger til DLQ automatisk.

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

Mønsteret med oppgavekø

Mønsteret med oppgavekø avkobler utløseren fra kjøringen: en enkel produsent legger oppgavebeskrivelser i køen, mens arbeidere henter og utfører dem. Flere arbeidere kan behandle oppgaver parallelt.

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: Oppgavekø på høyt nivå

Celery er den mest populære oppgavekøen for Python. Den støtter Redis og RabbitMQ som meldingsmeglere, automatiske nye forsøk, oppgaveruting og overvåking med Flower. Bruk den til agentarbeidsbelastninger i produksjon.

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

Overvåking av køens tilstand

Overvåk kødybden for å oppdage opphopninger. Hvis kødybden øker, trenger De flere arbeidere, eller agenten er for treg. Varsle når dybden overskrider en terskelverdi.

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

Velge mellom Redis og RabbitMQ

Bruk Redis når De trenger enkelhet, allerede bruker Redis, eller oppgavene har kort levetid og det er akseptabelt å miste noen meldinger ved en krasj. Bruk RabbitMQ når De trenger garantert levering, kompleks ruting, køer for meldinger som er sendt til side, eller flere konsumenter med ulike abonnementer.

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

Kunnskapssjekk: Meldingskøer

Test Deres forståelse av meldingskøer og meldingsmeglere for agenter.

Oppsummering av meldingskøer

Meldingskøer avkobler agentutløsere fra kjøringen: Redis-lister tilbyr enkle FIFO-køer med LPUSH/BRPOP; Redis Pub/Sub muliggjør kringkastingshendelser; RabbitMQ tilbyr garantert levering med kvitteringer og køer for meldinger som er sendt til side; Celery legger til planlegging, nye forsøk og overvåking på høyt nivå. Overvåking av kødybden varsler om opphopninger før de utvikler seg til driftsavbrudd.

Gratis å komme i gang

Lær deg AI-agenter med en AI-veileder – gratis

Skriv og kjør ekte kode i nettleseren, få umiddelbar hjelp fra en AI-veileder som er tilgjengelig døgnet rundt, og fortsett der du slapp – på nettet eller i appen.

Kurs
60
Leksjoner
239

Ofte stilte spørsmål

Er leksjonen «Hendelseskøer og meldingsmeglere» gratis?

Ja – hele teksten i «Hendelseskøer og meldingsmeglere» er gratis å lese her på nettet. For å øve interaktivt med en innebygd kodeeditor og en AI-veileder som er tilgjengelig døgnet rundt, og for å låse opp resten av AI-agenter-kurset, kan du oppgradere til CoddyKit PRO. Kurset i AI-agenter inneholder totalt 4 leksjoner.

Hva lærer jeg i «Hendelseskøer og meldingsmeglere»?

Redis-køer, RabbitMQ og Kafka for frakoblet kommunikasjon mellom agenter. Du øver på AI-agenter med praktisk kode som du kjører direkte i nettleseren, mens en AI-veileder som er tilgjengelig døgnet rundt, svarer på spørsmålene dine mens du jobber deg gjennom leksjonen.

Trenger jeg erfaring for å begynne med AI-agenter?

Ingen tidligere erfaring er nødvendig. AI-agenter på CoddyKit er lagt opp for både nybegynnere og viderekomne, så De kan begynne her eller helt fra start og lære i Deres eget tempo. Dette er leksjon 2 av 4.

Hvor lang tid tar leksjonen «Hendelseskøer og meldingsmeglere»?

De fleste CoddyKit-leksjoner tar omtrent 5–10 minutter. Hver leksjon er kort og interaktiv, slik at De gjør jevne fremskritt og kan fortsette akkurat der De slapp – både på nettet og i appen.

Kan jeg skrive og kjøre kode i denne AI-agenter-leksjonen?

Ja. Alle AI-agenter-leksjoner har en innebygd kodeeditor, slik at De kan skrive og kjøre ekte kode direkte i nettleseren og få umiddelbar tilbakemelding fra AI – uten lokal konfigurering.

Alle leksjonene i dette kurset

  1. Async Python for agentutviklere
  2. Hendelseskøer og meldingsmeglere
  3. Blokkeringsfri parallell kjøring av verktøy
  4. Asynkrone agentrammeverk: LangChain og videre
← Tilbake til AI-agenter