メッセージキュー

queue モジュールは、RabbitMQ やその他の AMQP 互換ブローカーを含む、構成済みの分散キューにメッセージを発行し、配信を処理します。

このページは API リファレンスです。発行のスニペットでは、キューエントリと権限がすでに存在することを前提としています。コンシューマのセクションは、queue.consumer によって呼び出されるハンドラーの部分的なレシピであり、単独で動作するキューのデプロイではありません。

キューの構成については、キューを参照してください。

ロード

local queue = require("queue")

メッセージのパブリッシュ

IDでキューにメッセージを送信します:

local ok, err = queue.publish("app:tasks", {
    action = "send_email",
    user_id = 456,
    template = "welcome"
})
if err then
    return nil, err
end
パラメータ 型 説明
queue_id string キュー識別子(形式: "namespace:name")
data any メッセージデータ(テーブル、文字列、数値、ブール値)
headers table オプションのメッセージヘッダー

戻り値: boolean, error

メッセージヘッダー

ヘッダーはルーティング、優先度、トレースのメタデータを保持します。キーは文字列でなければならず、発行側の値には文字列、整数、数値、ブール値を使用できます:

local ok, err = queue.publish("app:notifications", {
    type = "order_shipped",
    order_id = order.id
}, {
    priority = 5,
    correlation_id = request_id
})
if err then return nil, err end

コンシューマはすべてのヘッダー値を文字列として受け取ります。x_original_queue、x_dead_letter_reason、x_dead_letter_time、attempts の各キーは、配信とデッドレターの記録用に予約されているため、発行側で設定してはいけません。

デリバリーコンテキストへのアクセス

キューコンシューマ内で、現在のメッセージにアクセスします:

local msg, err = queue.message()
if err then
    return nil, err
end

local msg_id, id_err = msg:id()
if id_err then return nil, id_err end
local priority, header_err = msg:header("priority")
if header_err then return nil, header_err end
local all_headers, headers_err = msg:headers()
if headers_err then return nil, headers_err end

戻り値: Message, error

コンシューマコンテキストでキューメッセージを処理する場合のみ利用可能です。

メッセージメソッド

メソッド 戻り値 説明
id() string, error 一意のメッセージ識別子
header(key) string, error 単一のヘッダー値を文字列で返す(存在しない場合nil)
headers() table, error すべてのメッセージヘッダー
ack() boolean, error 処理を確認 (single-shot)
nack() boolean, error 再配信またはデッドレターの失敗を通知 (single-shot)

ランタイムはハンドラーの成功時に自動で ack し、ハンドラーのエラー時に自動で nack します。早期に確定するときだけ ack/nack を呼び出してください。確定は 1 回限りであり、コンシューマハンドラーが戻った後の Message は無効です。

キュー情報

local stats, err = queue.info("app:tasks")
if err then return nil, err end
-- stats may contain: message_count, consumer_count, ready (driver-dependent)

戻り値: table, error

コンシューマパターン

queue.consumerエントリは、キューをハンドラ関数(funcで参照)に結び付けます。ハンドラはメッセージペイロードを直接受け取ります:

entries:
  - kind: queue.consumer
    name: email_worker
    queue: app:emails
    func: app:email_handler

このフラグメントでは、app:emails と app:email_handler の関数エントリがすでに存在することを前提としています。次の関数ソースでは、アプリケーションが deliver_email(payload) と、それに必要な権限を提供することを前提としています。

-- app:email_handler
function handle_email(payload)
    local msg = queue.message()

    logger:info("Processing", {
        message_id = message_id,
        to = payload.to
    })

    local ok, send_err = deliver_email(payload)
    if send_err then return nil, send_err end
    return ok
end

return {main = main}

呼び出しエラーを返すと、コンシューマは未確定の配信を nack します。その後の再配信は選択したドライバーの動作に従います。このリリースでは、組み込みのデッドレター構成は適用されません。

権限

キュー操作はセキュリティポリシー評価の対象です。

アクション リソース 説明
queue.publish - メッセージをパブリッシュする一般的な権限
queue.publish.queue Queue ID 特定のキューへのパブリッシュ

ランタイムは、最初に一般権限、次にキュー固有の権限を確認します。

エラー

条件 種別 再試行可能
キューIDが空 errors.INVALID いいえ
メッセージ引数がない、または空のテーブル errors.INVALID いいえ
デリバリーコンテキストがない errors.INVALID いいえ
メッセージが解放済み、または確定済み errors.INVALID いいえ
パブリッシュ不許可 errors.INVALID いいえ
パブリッシュ失敗 errors.INTERNAL いいえ
info のキューまたはドライバーが見つからない errors.INTERNAL いいえ

エラーの処理については、エラー処理を参照してください。

関連項目