0Pricing
PHP Academy · レッスン

PHPでRabbitMQを扱う

RabbitMQでメッセージを公開・消費します。

「PHPでRabbitMQを扱う」はCoddyKit上の無料PHP Academyレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはPHP Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 PHP Academyコースには全4レッスンが含まれています。

PHPでのRabbitMQ

RabbitMQはAMQP 0-9-1を使用するブローカーです。PHPで標準的に使われるクライアントには、純粋なPHPで実装されたphp-amqplib/php-amqplibと、C実装のext-amqpがあります。このレッスンでは、Composerが動作する環境ならどこでも実行できるphp-amqplibを使います。

AMQPの基本モデルには3つの要素があります。exchangeはメッセージを受け取り、bindingはキーでルーティングし、queueはコンシューマーが処理するメッセージを保持します。これを理解すれば、残りは細部です。

クライアントのインストール

Composerでライブラリを追加します。CLI向けPHPビルドでは一般的な、socketsとbcmathの拡張機能が必要です。

composer require php-amqplib/php-amqplib
# Connection target, e.g. amqp://guest:guest@localhost:5672/

Exchange/Queue/Bindingモデル

プロデューサーはキューに直接ではなく、必ずexchangeにメッセージを発行します。ルーティングはexchangeの種類によって決まります。

  • direct — ルーティングキーの完全一致。
  • topic — order.*.euのようなワイルドカードパターン。
  • fanout — バインドされたすべてのキューへのブロードキャスト。
  • headers — ヘッダー属性による一致。

bindingは、ルーティングパターンを使ってキューとexchangeを接続します。プロデューサーをキューのトポロジーから分離することが、exchangeレイヤーの目的です。

接続と宣言

接続を開き、チャネルを取得して、トポロジーを宣言します。durable: trueを指定すると、exchange/queueはブローカーの再起動後も存続します。宣言はべき等です — 一致する引数で既存のエンティティを宣言しても、何も実行されません。

<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;

$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();

$ch->exchange_declare('orders', 'topic', false, true, false);
$ch->queue_declare('orders.email', false, true, false, false);
$ch->queue_bind('orders.email', 'orders', 'order.created');

echo "Topology ready\n";
$ch->close();
$conn->close();

メッセージの公開

本文をAMQPMessageでラップします。delivery_mode = 2を設定すると、メッセージが永続的になります。再起動後も保持されるのは、durable queueとpersistent messageを組み合わせた場合です。どちらか一方だけでは保持されません。

<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;

$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();

$payload = json_encode(['orderId' => 42, 'total' => 19.90]);
$msg = new AMQPMessage($payload, [
    'content_type'  => 'application/json',
    'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
    'message_id'    => bin2hex(random_bytes(8)),
]);

$ch->basic_publish($msg, 'orders', 'order.created');
echo "Published\n";
$ch->close();
$conn->close();

手動Ackでの消費

basic_consumeでコールバックを登録します。確認応答を制御できるように、no_ack = falseを渡します。処理が成功した後にのみackします。失敗した場合は、requeueを指定したbasic_nackにより、RabbitMQに再配信させることができます。

<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;

$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->basic_qos(null, 10, null); // prefetch 10

$callback = function ($msg) {
    $data = json_decode($msg->getBody(), true);
    try {
        // ... do work ...
        $msg->ack();
    } catch (\Throwable $e) {
        $msg->nack(true); // requeue
    }
};

$ch->basic_consume('orders.email', '', false, false, false, false, $callback);
while ($ch->is_consuming()) {
    $ch->wait();
}

プリフェッチと公平なディスパッチ

RabbitMQはデフォルトでは、各コンシューマーの負荷を考慮せず、ラウンドロビン方式でメッセージを配信します。そのため、処理の遅いコンシューマーに作業が積み上がります。basic_qos(null, prefetch, null)を使うと、コンシューマーが保持できる未確認応答メッセージ数に上限を設定できます。

処理負荷が高く、ばらつきのあるタスクでは、prefetchを小さな値(例:1〜10)に設定します。そうすると、ブローカーは空き容量のあるコンシューマーにだけ新しい作業を送ります。これが公平なディスパッチです。

パブリッシャー確認

basic_publishはすぐに戻るため、ブローカーがメッセージを受け入れたかどうかは通知されません。確実に公開するには、パブリッシャー確認を有効にします。ブローカーは、メッセージが安全に永続化またはルーティングされた時点でackを送信します。

