Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-04)
文章目录每日一句正能量6.4 Kafka生产者消费者实例6.4.1 基于命令行方式使用Kafka6.4.2 基于Java API方式使用Kafka每日一句正能量懂得感恩的人才能懂得生活最美好之处也能常与安然相伴。以“感恩”安顿内心以“归零”保持活力以“不悔”笃定前行。愿你带着这份心境在属于自己的节奏里一边接纳一边前行。6.4 Kafka生产者消费者实例6.4.1 基于命令行方式使用Kafka命令行操作是使用Kafka最基本的方式也是便于初学者入门使用。要想建立生产者和消费者互相通信就必须先创建一个“公共频道” 它就是我们所说的主题(Topic) 在Kafka解压包的bin目录下 有一个kafka-topics.sh文件通过该文件就可以操作与主题组件相关的功能由于前面我们配置了环境变量所以可以在任何目录下访问bin目录下的所有文件。创建主题下面首先创建一个名为itcasttopic的主题 命令如下所示。kafka-topics.sh--create\--topicitcasttopic\--partitions3\--replication-factor2\--zookeeperhadoop01:2181, hadoop02:2181, hadoop03:2181上述命令创建了一个名为itcasttopic的主题, 该主题的分区数为3,副本数为2。关于上述命令参数的说明如下:–create创建一个主题。–topic定义主题名称。–partitions定义分区数。–replication-factor定义副本数replication-factortopic副本个数不能超过broker服务器的个数。–zookeeper指定Zookeeper服务IP地址与端口号。结果如下图所示向主题中发送消息数据主题创建成功后就可以创建生产者生产消息用来模拟生产环境中源源不断的消息bin目 录中的kafka-console-producer.sh文件可以使用生产者组件相关的功能例如向主题中发送消息数据的功能命令如下所示。kafka-console-producer.sh\--broker-list hadoop01:9092, hadoop02:9092 , hadoop03:9092\--topicitcasttepic结构如下图所示消费主题中的消息当光标出现闪烁表示在等待输入这时切换hadoop02终端创建消费者消费消息bin目录kafka-console- consumer.sh文件可以使用消费者组件相关的功能例如消费主题中的消息数据的功能命令如下所示。kafka-console-consumer.sh\--from-beginning--topicitcasttopic\--bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092上述命令中参数–from-beginning 表示要读取itcasttopic主题中的全部内容, 我们可以根据业务需求判断是否需要添加该参数。结果如下图所示查看所有的主题Kafka常用命令行操作中还可以使用“–list”参数可以查看所有的主题具体指令如下克隆一个hadoop01会话测试下面的指令。kafka-topics.sh--list\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181结果如下图所示删除当前主题当想要删除当前主题时只需要输入以下命令。kafka-topics.sh--delete\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181\--topicitcasttopic再用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(ProducerRecordK,V record)异步发送记录到主题表6-2 KafkaProducer常用API生产者客户端用来向Kafka集群中发送消息消费者客户端则是从Kafka集群中消费消息。作为分布式消息系统,Kafka支持多个生产者和多个消费者,生产者可以将消息发布到集群中不同节点的不同分区上,消费者也可以消费集群中多个节点的多个分区上的消息消费者应用程序是由KafkaConsumer对象代表一个消费者客户端进程,KafkaConsumer类常用的方法如表所示。方法名称相关说明abortTransaction()终止正在进行的事物close()关闭这个生产者flush()调用此方法使所有缓冲的记录立即发送partitionsFor(java.lang.String topic)获取给定主题的分区元数据send(ProducerRecordK,V record)异步发送记录到主题表6-3 KafkaConsumer常用API接下来我们以实例演示的方式分步骤介绍Kafka的Java API操作方式。创建工程,添加依赖创建一个名为“spark_ chapter06的Maven工程, 在pom.xml文件中添加Kafka依赖,需要注意的是Kafka依赖需要与虚拟机安装的Kafka版本保持一致, 配置参数如下所示。文件6-2 pom.xml?xml version1.0 encodingUTF-8?projectxmlnshttp://maven.apache.org/POM/4.0.0xmlns:xsihttp://www.w3.org/2001/XMLSchema-instancexsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsdmodelVersion4.0.0/modelVersiongroupIdcn.itcast/groupIdartifactIdspark_chapter06/artifactIdversion1.0-SNAPSHOT/versionbuildpluginsplugingroupIdorg.apache.maven.plugins/groupIdartifactIdmaven-compiler-plugin/artifactIdconfigurationsource1.8/sourcetarget1.8/target/configuration/plugin/plugins/builddependenciesdependencygroupIdorg.apache.kafka/groupIdartifactIdkafka-clients/artifactIdversion2.0.0/version/dependency/project添加完毕后IDEA工具会自动下载相关Jar包。结果如下所示编写生产者客户端打开spark_ _chapter06工程下的Java目录创建KafkaProducerTest文件用来实现生产消息数据并将数据发送到Kafka集群如文件6-2所示。文件6-3 KafkaProducerTest.javaimportorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerRecord;importjava.util.Properties;publicclassKafkaProducerTest{publicstaticvoidmain(String[]args){PropertiespropsnewProperties();// 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、生产数据KafkaProducerString,StringproducernewKafkaProducerString,String(props);for(inti0;i50;i){producer.send(newProducerRecordString,String(itcasttopic,Integer.toString(i),hello world-i));}producer.close();}}结果如下图所示编写消费者客户端接下来通过Kafka API创建KafkaConsumer对象用来消费Kafka集群中名为itcasttopic主题的消息数据。在工程下创建KafkaConsumerTest.java文件 代码如文件所示。文件6-4 KafkaConsumerTest.javaimportorg.apache.kafka.clients.consumer.ConsumerRecord;importorg.apache.kafka.clients.consumer.ConsumerRecords;importorg.apache.kafka.clients.consumer.KafkaConsumer;importorg.apache.kafka.clients.producer.Callback;importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerRecord;importorg.apache.kafka.clients.producer.RecordMetadata;importjava.util.Arrays;importjava.util.Properties;publicclassKafkaConsumerTest{publicstaticvoidmain(String[]args){// 1、准备配置文件![在这里插入图片描述](https://img-blog.csdnimg.cn/aab57a9cf0e7430fbad09230dde180a1.png#pic_center)PropertiespropsnewProperties();// 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);KafkaConsumerString,StringkafkaConsumernewKafkaConsumerString,String(props);// 6、订阅数据这里的topic可以是多个kafkaConsumer.subscribe(Arrays.asList(itcasttopic));// 7、获取数据while(true){//每隔100ms就拉去一次ConsumerRecordsString,StringrecordskafkaConsumer.poll(100);for(ConsumerRecordString,Stringrecord: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欢迎 点赞✍评论⭐收藏欢迎指正