恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05)

  • 首页
  • 资讯中心
  • /
  • Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05)

相关资讯

形式化证明:自指描述是非平衡稳态系统的热力学必然 2026/8/16 13:59:36
玉米种子色选的视觉识别逻辑:霉粒与破损粒如何稳定剔除 2026/8/16 13:59:36
认知几何学:逻辑推理与概念冲突的微分几何形式化研究报告 2026/8/16 13:59:36

最新资讯

数学建模论文写作全攻略:从模型构建到团队协作的实战指南
国赛A题定日镜场优化:从物理建模到遗传算法实现全解析
Silk v3解码完整指南:把打不开的微信语音变成MP3,从零编译到批量转换全流程
数学建模竞赛:从问题结构化到模型求解的完整实战指南
换机不丢档:BotW-Save-Manager让Switch与WiiU存档互转只需5分钟
基于Flask与Windows文件监控实现微信个人收款自动化处理

今日推荐

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码
【双层规划,节点出清价,绿证交易,CVaR方法】两级电力市场环境下计及风险的省间交易商最优购电模型附Matlab代码
隐式mpc+自适应mpc+时变mpc,线性时变模型预测控制附Simulink仿真

本周热门

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码
【双层规划,节点出清价,绿证交易,CVaR方法】两级电力市场环境下计及风险的省间交易商最优购电模型附Matlab代码
隐式mpc+自适应mpc+时变mpc,线性时变模型预测控制附Simulink仿真

本月精选

如何用DamaiHelper实现演唱会门票的智能自动化抢购:完整技术解决方案指南
第4篇:59 倍性能差距的索引瓶颈定位——一次教科书级的全表扫描调优
终极歌词批量下载神器:5分钟解决离线音乐库歌词同步难题

Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05)

