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欢迎 点赞✍评论⭐收藏欢迎指正

相关新闻

最新新闻

别再重复格式化U盘了!这款开源多系统启动盘工具,一个U盘装下所有系统

别再重复格式化U盘了!这款开源多系统启动盘工具,一个U盘装下所有系统

别再重复格式化U盘了!这款开源多系统启动盘工具,一个U盘装下所有系统 【免费下载链接】Ventoy A new bootable USB solution. 项目地址: https://gitcode.com/GitHub_Trending/ve/Ventoy 上周帮同事重装电脑,我在抽屉里翻出三个U盘&am…

2026/8/16 14:59:29
把十年的QQ空间说说装进Excel,GetQzonehistory一个工具就够了

把十年的QQ空间说说装进Excel,GetQzonehistory一个工具就够了

把十年的QQ空间说说装进Excel,GetQzonehistory一个工具就够了 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 凌晨两点,我对着QQ空间一页页往前翻,想…

2026/8/16 14:59:29
4.7 网络配置练习

4.7 网络配置练习

1、按照图示的VLAN及IP地址需求,完成相关配置2、要求SW1为VLAN 2/3的主根及主网关,SW2为vlan 20/30的主根及主网关,SW1和SW2互为备份3、可以使用super vian4、上层通过静态路由协议完成数据通信过5、AR1为企业出口路由器6、要求全网可达二层配…

2026/8/16 14:59:29
Vue3 + Esmx 实战:从零构建 SSR 微前端模块的完整教程

Vue3 + Esmx 实战:从零构建 SSR 微前端模块的完整教程

Vue3 Esmx 实战:从零构建 SSR 微前端模块的完整教程 【免费下载链接】genesis Next-generation micro-frontend framework based on ESM, sandbox-free with zero runtime overhead, supporting multi-framework hybrid development 项目地址: https://gitcode.c…

2026/8/16 14:59:29
【Scrapy】Scrapy教程11——XPath详解

【Scrapy】Scrapy教程11——XPath详解

前面我们简单的先了解XPath的应用,这节我们来详细学习下XPath,XPath主要是XML的查询语言,因此这节内容和 Scrapy 关系不大,不过XPath确在scrapy、BS4等工具中有很广泛的应用,本文参考菜鸟教程的XPath教程编写,由于本人能力有限,如果文中有什么错误欢迎指正。 简介 XPa…

2026/8/16 14:59:29
Excel理财应用07-Excel 投资表怎么防错?数据验证+条件格式+保护三件套拦住 90% 低级事故

Excel理财应用07-Excel 投资表怎么防错?数据验证+条件格式+保护三件套拦住 90% 低级事故

本篇定位:Excel 投资系列第 07 篇。投资表里 90% 的事故不是"算错",而是"录错"——本文给你 3 件防错工具(数据验证 条件格式 工作表保护),把低级事故消灭在萌芽。 🔥 黄金 100 字开…

2026/8/16 14:54:28