docker-compose.yml内容:
version: '3.1'
services:
zookeeper:
container_name: zookeeper
image: zookeeper:3.6
ports:
- 2181:2181
kafka:
image: wurstmeister/kafka
container_name: kafka
depends_on:
- zookeeper
ports:
- 9092:9092
volumes:
- /docker/kafka:/kafka
environment:
ALLOW_PLAINTEXT_LISTENER: 'yes'
KAFKA_BROKER_NO: 0
KAFKA_LISTENERS: PLAINTEXT://:9092
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://192.168.153.128:9092
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
kafka-manager:
container_name: kafka-manager
image: kafkamanager/kafka-manager:latest
ports:
- 9000:9000
environment:
ZK_HOSTS: zookeeper:2181
KAFKA_MANAGER_AUTH_ENABLED: "true"
KAFKA_MANAGER_USERNAME: admin
KAFKA_MANAGER_PASSWORD: admin
depends_on:
- kafka
- zookeeper
web管理页面:
http://192.168.153.128:9000
php生产者:
<?php
namespace app\index\controller;
class KafkaProducer
{
public function index()
{
$rk = new \RdKafka\Producer();
$rk->addBrokers("192.168.153.128"); //kafka服务器地址
$topic = $rk->newTopic("topic1"); //topic名称
for ($i = 0; $i < 5; $i++) {
$topic->produce(RD_KAFKA_PARTITION_UA, 0, "发送信息成功: $i");
$rk->poll(0);
}
while ($rk->getOutQLen() > 0) {
$rk->poll(50);
}
}
}
生产者路径:http://tp6api.local/index/kafkaproducer
php消费者:
<?php
$conf = new \RdKafka\Conf();
$conf->set('group.id', 'myConsumerGroup');
$rk = new \RdKafka\Consumer($conf);
$rk->addBrokers("192.168.153.128:9092");
$topicConf = new \RdKafka\TopicConf();
$topicConf->set('auto.commit.interval.ms', 100);
$topicConf->set('offset.store.method', 'file');
$topicConf->set('offset.store.path', sys_get_temp_dir());
$topicConf->set('auto.offset.reset', 'smallest');
$topic = $rk->newTopic("topic1", $topicConf);
// Start consuming partition 0
$topic->consumeStart(0, RD_KAFKA_OFFSET_STORED);
while (true) {
$message = $topic->consume(0, 120*10000);
switch ($message->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
//没有错误打印信息
var_dump($message);
break;
case RD_KAFKA_RESP_ERR__PARTITION_EOF:
echo "等待接收信息\n";
break;
case RD_KAFKA_RESP_ERR__TIMED_OUT:
echo "超时\n";
break;
default:
throw new \Exception($message->errstr(), $message->err);
break;
}
}
消费者路径: