Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-04)
每日一句正能量
懂得感恩的人,才能懂得生活最美好之处,也能常与安然相伴。
以“感恩”安顿内心,以“归零”保持活力,以“不悔”笃定前行。愿你带着这份心境,在属于自己的节奏里,一边接纳,一边前行。
6.4 Kafka生产者消费者实例
6.4.1 基于命令行方式使用Kafka
命令行操作是使用Kafka最基本的方式,也是便于初学者入门使用。要想建立生产者和消费者互相通信,就必须先创建一个“公共频道”, 它就是我们所说的主题(Topic), 在Kafka解压包的bin目录下, 有一个kafka-topics.sh文件,通过该文件就可以操作与主题组件相关的功能,由于前面我们配置了环境变量,所以可以在任何目录下访问bin目录下的所有文件。
- 创建主题
下面首先创建一个名为"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地址与端口号。
结果如下图所示:

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

- 消费主题中的消息
当光标出现闪烁,表示在等待输入,这时,切换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"主题中的全部内容, 我们可以根据业务需求判断是否需要添加该参数。
结果如下图所示:



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

- 删除当前主题
当想要删除当前主题时,只需要输入以下命令。
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操作方式。
- 创建工程,添加依赖
创建一个名为“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包。结果如下所示:




- 编写生产者客户端
打开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();
}
}
结果如下图所示:


- 编写消费者客户端
接下来,通过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、准备配置文件
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
欢迎 👍点赞✍评论⭐收藏,欢迎指正
更多推荐
所有评论(0)