Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05) 📅 2026/8/16 13:59:22 文章目录每日一句正能量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欢迎 点赞✍评论⭐收藏欢迎指正