Linux安装Kafka
作者QQ:67065435 QQ群:669756510
本站内容全部为作者原创,转载请注明出处!
开始安装
安装kafka
cd /root wget --no-check-certificate -c https://dlcdn.apache.org/kafka/3.9.0/kafka_2.13-3.9.0.tgz tar -xf kafka_2.13-3.9.0.tgz mv /root/kafka_2.13-3.9.0 /usr/local/kafka vim /etc/profile export PATH=$PATH:/usr/local/kafka/bin ESC :wq修改zookeeper配置
vim /usr/local/kafka/config/zookeeper.properties dataDir=/tmp/zookeeper clientPort=2181 maxClientCnxns=0 admin.enableServer=false # admin.serverPort=8080 ESC :wq启动服务
# 启动zookeeper /usr/local/kafka/bin/zookeeper-server-start.sh /usr/local/kafka/config/zookeeper.properties # 启动kafka /usr/local/kafka/bin/kafka-server-start.sh /usr/local/kafka/config/server.properties生产者
<?php $num = 10; $brokers = "127.0.0.1:9092";//","分隔多个 $topic_name = "user_topic"; $msg_key = null; //生产者配置 $conf = new \RdKafka\Conf(); $conf->set("metadata.broker.list", $brokers); $conf->set("queue.buffering.max.kbytes", intval(bcpow(1024, 2, 0))); $conf->set("queue.buffering.max.messages", intval(bcpow(1024, 2, 0))); //生产者初始化 $producer = new \RdKafka\Producer($conf); //主题配置 $topic_conf = new \RdKafka\TopicConf(); $topic_conf->set("auto.offset.reset", "earliest"); //主题初始化 $topic = $producer->newTopic($topic_name, $topic_conf); //循环发送$num条消息 for ($i = 0; $i < $num; $i++) { $body = json_encode([ "date_time" => date("Y-m-d H:i:s"), "timestamp" => sprintf("%.6f", microtime(true)), ], JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES); $topic->produce(RD_KAFKA_PARTITION_UA, $msg_key, $body); } while ($producer->getOutQLen() > 0) { $producer->poll(0); } jump: $result = $producer->flush(10000); if ($result !== RD_KAFKA_RESP_ERR_NO_ERROR) { goto jump; }消费者
<?php $brokers = "127.0.0.1:9092";//","分隔多个 $topic_name = "user_topic"; $group_id = "user_group"; //等待消息超时时间设为10000毫秒 $time_out = 10000; //broker配置 $conf = new \RdKafka\Conf(); $conf->set("metadata.broker.list", $brokers); $conf->set("group.id", $group_id); $conf->set("socket.keepalive.enable", "true"); $conf->set("enable.partition.eof", "false"); $conf->set("enable.auto.commit", "false"); //仅当enable.auto.commit=true时有效↓ //$conf->set("auto.commit.interval.ms", 1); $conf->set("auto.offset.reset", "earliest"); //消费者 $consumer = new \RdKafka\KafkaConsumer($conf); $consumer->subscribe([$topic_name]); while (true) { //0-永不超时 $time_out-指定超时毫秒 $message = $consumer->consume(0); //$message = $consumer->consume($time_out); if (!isset($message->err)) { continue; } if ($message->err == RD_KAFKA_RESP_ERR_NO_ERROR) { $body = $message->payload; $time = microtime(true); //开始处理消息,设置失败率为10% echo "\r\e[1;32m{$body}开始{$time}\e[0m"; if (rand(1, 100) > 10) { $time = microtime(true); echo "\r\e[1;32m{$body}完成{$time}\e[0m\n"; } else { //将消息记录到日志 //saveMsgErr($topic_name, $body); echo "\r\e[1;31m{$body}失败\e[0m\n"; } //手动提交偏移量 $consumer->commit($message); } }