Zum Inhalt springen

Nachrichten-Queue

Veröffentlichen und Konsumieren von Nachrichten aus verteilten Queues. Unterstützt mehrere Backends einschließlich RabbitMQ und andere AMQP-kompatible Broker.

Für Queue-Konfiguration siehe Queue.

local queue = require("queue")

Senden Sie Nachrichten an eine Queue per ID:

local ok, err = queue.publish("app:tasks", {
action = "send_email",
user_id = 456,
template = "welcome"
})
if err then
return nil, err
end
ParameterTypBeschreibung
queue_idstringQueue-Identifikator (Format: “namespace:name”)
dataanyNachrichtendaten (Tables, Strings, Zahlen, Booleans)
headerstableOptionale Nachrichten-Header

Gibt zurück: boolean, error

Header ermöglichen Routing, Priorisierung und Tracing:

queue.publish("app:notifications", {
type = "order_shipped",
order_id = order.id
}, {
priority = "high",
correlation_id = request_id
})

Innerhalb eines Queue-Consumers auf die aktuelle Nachricht zugreifen:

local msg, err = queue.message()
if err then
return nil, err
end
local msg_id = msg:id()
local priority = msg:header("priority")
local all_headers = msg:headers()

Gibt zurück: Message, error

Nur verfügbar beim Verarbeiten von Queue-Nachrichten im Consumer-Kontext.

MethodeGibt zurückBeschreibung
id()string, errorEindeutiger Nachrichten-Identifikator
header(key)any, errorEinzelner Header-Wert (nil wenn fehlend)
headers()table, errorAlle Nachrichten-Header
ack()boolean, errorVerarbeitung bestaetigen (single-shot)
nack()boolean, errorFehlschlag fuer Redelivery oder Dead-Letter melden (single-shot)

Die Runtime fuehrt bei Handler-Erfolg automatisch ack aus und bei Handler-Fehler automatisch nack. Rufen Sie ack/nack nur auf, um frueher abzuschliessen.

local stats, err = queue.info("app:tasks")
-- stats kann enthalten: message_count, consumer_count, ready (treiberabhaengig)

Gibt zurueck: table, error

Queue-Consumer werden als Entry-Points definiert, die den Payload direkt empfangen:

entries:
- kind: queue.consumer
id: email_worker
queue: app:emails
method: handle_email
function handle_email(payload)
local msg = queue.message()
logger:info("Processing", {
message_id = msg:id(),
to = payload.to
})
local ok, err = email.send(payload.to, payload.template, payload.data)
if err then
return nil, err -- Nachricht wird erneut eingereiht oder dead-lettered
end
end

Queue-Operationen unterliegen der Sicherheitsrichtlinienauswertung.

AktionRessourceBeschreibung
queue.publish-Allgemeine Berechtigung zum Veröffentlichen von Nachrichten
queue.publish.queueQueue-IDZu spezifischer Queue veröffentlichen

Beide Berechtigungen werden geprüft: zuerst die allgemeine Berechtigung, dann die queue-spezifische.

BedingungArtWiederholbar
Queue-ID leererrors.INVALIDnein
Nachrichtendaten leererrors.INVALIDnein
Kein Zustellungskontexterrors.INVALIDnein
Veröffentlichung nicht erlaubterrors.INVALIDnein
Veröffentlichung fehlgeschlagenerrors.INTERNALnein

Siehe Fehlerbehandlung für die Arbeit mit Fehlern.