发布时间:2026/8/16 13:59:36
Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05) 文章目录每日一句正能量6.5 Kafka Streams6.5.1 Kafka Streams概述6.5.2 Kafka Streams开发单词计数每日一句正能量真正有格局的人遇事从不会急于反驳而是先会理解。用认知的宽度代替情绪的本能。反驳是动物的防御本能而理解是人类的理性光辉。先理解意味着我们愿意走出自己的视角去看见更大的世界。这些文案像一面镜子照见内心宽广、懂得与世界温柔相待。带着这样的心境前行无论遇到怎样的风景相信都能安然欣赏从容经过。6.5 Kafka Streams6.5.1 Kafka Streams概述Kafka Streams是Apache Kafka开源项目的一个流处理框架它是基于Kafka的生产者和消费者,为开发者提供了流式处理的能力具有低延迟性、高扩展性、弹性、容错的特点易于集成到现有的应用程序中。Kafka Streams是一套处理分析Kafka中存储数据的客户端类库 处理完的数据可以重新写回Kafka,也可以发送给外部存储系统。作为类库可以非常方便的嵌入到应用程序中直接提供具体的类供开发者调用而且在打包和部署的过程中基本没有任何要求整个应用的运行方式主要由开发者控制方便使用和调试。在流式计算框架的模型中通常需要构建数据流的拓扑结构例如生产数据源、分析数据的处理器以及处理完成后发送的目标节点, Kafka流处理框架同样是将“输入主题-自定义处理器-输出主题’抽象成一个DAG拓扑图 如图6-15所示。图6-15 计算流程拓扑图在图6-15中生产者作为数据源不断生产和发送消息至Kafka的testStreams1主题中然后通过自定义处理器(Processor)对每条消息执行相应计算逻辑最后将结果发送到Kafka的testStreams2主题中供消费者消费消息数据。需要注意的是任务的执行拓扑图是一张有向无环图(DAG) 。有向表示从一个处理节点到另一个处理节点是具有方向性的无环表示不能有环路因为一旦有环路就会陷入死循环状态,任务将无法结束。6.5.2 Kafka Streams开发单词计数本节将通过实时计算单词出现的次数的经典案例分步骤讲解开发流程。处理流程是这样的添加依赖在spark_chapter06项目中 打开pom.xm文件添加Kafka Streams依赖配置参数如下所示。文件6-5 pom.xmldependencygroupIdorg.apache.kafka/groupIdartifactIdkafka-streams/artifactIdversion2.0.0/version/dependency添加相关依赖时要注意选择匹配当前版本号避免兼容性问题。结果如下图所示编写代码根据上述业务流程分析得出单词数据通过自定义处理醋接收并执行相应业务计算因此创建LogProcessor类 并且继承Streams API中的Processor接口在Processor接口中 定义了以下三个方法:Init(ProcessorContext processorContext):初始化上下文对象。process(Key, Value): 每按收到一条消息时都会洞用该方法处理并更新状态进行存储。close(): 关闭处理器这里可以做一些资源清理工作。Kafka Strearms单词计数详田代码如文件所示。文件6-6 LogProcessor.javapackagecn.itcast.Streams;importorg.apache.kafka.streams.processor.Processor;importorg.apache.kafka.streams.processor.ProcessorContext;importjava.util.HashMap;publicclassLogProcessorimplementsProcessorbyte[],byte[]{//上下文对象privateProcessorContextprocessorContext;Overridepublicvoidinit(ProcessorContextprocessorContext){//初始化方法this.processorContextprocessorContext;}Overridepublicvoidprocess(byte[]key,byte[]value){//处理一条消息StringinputOrinewString(value);HashMapString,IntegermapnewHashMapString,Integer();inttimes1;if(inputOri.contains( )){//截取字段String[]wordsinputOri.split( );for(Stringword:words){if(map.containsKey(word)){map.put(word,map.get(word)1);}else{map.put(word,times);}}}inputOrimap.toString();processorContext.forward(key,inputOri.getBytes());}Overridepublicvoidclose(){}}结果如下图所示单词计数的业务功能开发完成后Kafka Streams需要编写一个运行主程序的类App 来测试LogProcessor业务程序具体代码如文件所示。文件6-7 App.javapackagecn.itcast.Streams;importorg.apache.kafka.streams.KafkaStreams;importorg.apache.kafka.streams.StreamsConfig;importorg.apache.kafka.streams.Topology;importorg.apache.kafka.streams.processor.Processor;importorg.apache.kafka.streams.processor.ProcessorSupplier;importjava.util.Properties;publicclassApp{publicstaticvoidmain(String[]args){//声明来源主题StringfromTopictestStreams1;//声明目标主题StringtoTopictestStreams2;//设置参数PropertiespropsnewProperties();props.put(StreamsConfig.APPLICATION_ID_CONFIG,logProcessor);props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,hadoop01:9092,hadoop02:9092,hadoop03:9092);//实例化StreamsConfigStreamsConfigconfignewStreamsConfig(props);//构建拓扑结构TopologytopologynewTopology();//添加源处理节点为源处理节点指定名称和它订阅的主题topology.addSource(SOURCE,fromTopic)//添加自定义处理节点指定名称处理器类和上一个节点的名称.addProcessor(PROCESSOR,newProcessorSupplier(){OverridepublicProcessorget(){//调用这个方法就知道这条数据用哪个process处理returnnewLogProcessor();}},SOURCE)//添加目标处理节点需要指定目标处理节点的名称和上一个节点名称。.addSink(SINK,toTopic,PROCESSOR);//最后给SINK//实例化KafkaStreamsKafkaStreamsstreamsnewKafkaStreams(topology,config);streams.start();}}结果如下图所示执行测试代码编写完成后在hadoop01节点创建testStreams1和testStreams2主题 合令如下所示。#创建来源主题kafka-topics.sh--create\--topictestStreams1\--partitions3\--replication-factor1\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181结果如下图所示#创建目标主题kafka-topics.sh--create\--topictestStreams1\--partitions3\--replication-factor1\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181结果如下图所示成功创建好目标主题后分别在hadoop01和hadoop02 节点启动生产者服务和消费者服务。启动生产者服务的命令如下kafka-console-producer.sh\--broker-list hadoop01:9092,hadoop02:9092,hadoop03:9092\--topictestStreams1结果如下图所示在hadoop02启动消费者服务的命令如下kafka-console-consumer.sh\--from-beginning\--topicteststreams2\--bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092最后运行App主程序类。至此我们就完成了Kafka Streams所需环境的测试。在生产者服务节点(hadoop01) 中输入hello itcast hello spark hello kafka语句,返回消费者服务节点(hadoop02)中查看执行效果。转载自https://blog.csdn.net/u014727709/article/details/132865048欢迎 点赞✍评论⭐收藏欢迎指正

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号