0Pricing
PHP Academy · 课时

使用 PHP 的 Apache Kafka

使用 Kafka 流式处理高吞吐量事件

使用 PHP 的 Apache Kafka 是 CoddyKit 上的免费 PHP Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 PHP Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 PHP Academy 课程共包含 4 节课。

使用 PHP 操作 Kafka

Kafka 不是任务队列——它是分布式、仅追加的提交日志。生产者将记录追加到主题;消费者从自己的偏移量读取,并可以重放历史记录。这使 Kafka 非常适合高吞吐事件流、事件溯源,以及从一个事件流向多个相互独立的消费者组提供数据。

在 PHP 中,您可以通过 ext-rdkafka 与 Kafka 通信;它是经过充分验证的 C 库 librdkafka 的绑定。

是日志,而非队列

从 RabbitMQ 转向 Kafka 时需要改变思维方式:

  • 消息不会在消费后删除;它们会根据保留策略(时间或大小)过期。
  • 每个消费者都会跟踪自己的偏移量——也就是它在日志中的位置。
  • 主题会拆分为分区;只有在同一分区内才能保证顺序。
  • 多个消费者组可以彼此独立地读取 the 同一主题。

如果您需要重放、向许多读取者扇出,或需要极高吞吐量,Kafka 很适合。如果您需要逐消息路由和过期时间,RabbitMQ 更适合。

安装 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()。the 键决定记录进入哪个分区——相同键进入相同分区,顺序得以保留。退出前始终调用 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 分区 → 某个客户的所有事件都保持有序并位于同一分区。
  • 空键 → 在各分区之间轮询(吞吐量最高,但不保证顺序)。

永远不能减少主题的分区数,添加分区会重新调整哈希映射——因此应预先按峰值并行度设置分区数量。

消费者组与偏移量

使用带 group.id 的高级 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 {}

何时提交

偏移量提交定义了交付语义:

  • 在处理之后提交 → 至少一次(提交前发生崩溃会重新播放记录)。
  • 在处理之前提交 → 至多一次(发生崩溃会丢失记录)。

自动提交(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 — 短暂等待,以便每个请求批量处理更多记录(例如 5–20ms)。
  • 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 请求中运行。
  • 重新平衡会暂停消费;请缩短每条消息的处理时间,或适当增大 max.poll.interval.ms,以免被移出消费者组。
  • 设置 log_level,并注册生产者的投递报告回调,以捕获静默失败。

快速检查

Kafka 中的顺序保证。

回顾

以 PHP 的方式使用 Kafka:

  • Kafka 是可重放日志,而不是队列;消费者会跟踪偏移量。
  • 按键分区可以实现按键排序和扩展。
  • 消费者组会拆分分区并自动重新平衡。
  • 在工作之后提交偏移量,以实现至少一次;禁用自动提交以便自行控制。
  • 调节 linger.ms、batch.size 和 compression.type;始终调用 flush()。

接下来,将这些基础机制连接起来,构建可靠的事件驱动工作流。

常见问题解答

「使用 PHP 的 Apache Kafka」课时是免费的吗?

是的 — 「使用 PHP 的 Apache Kafka」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 PHP Academy 课程的其余内容,请升级到 CoddyKit PRO。 PHP Academy 课程共包含 4 节课。

「使用 PHP 的 Apache Kafka」这节课中我会学到什么?

使用 Kafka 流式处理高吞吐量事件 你通过在浏览器中直接运行的动手代码来练习 PHP Academy,全天候 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