rdkafka是php唯一稳定、活跃、支持sasl/ssl/exactly-once的生产级kafka扩展,需手动安装librdkafka和rdkafka扩展,禁用自动提交offset并用cli常驻进程运行消费者。

PHP读Kafka用rdkafka扩展,不是php-kafka或kafka-php
PHP原生不支持Kafka,必须依赖C扩展。rdkafka是目前唯一稳定、活跃、支持SASL/SSL/Exactly-Once语义的生产级扩展,php-kafka和kafka-php都是纯PHP实现,吞吐低、不维护、不兼容新协议,线上别碰。
安装方式(Ubuntu/Debian):
sudo apt-get install librdkafka1-dev pecl install rdkafka echo "extension=rdkafka.so" >> /etc/php/*/cli/php.ini echo "extension=rdkafka.so" >> /etc/php/*/fpm/php.ini
验证是否加载成功:
php -m | grep rdkafka —— 有输出即为成功
注意:PHP 8.1+ 需用 rdkafka ≥6.0.1;若用Confluent Cloud,必须启用sasl.mechanisms=PLAIN和security.protocol=sasl_ssl,否则连不上。
消费者怎么拉取并正确提交offset
rdkafka消费者默认是自动提交(enable.auto.commit=true),但实际中容易丢数据——比如消息处理到一半进程崩溃,offset已提交,下次启动就跳过这条。
推荐手动控制提交时机,示例关键逻辑:
$conf = new \RdKafka\Conf();
$conf->set('group.id', 'my-group');
$conf->set('auto.offset.reset', 'latest');
$conf->set('enable.auto.commit', 'false'); // 关键:关掉自动提交
$conf->set('bootstrap.servers', 'kafka-broker:9092');
$consumer = new \RdKafka\Consumer($conf);
$consumer->subscribe(['my-topic']);
while (true) {
$message = $consumer->consume(1000); // 1s超时
if ($message->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
// 处理业务逻辑(这里不能抛异常!)
process_message($message->payload);
// 处理成功后才提交当前offset
$consumer->commit($message);
} elseif ($message->err === RD_KAFKA_RESP_ERR__PARTITION_EOF) {
continue; // 到分区尾部,继续轮询
} else {
throw new \Exception($message->errstr());
}
}
常见坑:
-
commit()传$message只提交该条,传null会提交所有已拉取但未提交的offset,慎用 - 不要在
process_message()里sleep或阻塞太久,否则触发session.timeout.ms(默认10s),消费者被踢出group - 如果处理失败需重试,不要
commit(),下次consume()还会拿到同一条
生产者发消息要注意ack和重试配置
默认acks=1(leader写入即返回),但网络抖动或leader切换时可能丢数据。关键配置项:
-
acks=all:等ISR中所有副本写入才返回,最安全,但延迟略高 -
retries=MAX_INT:rdkafka内部自动重试,别自己封装for循环重发 -
retry.backoff.ms=100:重试间隔,避免高频打满broker -
message.timeout.ms=30000:单条消息从发送到确认的总时限,超时会进error_cb
基础发送示例:
$conf = new \RdKafka\Conf();
$conf->set('bootstrap.servers', 'kafka-broker:9092');
$conf->set('acks', 'all');
$conf->set('retries', '2147483647');
$conf->set('retry.backoff.ms', '100');
$conf->set('message.timeout.ms', '30000');
$producer = new \RdKafka\Producer($conf);
$topic = $producer->newTopic('my-topic');
// 发送(异步,立即返回)
$topic->produce(RD_KAFKA_PARTITION_UA, 0, 'my-message', null, 0);
// 强制刷出缓冲区(可选,调试用)
$producer->flush(1000);
注意:produce()不抛异常,错误走error_cb回调,必须设:
$conf->setErrorCb(function ($kafka, $err, $reason) { error_log("Kafka error: $err $reason"); });
PHP-FPM环境下不能长期运行消费者进程
FPM是短生命周期模型,worker进程处理完请求就回收,不适合跑常驻的while(true)消费者。强行跑会导致:
- 进程被FPM的
max_execution_time或request_terminate_timeout杀死 - 每次HTTP请求都新建消费者,重复订阅、重复拉取、offset混乱
- 无法优雅重启,信号处理失效
正确做法只有两种:
① 用systemd或supervisord管理独立PHP CLI脚本,例如:
php /var/www/kafka-consumer.php
② 用消息队列桥接:Nginx → PHP接口入Redis/RabbitMQ → 后台Worker消费 → 再发Kafka。PHP只做“入口”,不碰Kafka长连接。
rdkafka的RdKafka\KafkaConsumer对象内部持有C层连接和线程,一旦被FPM回收,资源不释放,容易导致句柄泄漏或core dump。
php免费学习视频:立即使用
踏上前端学习之旅,开启通往精通之路!从前端基础到项目实战,循序渐进,一步一个脚印,迈向巅峰!











