如何在Laradock中使用rdkafka功能?求使用示例
Laradock中使用rdkafka的实操示例
前提确认
确保Laradock的Kafka、Zookeeper容器已启动,且rdkafka扩展安装成功:
- 进入Workspace容器:
docker-compose exec workspace bash - 执行
php -m | grep rdkafka,若输出rdkafka则扩展正常可用
1. 原生PHP生产者示例
创建kafka-producer.php文件,用于发送消息到Kafka主题:
<?php // 配置Kafka broker地址(Laradock中默认kafka容器地址为kafka:9092) $conf = new RdKafka\Conf(); $conf->set('bootstrap.servers', 'kafka:9092'); // 创建生产者实例 $producer = new RdKafka\Producer($conf); // 创建主题实例 $topic = $producer->newTopic('test-topic'); // 发送消息(RD_KAFKA_PARTITION_UA表示自动选择分区) for ($i = 0; $i < 5; $i++) { $message = "测试消息 " . $i; $topic->produce(RD_KAFKA_PARTITION_UA, 0, $message); $producer->poll(0); echo "已发送消息: {$message}\n"; } // 等待消息发送完成 for ($flushRetries = 0; $flushRetries < 10; $flushRetries++) { $result = $producer->flush(1000); if (RD_KAFKA_RESP_ERR_NO_ERROR === $result) { break; } } if (RD_KAFKA_RESP_ERR_NO_ERROR !== $result) { trigger_error('消息发送失败', E_USER_ERROR); }
运行方式:在Workspace容器中执行php kafka-producer.php
2. 原生PHP消费者示例
创建kafka-consumer.php文件,用于订阅并消费Kafka主题的消息:
<?php $conf = new RdKafka\Conf(); $conf->set('bootstrap.servers', 'kafka:9092'); // 设置消费者组ID $conf->set('group.id', 'test-group'); // 创建消费者实例 $consumer = new RdKafka\KafkaConsumer($conf); // 订阅主题 $consumer->subscribe(['test-topic']); echo "开始监听test-topic主题的消息...\n"; while (true) { $message = $consumer->consume(120*1000); // 超时时间120秒 switch ($message->err) { case RD_KAFKA_RESP_ERR_NO_ERROR: echo "收到消息: {$message->payload}\n"; break; case RD_KAFKA_RESP_ERR__PARTITION_EOF: echo "当前分区无新消息\n"; break; case RD_KAFKA_RESP_ERR__TIMED_OUT: echo "消费超时\n"; break; default: trigger_error($message->errstr(), E_USER_ERROR); break; } }
运行方式:在Workspace容器中执行php kafka-consumer.php,保持进程运行即可持续消费
3. Laravel框架中的使用示例(Artisan命令)
如果你的项目是Laravel,可以封装成Artisan命令来使用:
生产者命令
创建app/Console/Commands/KafkaProducer.php:
<?php namespace App\Console\Commands; use Illuminate\Console\Command; use RdKafka\Producer; use RdKafka\Conf; class KafkaProducer extends Command { protected $signature = 'kafka:produce {message}'; protected $description = '发送消息到Kafka主题'; public function handle() { $conf = new Conf(); $conf->set('bootstrap.servers', 'kafka:9092'); $producer = new Producer($conf); $topic = $producer->newTopic('test-topic'); $message = $this->argument('message'); $topic->produce(RD_KAFKA_PARTITION_UA, 0, $message); $producer->poll(0); $producer->flush(1000); $this->info("消息已发送: {$message}"); } }
运行方式:在Workspace容器中执行php artisan kafka:produce "你的测试消息"
消费者命令
创建app/Console/Commands/KafkaConsumer.php:
<?php namespace App\Console\Commands; use Illuminate\Console\Command; use RdKafka\KafkaConsumer; use RdKafka\Conf; class KafkaConsumer extends Command { protected $signature = 'kafka:consume'; protected $description = '消费Kafka主题的消息'; public function handle() { $conf = new Conf(); $conf->set('bootstrap.servers', 'kafka:9092'); $conf->set('group.id', 'laravel-test-group'); $consumer = new KafkaConsumer($conf); $consumer->subscribe(['test-topic']); $this->info('开始监听test-topic主题...'); while (true) { $message = $consumer->consume(60*1000); switch ($message->err) { case RD_KAFKA_RESP_ERR_NO_ERROR: $this->line("收到消息: {$message->payload}"); break; case RD_KAFKA_RESP_ERR__PARTITION_EOF: $this->line('当前分区无新消息'); break; default: $this->error($message->errstr()); break; } } } }
运行方式:在Workspace容器中执行php artisan kafka:consume
内容的提问来源于stack exchange,提问作者wuku
相关产品推荐
相关产品推荐

