Linux安装Kafka

  1. 作者QQ:67065435 QQ群:669756510

  2. 本站内容全部为作者原创,转载请注明出处!

开始安装

  1. 安装java

  2. 安装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
    
  3. 修改zookeeper配置

    vim /usr/local/kafka/config/zookeeper.properties
    
    dataDir=/tmp/zookeeper
    clientPort=2181
    maxClientCnxns=0
    admin.enableServer=false
    # admin.serverPort=8080
    
    ESC
    :wq
    
  4. 启动服务

    # 启动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
    
  5. 生产者

    <?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;
    }
    
  6. 消费者

    <?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);
        }
    }
    
Copyright © 鲸小鱼 2012-∞ all right reserved,powered by Gitbook修订: 2025-08-27 11:36

results matching ""

    No results matching ""