0Pricing
PHP Academy · レッスン

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フィードバックを取得できます。ローカル設定は不要です。

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

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