在这里插入图片描述

每日一句正能量

懂得感恩的人,才能懂得生活最美好之处,也能常与安然相伴。
以“感恩”安顿内心,以“归零”保持活力,以“不悔”笃定前行。愿你带着这份心境,在属于自己的节奏里,一边接纳,一边前行。

6.4 Kafka生产者消费者实例

6.4.1 基于命令行方式使用Kafka

命令行操作是使用Kafka最基本的方式,也是便于初学者入门使用。要想建立生产者和消费者互相通信,就必须先创建一个“公共频道”, 它就是我们所说的主题(Topic), 在Kafka解压包的bin目录下, 有一个kafka-topics.sh文件,通过该文件就可以操作与主题组件相关的功能,由于前面我们配置了环境变量,所以可以在任何目录下访问bin目录下的所有文件。

  1. 创建主题
    下面首先创建一个名为"itcasttopic"的主题, 命令如下所示。
kafka-topics.sh --create \
--topic itcasttopic \
--partitions 3 \
--replication-factor 2 \
--zookeeper hadoop01:2181, hadoop02:2181, hadoop03:2181

上述命令创建了一个名为"itcasttopic"的主题, 该主题的分区数为3,副本数为2。关于上述命令参数的说明如下:
–create:创建一个主题。
–topic:定义主题名称。
–partitions:定义分区数。
–replication-factor:定义副本数(replication-factor(topic副本)个数不能超过broker(服务器)的个数)。
–zookeeper:指定Zookeeper服务IP地址与端口号。

结果如下图所示:
在这里插入图片描述

  1. 向主题中发送消息数据
    主题创建成功后,就可以创建生产者生产消息,用来模拟生产环境中源源不断的消息,bin目 录中的kafka-console-producer.sh文件,可以使用生产者组件相关的功能,例如向主题中发送消息数据的功能,命令如下所示。
kafka-console-producer.sh \
--broker-list hadoop01:9092, hadoop02:9092 , hadoop03:9092 \
--topic itcasttepic

结构如下图所示:
在这里插入图片描述

  1. 消费主题中的消息
    当光标出现闪烁,表示在等待输入,这时,切换hadoop02终端,创建消费者消费消息,bin目录kafka-console- consumer.sh文件,可以使用消费者组件相关的功能,例如消费主题中的消息数据的功能,命令如下所示。
kafka-console-consumer.sh \
--from-beginning --topic itcasttopic \
--bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092

上述命令中,参数–from-beginning "表示要读取"itcasttopic"主题中的全部内容, 我们可以根据业务需求判断是否需要添加该参数。

结果如下图所示:
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述

  1. 查看所有的主题
    Kafka常用命令行操作中还可以使用“–list”参数可以查看所有的主题,具体指令如下(克隆一个hadoop01会话,测试下面的指令)。
kafka-topics.sh --list \
--zookeeper hadoop01:2181,hadoop02:2181,hadoop03:2181

结果如下图所示:
在这里插入图片描述

  1. 删除当前主题
    当想要删除当前主题时,只需要输入以下命令。
kafka-topics.sh --delete \
--zookeeper hadoop01:2181,hadoop02:2181,hadoop03:2181 \
--topic itcasttopic

再用list查看,如果还能看到,表示正在使用中的主题是不能被删除的。停掉后再执行删除即可。

结果如下图所示:
在这里插入图片描述
注:删除之前切记要将生产者和消费者先关闭,否则占用资源会删除失败

6.4.2 基于Java API方式使用Kafka

用户不仅能够通过命令行的形式操作Kafka服务, Kafka还提供 了许多编程语言的客户端工具,用户在开发独立项目时,通过调用Kafka API来操作Kafka集群,其核心API主要有以下5种。

  • Producer API: 构建应用種序发送数据流到Kafka集群中的主题。
  • Consumer API:构建应用程序从Kafka集群中的主题读取数据流。
  • Streams API: 构建流处理程序的库,能够处理流式数据。
  • ConnectAPI: 实现连接器,用于在Kafka和其他系统之间可扩展的、可靠的流式传输数据的工具。
  • AdminClientAPI: 构建集群管理工具, 查看Kafka集群组件信息。

在开发生产者客户端时,Producer API提供了KafkaProducer类,该类的实例化对象用来代表一个生产者进程, 生产者发送消息时,并不是直接发送给服务端,而是先在客户端中把消息存入队列中,然后由一个发送线程从队列中消费消息,并以批量的方式发送消息给服务端。

方法名称相关说明
abortTransaction()终止正在进行的事物
close()关闭这个生产者
flush()调用此方法使所有缓冲的记录立即发送
partitionsFor(java.lang.String topic)获取给定主题的分区元数据
send(ProducerRecord<K,V> record)异步发送记录到主题

表6-2 KafkaProducer常用API

