Un endpoint di checkout che ridimensiona l’immagine di un prodotto, invia un’email di conferma e aggiorna un motore di raccomandazione prima di rispondere con 200 OK è veloce quanto il suo passaggio più lento. Se anche uno solo va in timeout, l’intera richiesta fallisce. Una coda di messaggi permette all’endpoint di passare quel lavoro come messaggi e rispondere subito, mentre worker separati lo prendono in carico con i propri tempi.

Cosa fa davvero una coda di messaggi

Una coda di messaggi si mette tra due parti di un sistema che non hanno bisogno di essere attive nello stesso istante. Un producer pubblica un messaggio che descrive un lavoro da fare. La coda lo conserva. Un consumer lo preleva, lo elabora e lo conferma. I due non si parlano mai direttamente.

Pensa a una cassetta delle lettere. Chi spedisce lascia la lettera e se ne va; chi legge la controlla quando è libero, e se è fuori per un’ora, le lettere aspettano.

Da questo derivano tre cose che una chiamata di funzione diretta non può dare:

  • Disaccoppiamento. L’endpoint di checkout non ha bisogno di sapere come funziona il ridimensionamento delle immagini, né che è lento. Pubblica {"event": "order_placed", "order_id": 4821} e continua.
  • Bufferizzazione. Se arrivano 500 ordini nello stesso secondo, la coda assorbe il picco. I worker li elaborano al ritmo che riescono a sostenere, invece di avere 500 richieste bloccate insieme.
  • Scalabilità e guasti indipendenti. Esegui tre worker per il ridimensionamento e uno per le email, riavvia uno senza toccare l’altro, distribuisci il servizio di checkout senza ridistribuire i suoi worker.

Eseguire RabbitMQ in Docker

RabbitMQ è il message broker general-purpose più diffuso ed è una scelta ragionevole se non hai già una coda. Scarica l’immagine management per avere anche una UI web accanto al broker:

1
2
3
4
5
docker run -d --hostname mq --name rabbitmq \
  -e RABBITMQ_DEFAULT_USER=app \
  -e RABBITMQ_DEFAULT_PASS=changeme \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:4-management

La porta 5672 è il protocollo AMQP a cui si connette la tua applicazione. La porta 15672 è la UI di gestione: apri http://localhost:15672, accedi con app / changeme e puoi osservare in diretta code, frequenza dei messaggi e connessioni. Il login predefinito guest/guest funziona solo dall’interno del container, quindi imposta un utente tuo per tutto ciò a cui ti collegherai dall’host.

Se non hai mai usato Docker, Cos’è Docker e a cosa serve: guida per chi inizia spiega immagini, container e porte. Per uno stack reale, esegui il broker insieme alla tua applicazione con Compose — vedi Cos’è Docker Compose e Come si Usa:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
services:
  rabbitmq:
    image: rabbitmq:4-management
    environment:
      - RABBITMQ_DEFAULT_USER=app
      - RABBITMQ_DEFAULT_PASS=changeme
    ports:
      - "5672:5672"
      - "15672:15672"

  worker:
    build: ./worker
    depends_on:
      - rabbitmq
    environment:
      - RABBITMQ_URL=amqp://app:changeme@rabbitmq:5672/

worker raggiunge il broker all’hostname rabbitmq. Compose mette entrambi i servizi sulla stessa rete user-defined, quindi il nome del servizio funziona anche come hostname risolvibile: il meccanismo descritto in Cos’è una Rete Docker e Come si Usa.

Pubblicare e consumare un messaggio

pika è il client Python standard per il protocollo di RabbitMQ (AMQP 0-9-1). Installalo con pip install pika. Un producer che pubblica un task ha questo aspetto:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host="localhost",
        credentials=pika.PlainCredentials("app", "changeme"),
    )
)
channel = connection.channel()

channel.queue_declare(queue="resize_image", durable=True)

channel.basic_publish(
    exchange="",
    routing_key="resize_image",
    body='{"order_id": 4821, "image_url": "https://example.com/p/4821.jpg"}',
    properties=pika.BasicProperties(delivery_mode=pika.DeliveryMode.Persistent),
)
connection.close()

durable=True su queue_declare dice a RabbitMQ di mantenere la coda stessa oltre un riavvio del broker. delivery_mode=pika.DeliveryMode.Persistent gli dice di scrivere ogni messaggio su disco invece di tenerlo solo in memoria. Senza entrambi, un docker restart sul broker fa sparire in silenzio tutto ciò che non era ancora stato consumato.