<?php
require 'vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;

$conn = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$ch = $conn->channel();
$ch->confirm_select(); // enable confirms on this channel

$ch->set_ack_handler(fn($m) => print("confirmed\n"));
$ch->set_nack_handler(fn($m) => print("REJECTED\n"));

$ch->basic_publish(new AMQPMessage('hi'), 'orders', 'order.created');
$ch->wait_for_pending_acks(5.0); // block until confirmed or timeout

デッドレタリング

キューのx-dead-letter-exchange引数を設定すると、拒否されたメッセージ(requeueなしのnack)や期限切れのメッセージをDLXへルーティングできます。クォーラムキューでx-delivery-limitと組み合わせると、再試行回数に自動で上限を設定できます。

<?php
use PhpAmqpLib\Wire\AMQPTable;

$args = new AMQPTable([
    'x-dead-letter-exchange'    => 'orders.dlx',
    'x-dead-letter-routing-key' => 'order.failed',
    'x-message-ttl'             => 60000, // ms before expiry
]);
// false=passive, true=durable, false=exclusive, false=autodelete, args
$ch->queue_declare('orders.email', false, true, false, false, false, $args);
$ch->queue_declare('orders.dead', false, true, false, false);
$ch->queue_bind('orders.dead', 'orders.dlx', 'order.failed');

本番環境でのワーカー実行

PHPでRabbitMQワーカーを運用するための、いくつかの実践的なルールを紹介します。

  • コンシューマーは、終了時に再起動するスーパーバイザー(systemd / Supervisor)の管理下で、長時間実行するCLIプロセスとして動かします。
  • PHPは時間の経過とともにメモリを消費するため、N件のメッセージを処理した後、またはメモリ使用量がしきい値を超えた時点でワーカーを再起動します。
  • AMQPのハートビートを送信し、SIGTERMを処理して正常にシャットダウンします(処理中のメッセージを完了してから停止します)。
  • 高可用性(HA)には、レガシーなミラーキューではなくクォーラムキューを使用します。

正常なシャットダウン

デプロイ時にはSIGTERMが送信されます。単純なワーカーはメッセージの処理途中で終了し、再配信を引き起こします。シグナルハンドラーを登録してフラグを切り替え、現在のメッセージを完了してackを送信してから、consumeループを正常に終了します。pcntl_async_signals(true)を使うと、AMQPの待機処理の合間にPHPがシグナルを配信できます。

<?php
pcntl_async_signals(true);
$running = true;
pcntl_signal(SIGTERM, function () use (&$running) {
    $running = false; // stop after the current message
    echo "SIGTERM: draining...\n";
});

while ($running && $ch->is_consuming()) {
    try {
        $ch->wait(null, false, 5); // wakes for signals
    } catch (\PhpAmqpLib\Exception\AMQPTimeoutException $e) {
        // idle tick - loop and re-check $running
    }
}
$ch->close();
echo "Stopped cleanly\n";

クイックチェック

ブローカーの再起動後も維持されるもの

まとめ

これで、PHPで実用的なRabbitMQパイプラインを構築できるようになりました。

  • exchangeに公開し、bindingを経由してqueueへルーティングします。
  • durable queue + persistent messageの組み合わせで、再起動後も維持できます。
  • 手動ackで消費し、basic_qosのprefetchを調整して公平なディスパッチを実現します。
  • 確実な送信にはパブリッシャー確認を使用し、問題のあるメッセージにはDLXを使用します。
  • ハートビートと正常なシャットダウンに対応したスーパーバイザーの下でワーカーを実行します。

次は、タスクキューではなく高スループットのストリーミングが必要な場合に使うKafkaです。

よくある質問

「PHPでRabbitMQを扱う」レッスンは無料ですか?

はい。「PHPでRabbitMQを扱う」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、PHP Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 PHP Academyコースには全4レッスンが含まれています。

「PHPでRabbitMQを扱う」で何を学びますか?

RabbitMQでメッセージを公開・消費します。 ブラウザで直接実行するハンズオンコードでPHP Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

PHP Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのPHP Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。

「PHPでRabbitMQを扱う」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このPHP Academyレッスンでコードを書いて実行できますか?

はい。すべてのPHP Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. 非同期メッセージングが必要な理由
  2. PHPでRabbitMQを扱う
  3. PHPでApache Kafkaを扱う
  4. イベント駆動ワークフローの構築
← PHP Academyに戻る