生产者客户端用来向Kafka集群中发送消息,消费者客户端则是从Kafka集群中消费消息。作为分布式消息系统,Kafka支持多个生产者和多个消费者,生产者可以将消息发布到集群中不同节点的不同分区上,消费者也可以消费集群中多个节点的多个分区上的消息,消费者应用程序是由KafkaConsumer对象代表一个消费者客户端进程,KafkaConsumer类常用的方法如表所示。

方法名称相关说明
abortTransaction()终止正在进行的事物
close()关闭这个生产者
flush()调用此方法使所有缓冲的记录立即发送
partitionsFor(java.lang.String topic)获取给定主题的分区元数据
send(ProducerRecord<K,V> record)异步发送记录到主题

表6-3 KafkaConsumer常用API

接下来,我们以实例演示的方式,分步骤介绍Kafka的Java API操作方式。

  1. 创建工程,添加依赖
    创建一个名为“spark_ chapter06"的Maven工程, 在pom.xml文件中添加Kafka依赖,需要注意的是,Kafka依赖需要与虚拟机安装的Kafka版本保持一致, 配置参数如下所示。
    文件6-2 pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>cn.itcast</groupId>
    <artifactId>spark_chapter06</artifactId>
    <version>1.0-SNAPSHOT</version>
    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                </configuration>
            </plugin>
        </plugins>
    </build>

    <dependencies>
        <dependency>
            <groupId>org.apache.kafka</groupId>
            <artifactId>kafka-clients</artifactId>
            <version>2.0.0</version>
        </dependency>
</project>

添加完毕后,IDEA工具会自动下载相关Jar包。结果如下所示:
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述

  1. 编写生产者客户端
    打开spark_ _chapter06工程下的Java目录,创建KafkaProducerTest文件用来实现生产消息数据并将数据发送到Kafka集群,如文件6-2所示。
    文件6-3 KafkaProducerTest.java
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;

import java.util.Properties;

public class KafkaProducerTest {
    public static void main(String[] args) {
        Properties props = new Properties();
        // 1、指定Kafka集群的主机名和端口号
        props.put("bootstrap.servers", "hadoop01:9092,hadoop02:9092,hadoop03:9092");
        // 2、指定等待所有副本节点的应答
        props.put("acks", "all");
        // 3、指定消息发送最大尝试次数
        props.put("retries", 0);
        // 4、指定一批消息处理大小
        props.put("batch.size", 16384);
        // 5、指定请求延时
        props.put("linger.ms", 1);
        // 6、指定缓存区内存大小
        props.put("buffer.memory", 33554432);
        // 7、设置key序列化
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        // 8、设置value序列化
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        // 9、生产数据
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(props);
        for (int i = 0; i < 50; i++) {
            producer.send(new ProducerRecord<String, String>("itcasttopic", Integer.toString(i), "hello world-" + i));
        }
        producer.close();
    }
}

结果如下图所示:
在这里插入图片描述
在这里插入图片描述

  1. 编写消费者客户端
    接下来,通过Kafka API创建KafkaConsumer对象,用来消费Kafka集群中名为"itcasttopic"主题的消息数据。在工程下创建KafkaConsumerTest.java文件, 代码如文件所示。
    文件6-4 KafkaConsumerTest.java
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.Callback;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;

import java.util.Arrays;
import java.util.Properties;

public class KafkaConsumerTest {
    public static void main(String[] args) {
        // 1、准备配置文件![在这里插入图片描述](https://img-blog.csdnimg.cn/aab57a9cf0e7430fbad09230dde180a1.png#pic_center)

        Properties props = new Properties();
        // 2、指定Kafka集群主机名和端口号
        props.put("bootstrap.servers", "hadoop01:9092,hadoop02:9092,hadoop03:9092");
        // 3、指定消费者组ID,在同一时刻同一消费组中只有一个线程可以去消费一个分区数据,不同的消费组可以去消费同一个分区的数据。
        props.put("group.id", "itcasttopic");
        // 4、自动提交偏移量
        props.put("enable.auto.commit", "true");
        // 5、自动提交时间间隔,每秒提交一次
        props.put("auto.commit.interval.ms", "1000");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(props);
        // 6、订阅数据,这里的topic可以是多个
        kafkaConsumer.subscribe(Arrays.asList("itcasttopic"));
        // 7、获取数据
        while (true) {
            //每隔100ms就拉去一次
            ConsumerRecords<String, String> records = kafkaConsumer.poll(100);
            for (ConsumerRecord<String, String> record : records) {
                System.out.printf("topic = %s,offset = %d, key = %s, value = %s%n", record.topic(), record.offset(), record.key(), record.value());
            }
        }
    }
}

结果如下图所示:
在这里插入图片描述
先启动生产者客户端程序,然后再启动清理费者端程序,这里会卡住,再启动一次生产者客户端程序,就可以看到接收到消费的信息了。
结果如下图所示:
在这里插入图片描述

注:这里之前创建的主题已经被删除了,要先创建主题。然后再启动生产者、消费者,卡住之后再次运行生产者


转载自:https://blog.csdn.net/u014727709/article/details/151689446
欢迎 👍点赞✍评论⭐收藏,欢迎指正

更多推荐