Zum Inhalt springen

Queue-Konsumenten

Queue-Konsumenten verarbeiten Nachrichten aus Queues mittels Worker-Pools.

flowchart LR
subgraph Consumer[Konsument]
QD[Queue-Treiber] --> DC[Zustellkanal<br/>prefetch=10]
DC --> WP[Worker-Pool<br/>concurrency]
WP --> FH[Funktions-Handler]
FH --> AN[Ack/Nack]
end
OptionStandardMaxBeschreibung
queueErforderlich-Queue-Registry-ID
funcErforderlich-Handler-Funktions-Registry-ID
concurrency11000Worker-Anzahl
prefetch1010000Nachrichtenpuffer-Größe
auto_ackfalse-Automatisch bestätigen, bevor der Handler ausgeführt wird
driver_options{}-Treiberspezifische Consumer-Optionen
- name: order_consumer
kind: queue.consumer
queue: app:orders
func: app:process_order
concurrency: 5
prefetch: 20
lifecycle:
auto_start: true
depends_on:
- app:orders

Die Handler-Funktion empfängt den Nachrichteninhalt:

-- process_order.lua
local json = require("json")
local function handler(body)
local order = json.decode(body)
-- Bestellung verarbeiten
local result, err = process_order(order)
if err then
-- Fehler zurückgeben löst Nack aus (Requeue)
return nil, err
end
-- Erfolg löst Ack aus
return result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
ErgebnisAktionEffekt
ErfolgAckNachricht aus Queue entfernt
FehlerNackNachricht erneut eingereiht (treiberabhängig)
  • Worker laufen als nebenläufige Goroutinen
  • Jeder Worker verarbeitet eine Nachricht auf einmal
  • Nachrichten werden Round-Robin aus dem Delivery-Channel verteilt
  • Prefetch-Puffer ermöglicht es dem Treiber, voraus zu liefern
concurrency: 3
prefetch: 10
Ablauf:
1. Treiber liefert bis zu 10 Nachrichten in den Puffer
2. 3 Worker holen nebenläufig aus dem Puffer
3. Wenn Worker fertig sind, füllt sich der Puffer nach
4. Gegendruck wenn alle Worker beschäftigt und Puffer voll

Beim Stoppen:

  1. Keine neuen Lieferungen mehr annehmen
  2. Worker-Kontexte abbrechen
  3. Auf laufende Nachrichten warten (mit Timeout)
  4. Timeout-Fehler zurückgeben wenn Worker nicht fertig werden
# Queue-Treiber (Memory für Dev/Test)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# Queue-Definition
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # Namen überschreiben (Standard: Entry-Name)
codec: json # Payload-Codec (optional)
dead_letter: # Dead-Letter-Behandlung (optional)
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # Memory-Treiber: begrenzte Queue-Größe
FeldBeschreibung
queue_nameQueue-Namen überschreiben (Standard: Entry-ID-Name)
codecName des Payload-Codecs
dead_letter.queueRegistry-ID der Dead-Letter-Queue
dead_letter.max_attemptsMaximale Zustellversuche vor Weiterleitung an die DLQ
driver_optionsTreiberspezifische Einstellungen, nach Treibernamen geschlüsselt
Dead-Letter-Routing ist treiberabhängig. AMQP respektiert die DLX-Konfiguration auf Broker-Ebene; der Memory-Treiber leitet nicht an eine DLQ weiter.

Eingebaute In-Memory-Queue für Entwicklung/Tests:

  • Kind: queue.driver.memory
  • Nachrichten im Speicher gehalten
  • Nack reiht die Nachricht wieder am Ende der Queue ein
  • Keine Persistenz über Neustarts hinweg