PHPでApache Kafkaを扱う
Kafkaを使って高スループットのイベントをストリーミングします。
「PHPでApache Kafkaを扱う」はCoddyKit上の無料PHP Academyレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはPHP Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 PHP Academyコースには全4レッスンが含まれています。
PHPでKafkaを使う
Kafkaはタスクキューではなく、分散型の追記専用コミットログです。プロデューサーはトピックにレコードを追記し、コンシューマーはそれぞれのオフセットから読み取って、履歴を再生できます。そのためKafkaは、高スループットのイベントストリーム、イベントソーシング、そして1つのストリームを複数の独立したコンシューマーグループに供給する用途に適しています。
PHPでは、実績のあるCライブラリlibrdkafkaを基盤とするバインディング、ext-rdkafkaを通じてKafkaと通信します。
キューではなくログ
RabbitMQからKafkaへ移行する際には、考え方を次のように切り替えます。
- メッセージは消費されても削除されません。保持ポリシー(時間またはサイズ)によって期限切れになります。
- 各コンシューマーは独自のオフセット、つまりログ内の位置を管理します。
- トピックはパーティションに分割され、順序が保証されるのはパーティション内だけです。
- 複数のコンシューマーグループが、同じトピックをそれぞれ独立して読み取ります。
再生、複数の読み取り側へのファンアウト、または大規模なスループットが必要ならKafkaが適しています。メッセージ単位のルーティングやTTLが必要ならRabbitMQが適しています。
ext-rdkafkaのインストール
まずネイティブライブラリをインストールし、次にPECL拡張をインストールします。その後、必要に応じて高レベルのラッパーを追加します。
# Debian/Ubuntu
apt-get install -y librdkafka-dev
pecl install rdkafka
echo "extension=rdkafka.so" >> php.ini
# Optional ergonomic wrapper
composer require longlang/phpkafka # or arnaud-lb/php-rdkafka-stubs for IDEレコードの生成
RdKafka\Producerを作成し、トピックハンドルを取得してproduce()を呼び出します。キーによってレコードが配置されるパーティションが決まります。同じキーなら同じパーティションになり、順序が保持されます。終了前には必ずflush()を呼び出してください。そうしないと、バッファー内のレコードが失われます。
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('orders');
// RD_KAFKA_PARTITION_UA = let the partitioner choose by key
$key = 'order-42';
$topic->produce(RD_KAFKA_PARTITION_UA, 0, json_encode(['id' => 42]), $key);
$producer->poll(0);
$result = $producer->flush(10000); // wait up to 10s
if ($result !== RD_KAFKA_RESP_ERR_NO_ERROR) {
throw new RuntimeException('Failed to flush');
}キーによるパーティショニング
パーティショニングは、Kafkaのスケーラビリティと順序性の中心です。デフォルトのパーティショナーはレコードキーをハッシュ化します。partition = hash(key) % numPartitionsという計算です。適切なキーを選ぶことが重要です。
customerIdをキーにすると、顧客ごとのすべてのイベントが順序を保ったまま1つのパーティションに配置されます。- キーがnullの場合は、パーティション間でラウンドロビンになります(最大スループットですが、順序は保証されません)。
トピックのパーティション数を減らすことはできません。また、パーティションを追加するとハッシュの対応関係が変わります。そのため、ピーク時の並列性を考慮して、最初にパーティション数を決めておく必要があります。
コンシューマーグループとオフセット
group.idを設定して、高レベルAPIのKafkaConsumerを使用します。Kafkaはグループのメンバー間でパーティションを割り当て、メンバーが参加または離脱するとリバランスします。各メンバーは割り当てられたパーティションだけを読み取るため、水平スケーリングを簡単に実現できます。
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('group.id', 'order-emailers');
$conf->set('auto.offset.reset', 'earliest'); // start of log if no offset
$conf->set('enable.auto.commit', 'false'); // we commit manually
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['orders']);
while (true) {
$msg = $consumer->consume(5000);
if ($msg->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
handle($msg->payload);
$consumer->commit($msg); // commit offset AFTER work
}
}
function handle(string $p): void {}コミットするタイミング
オフセットのコミット方法によって、配信セマンティクスが決まります。
- 処理の後にコミットする → 少なくとも1回の配信です(コミット前にクラッシュすると、レコードが再生されます)。
- 処理の前にコミットする → 最大1回の配信です(クラッシュすると、レコードが失われます)。
自動コミット(enable.auto.commit=true)は、処理が完了したかどうかに関係なくタイマーでコミットします。便利ですが、クラッシュ時にレコードを気付かないまま失う可能性があります。正確性が重要な場合は無効にし、手動でコミットしてください。
消費エラーの処理
consume()から返されるものがすべてメッセージとは限りません。エラーコードに応じて分岐する必要があります。__PARTITION_EOFと__TIMED_OUTは失敗ではなく、正常な制御シグナルです。
<?php
$msg = $consumer->consume(2000);
switch ($msg->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
echo "Got: {$msg->payload} @ offset {$msg->offset}\n";
break;
case RD_KAFKA_RESP_ERR__PARTITION_EOF:
echo "Reached end of partition\n"; // caught up, keep polling
break;
case RD_KAFKA_RESP_ERR__TIMED_OUT:
echo "No message this poll\n";
break;
default:
throw new \Exception($msg->errstr(), $msg->err);
}スループットの調整
Kafkaのスループットはバッチ処理によって向上します。主なプロデューサー設定は次のとおりです。
linger.ms— 少し待って、1回のリクエストにより多くのレコードをまとめます(例:5〜20ミリ秒)。batch.size/queue.buffering.max.messages— バッファーを大きくして、往復回数を減らします。compression.type—lz4またはzstdを使うと、ネットワークコストを大幅に削減できます。acks— 耐久性を重視する場合はall、レイテンシーを下げる場合は1を指定します。
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('compression.type', 'lz4');
$conf->set('linger.ms', '10');
$conf->set('batch.size', '65536');
$conf->set('acks', 'all');
$producer = new RdKafka\Producer($conf);配信レポートコールバック
produce()は非同期であるため、送信に失敗してもその場では例外が発生しません。プロデューサー設定に配信レポートコールバックを登録し、各レコードの結果を確認してください。これは、PHPでプロデューサーの失敗を確実に検出する唯一の方法です。
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->setDrMsgCb(function ($producer, $msg) {
if ($msg->err) {
fwrite(STDERR, 'Delivery FAILED: ' . rd_kafka_err2str($msg->err) . "\n");
} else {
echo "Delivered to partition {$msg->partition} @ offset {$msg->offset}\n";
}
});
$producer = new RdKafka\Producer($conf);
// poll() services the callback queue; call it after producing
$producer->poll(0);PHP固有の注意点
Kafkaは長時間稼働するクライアントを前提としていますが、PHPのリクエストライフサイクルはその前提と相性がよくありません。
- プロデューサーは非同期にバッファリングするため、スクリプトの終了前には必ず
flush()を呼び出してください。そうしないとレコードが失われます。 - コンシューマーはスーパーバイザーの管理下で永続的なCLIワーカーとして実行し、Webリクエスト内では実行しないでください。
- リバランス中は消費が一時停止します。1メッセージあたりの処理を短くするか、グループから除外されないように
max.poll.interval.msを十分大きく設定してください。 log_levelを設定し、プロデューサーの配信レポートコールバックを登録して、気付きにくい失敗を検出してください。
クイックチェック
Kafkaで保証される順序
まとめ
PHPでKafkaを使う場合の要点は次のとおりです。
- Kafkaはキューではなく、再生可能なログです。コンシューマーはオフセットを管理します。
- キーによるパーティショニングによって、キー単位の順序性とスケーリングを実現できます。
- コンシューマーグループはパーティションを分担し、自動的にリバランスします。
- 少なくとも1回の配信を実現するには処理の後にオフセットをコミットし、制御が必要な場合は自動コミットを無効にします。
linger.ms、batch.size、compression.typeを調整し、必ずflush()を呼び出します。
次は、これらの基本要素を信頼性の高いイベント駆動ワークフローに組み込む方法です。
よくある質問
「PHPでApache Kafkaを扱う」レッスンは無料ですか?
はい。「PHPでApache Kafkaを扱う」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、PHP Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 PHP Academyコースには全4レッスンが含まれています。
「PHPでApache Kafkaを扱う」で何を学びますか?
Kafkaを使って高スループットのイベントをストリーミングします。 ブラウザで直接実行するハンズオンコードでPHP Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
PHP Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのPHP Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。
「PHPでApache Kafkaを扱う」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このPHP Academyレッスンでコードを書いて実行できますか?
はい。すべてのPHP Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- 非同期メッセージングが必要な理由
- PHPでRabbitMQを扱う
- PHPでApache Kafkaを扱う
- イベント駆動ワークフローの構築