You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 21:17:12