尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
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欢迎 点赞✍评论⭐收藏欢迎指正
RELATED

相关推荐

5分钟搞懂res-downloader:一个代理开关,把全网视频音频“接“进本地

5分钟搞懂res-downloader:一个代理开关,把全网视频音频“接“进本地

5分钟搞懂res-downloader:一个代理开关,把全网视频音频"接"进本地 【免费下载链接】res-downloader 视频号、小程序、抖音、快手、小红书、直播流、m3u8、酷狗、QQ音乐等常见网络资源下载! 项目地址: https://gitcode.com/GitHub_Trending/r…

📅 2026/10/7 15:26:29
B站缓存M4S转MP4零基础教程:无损合并音视频文件,5秒找回你的离线视频

B站缓存M4S转MP4零基础教程:无损合并音视频文件,5秒找回你的离线视频

B站缓存M4S转MP4零基础教程:无损合并音视频文件,5秒找回你的离线视频 【免费下载链接】m4s-converter 一个跨平台小工具,将bilibili缓存的m4s格式音视频文件合并成mp4 项目地址: https://gitcode.com/gh_mirrors/m4/m4s-converter 你有…

📅 2026/10/8 4:02:08
玩日文游戏还要截图翻译?Translumo 一招让屏幕字幕原地“同声传译“

玩日文游戏还要截图翻译?Translumo 一招让屏幕字幕原地“同声传译“

玩日文游戏还要截图翻译?Translumo 一招让屏幕字幕原地"同声传译" 【免费下载链接】Translumo Advanced real-time screen translator for games, hardcoded subtitles in videos, static text and etc. 项目地址: https://gitcode.com/gh_mirrors/tr/T…

📅 2026/10/4 7:26:30
MORE NEWS

更多资讯

📰

个人智能体(Personal Agent)火爆背后:从“交互工具”到“代行中介”,重塑决策链的商业新博弈

免责声明:本文仅供产业观察与商业模式探讨,不构成任何投资建议,不涉及具体证券标的推荐。市场有风险,投资需谨慎。近期,Meta 推出 Personal Agent 产品 Muse,短时间内实现超 500 万下载量,引发市…

📰

基于FCN的腹部CT脊椎分割实战:从数据预处理到推理后处理

简介:本资源面向医学影像处理方向的深度学习学习者与研究人员,提供一套基于FCN全卷积神经网络的腹部脊椎自动分割完整方案,可用于医学图像分割入门实践与算法复现。压缩包共981个文件,约453.83MB,以jpg与png图像数据为…

📰

KNN实战ChineseMnist:15000张中文手写字的识别与调优

简介:这份资源面向机器学习入门者与中文字符识别方向的开发者,提供基于KNN算法的手写汉字识别完整实践素材。包内共2000个文件,以15000张jpg手写字符图像为主体,配合chinese_mnist.csv标签数据、Main.ipynb与Main.py主流程代码&am…

📰

《雷达原理》全套PPT课件(西安交通大学)

《雷达原理》全套PPT课件(西安交通大学) 课件内容: 第1章绪论.ppt 第2章 雷达作用距离.ppt 第3章雷达系统-ppt 第4章目标距离的测量.ppt 第5章角度测量.ppt 第6章运动目标检测及测速.ppt获取 《雷达原理》全套PPT课件(西安交通大…

📰

C# ASP.NET通讯录系统开发实战:从数据表设计到避坑指南

简介:基于C#与ASP.NET的Web通讯录管理系统源码,面向使用.NET技术栈的初学者,可作为课程设计或毕业设计的参考项目。系统以SQL Server 2005作为数据存储,实现了用户注册登录、联系人增删改查、分组树形展示、照片上传与个人信息修改…

📰

一文看懂 9 大 AI 模型:原理、落地与适用企业

如今人工智能已经不再只是会聊天的大语言模型,而是由向量模型、重排模型、语音模型、视觉模型、文生图、文生视频、图生视频等一系列专业模型共同组成的能力矩阵。不同模型各司其职,组合起来才能实现图文音视频的理解、检索、生成,下面逐一拆…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