コンテンツにスキップ

キューコンシューマ

キューコンシューマはワーカープールを使用してキューからメッセージを処理します。

flowchart LR
subgraph Consumer[コンシューマ]
QD[キュードライバ] --> DC[配信チャネル<br/>prefetch=10]
DC --> WP[ワーカープール<br/>concurrency]
WP --> FH[関数ハンドラ]
FH --> AN[Ack/Nack]
end
オプションデフォルト最大値説明
queue必須-キューレジストリID
func必須-ハンドラ関数レジストリID
concurrency11000ワーカー数
prefetch1010000メッセージバッファサイズ
auto_ackfalse-ハンドラ実行前に自動的にAckする
driver_options{}-ドライバ固有のコンシューマオプション
- 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

ハンドラ関数はメッセージボディを受け取ります:

-- process_order.lua
local json = require("json")
local function handler(body)
local order = json.decode(body)
-- 注文を処理
local result, err = process_order(order)
if err then
-- エラーを返すとNack(再キュー)がトリガーされる
return nil, err
end
-- 成功するとAckがトリガーされる
return result
end
return handler
- name: process_order
kind: function.lua
source: file://process_order.lua
modules:
- json
結果アクション効果
成功Ackメッセージがキューから削除される
エラーNackメッセージが再キューされる(ドライバ依存)
  • ワーカーは並行goroutineとして実行
  • 各ワーカーは一度に1つのメッセージを処理
  • メッセージは配信チャネルからラウンドロビン方式で分散
  • プリフェッチバッファによりドライバが先行して配信可能
concurrency: 3
prefetch: 10
フロー:
1. ドライバが最大10メッセージをバッファに配信
2. 3ワーカーがバッファから同時にプル
3. ワーカーが終了するとバッファが補充される
4. 全ワーカーがビジー状態でバッファが満杯のときバックプレッシャーが発生

停止時:

  1. 新しいデリバリーの受け入れを停止
  2. ワーカーコンテキストをキャンセル
  3. 処理中のメッセージを待機(タイムアウト付き)
  4. ワーカーが終了しない場合はタイムアウトエラーを返す
# キュードライバ(開発/テスト用メモリ)
- name: queue_driver
kind: queue.driver.memory
lifecycle:
auto_start: true
# キュー定義
- name: orders
kind: queue.queue
driver: app:queue_driver
queue_name: orders # 名前をオーバーライド(デフォルト: エントリ名)
codec: json # ペイロードコーデック(オプション)
dead_letter: # デッドレター処理(オプション)
queue: app:dlq
max_attempts: 5
driver_options:
memory:
max_length: 10000 # メモリドライバ: 有界キューサイズ
フィールド説明
queue_nameキュー名をオーバーライド(デフォルト: エントリID名)
codecペイロードコーデック名
dead_letter.queueデッドレターキューのレジストリID
dead_letter.max_attemptsDLQにルーティングされるまでの最大配信試行回数
driver_optionsドライバ名でキー付けされたドライバ固有の設定
デッドレタールーティングはドライバ依存です。AMQPはブローカーレベルのDLX設定を尊重します。メモリドライバはDLQへのルーティングを行いません。

開発/テスト用の組み込みインメモリキュー:

  • 種別: queue.driver.memory
  • メッセージはメモリに保存
  • Nackはメッセージをキューの末尾に再キュー
  • 再起動をまたいで永続化なし