Il consumer preleva dalla stessa coda e conferma ogni messaggio solo quando il lavoro è davvero finito:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host="localhost",
        credentials=pika.PlainCredentials("app", "changeme"),
    )
)
channel = connection.channel()
channel.queue_declare(queue="resize_image", durable=True)
channel.basic_qos(prefetch_count=1)

def handle_message(ch, method, properties, body):
    resize_image(body)  # il lavoro vero
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue="resize_image", on_message_callback=handle_message)
channel.start_consuming()

Due impostazioni fanno il lavoro sporco dell’affidabilità:

  • basic_qos(prefetch_count=1) impedisce a RabbitMQ di dare a un worker un secondo messaggio prima che abbia confermato il primo. Senza, un worker occupato può ritrovarsi con dieci messaggi mentre uno libero non riceve nulla.
  • basic_ack dopo il lavoro, non prima. Se il processo va in crash a metà del ridimensionamento, il messaggio non è mai stato confermato, quindi RabbitMQ lo consegna di nuovo a un altro worker invece di perderlo. Lascia auto_ack disattivato, che è già il default di pika.

Avvia tre copie dello script consumer e RabbitMQ divide la coda tra loro round-robin. Questa è tutta la storia della scalabilità per questo pattern: nessun codice di coordinamento, nessuno stato condiviso tra i worker.

Alcuni messaggi falliscono ogni volta che vengono consegnati: un payload malformato, una chiamata a un servizio ormai sparito per sempre. La riconsegna trasforma questo in un ciclo infinito tra i tuoi consumer. Imposta x-dead-letter-exchange sulla coda, e i messaggi che rifiuti senza rimetterli in coda (basic_nack con requeue=False) finiscono su un exchange separato dove li ispezioni a mano.

Come gli exchange di RabbitMQ instradano i messaggi

Gli esempi sopra pubblicano con exchange="", il default exchange, che instrada un messaggio direttamente alla coda indicata in routing_key. Copre la maggior parte dei casi d’uso da coda di task. Il modello reale di RabbitMQ mette un exchange davanti a ogni coda, e cambiare il tipo di exchange cambia il comportamento di instradamento:

Tipo di exchangeInstrada versoUsalo per
direct (il default exchange lo è)La coda che corrisponde esattamente alla routing keyCode di task, un gruppo di consumer per coda
fanoutOgni coda collegata, ignorando la routing keyTrasmettere un evento a più consumer indipendenti
topicLe code il cui pattern di binding corrisponde alla routing key (order.*.created)Trasmissione selettiva, quando i consumer vogliono un sottoinsieme di tipi di evento

Un exchange fanout è come ottieni il publish/subscribe da RabbitMQ: pubblichi order_placed una volta sola, e il servizio email, quello di analytics e quello anti-frode ricevono ciascuno la propria copia tramite la propria coda, invece di contendersi gli stessi messaggi.

Quando una coda di messaggi è lo strumento sbagliato

  • Chi chiama ha bisogno della risposta per rispondere a sua volta. Una coda serve per “fallo prima o poi”, non per “calcolalo adesso”. Se il tuo endpoint di checkout ha bisogno del costo di spedizione calcolato prima di poter rispondere, quella è una chiamata sincrona, non un job in coda.
  • Ti serve un ordine rigoroso su tutta la coda. RabbitMQ garantisce l’ordine per coda con un singolo consumer, ma aggiungi un secondo consumer per il throughput e due messaggi possono finire fuori ordine. Se la sequenza conta, come in un event log o in una macchina a stati, serve uno strumento pensato per stream ordinati, come le partizioni Kafka con chiave sull’ID dell’entità.
  • Il job è piccolo, sincrono e non gestisci già un broker. Un thread in background o un job runner in-process comportano meno superficie operativa che mettere su e monitorare RabbitMQ per un task da cron.
  • Hai già Redis e puoi tollerare qualche job perso ogni tanto. Redis Streams o una libreria come BullMQ ti danno una coda senza nuova infrastruttura, al costo di garanzie di consegna più deboli in caso di crash del broker. Controlla prima cosa fa già Redis nel tuo stack: Come Funziona il Caching con Redis e Come Usarlo.

Spostare il tuo primo task su una coda

Scegli una cosa che blocca una richiesta senza doverlo fare: l’email di conferma, la thumbnail, il ping di analytics. Esegui RabbitMQ in Docker, dichiara una coda durable per quel task, spostalo. Imposta prefetch_count=1 e conferma manualmente fin dal primo giorno, perché aggiungere affidabilità dopo che un crash ha già mangiato un messaggio costa molto più che scrivere quelle due righe adesso. Ricorri a un exchange fanout o topic solo quando hai un secondo consumer che vuole lo stesso evento.

Articoli correlati