Spark大数据分析与实战笔记(第六章 Kafka分布式发布订阅消息系统-05)
2026/8/16 13:59:18 网站建设 项目流程

文章目录

  • 每日一句正能量
  • 6.5 Kafka Streams
    • 6.5.1 Kafka Streams概述
    • 6.5.2 Kafka Streams开发单词计数

每日一句正能量

真正有格局的人,遇事从不会急于反驳,而是先会理解。
用认知的宽度,代替情绪的本能。反驳是动物的防御本能,而理解是人类的理性光辉。先理解,意味着我们愿意走出自己的视角,去看见更大的世界。
这些文案像一面镜子,照见内心宽广、懂得与世界温柔相待。带着这样的心境前行,无论遇到怎样的风景,相信都能安然欣赏,从容经过。

6.5 Kafka Streams

6.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开发单词计数

本节,将通过实时计算单词出现的次数的经典案例,分步骤讲解开发流程。

处理流程是这样的

  1. 添加依赖
    在spark_chapter06项目中, 打开pom.xm文件,添加Kafka Streams依赖,配置参数如下所示。
    文件6-5 pom.xml
<dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-streams</artifactId><version>2.0.0</version></dependency>

添加相关依赖时,要注意选择匹配当前版本号,避免兼容性问题。

结果如下图所示:

  1. 编写代码
    根据上述业务流程分析得出,单词数据通过自定义处理醋接收并执行相应业务计算,因此创建LogProcessor类, 并且继承Streams API中的Processor接口,在Processor接口中, 定义了以下三个方法:
  • Init(ProcessorContext processorContext):初始化上下文对象。
  • process(Key, Value): 每按收到一条消息时,都会洞用该方法处理并更新状态进行存储。
  • close(): 关闭处理器,这里可以做一些资源清理工作。

Kafka Strearms单词计数详田代码如文件所示。

文件6-6 LogProcessor.java

packagecn.itcast.Streams;importorg.apache.kafka.streams.processor.Processor;importorg.apache.kafka.streams.processor.ProcessorContext;importjava.util.HashMap;publicclassLogProcessorimplementsProcessor<byte[],byte[]>{//上下文对象privateProcessorContextprocessorContext;@Overridepublicvoidinit(ProcessorContextprocessorContext){//初始化方法this.processorContext=processorContext;}@Overridepublicvoidprocess(byte[]key,byte[]value){//处理一条消息StringinputOri=newString(value);HashMap<String,Integer>map=newHashMap<String,Integer>();inttimes=1;if(inputOri.contains(" ")){//截取字段String[]words=inputOri.split(" ");for(Stringword:words){if(map.containsKey(word)){map.put(word,map.get(word)+1);}else{map.put(word,times);}}}inputOri=map.toString();processorContext.forward(key,inputOri.getBytes());}@Overridepublicvoidclose(){}}

结果如下图所示:

单词计数的业务功能开发完成后,Kafka Streams需要编写一个运行主程序的类App, 来测试LogProcessor业务程序,具体代码如文件所示。
文件6-7 App.java

packagecn.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){//声明来源主题StringfromTopic="testStreams1";//声明目标主题StringtoTopic="testStreams2";//设置参数Propertiesprops=newProperties();props.put(StreamsConfig.APPLICATION_ID_CONFIG,"logProcessor");props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,"hadoop01:9092,hadoop02:9092,hadoop03:9092");//实例化StreamsConfigStreamsConfigconfig=newStreamsConfig(props);//构建拓扑结构Topologytopology=newTopology();//添加源处理节点,为源处理节点指定名称和它订阅的主题topology.addSource("SOURCE",fromTopic)//添加自定义处理节点,指定名称,处理器类和上一个节点的名称.addProcessor("PROCESSOR",newProcessorSupplier(){@OverridepublicProcessorget(){//调用这个方法,就知道这条数据用哪个process处理,returnnewLogProcessor();}},"SOURCE")//添加目标处理节点,需要指定目标处理节点的名称,和上一个节点名称。.addSink("SINK",toTopic,"PROCESSOR");//最后给SINK//实例化KafkaStreamsKafkaStreamsstreams=newKafkaStreams(topology,config);streams.start();}}

结果如下图所示:

  1. 执行测试
    代码编写完成后,在hadoop01节点创建testStreams1和testStreams2主题, 合令如下所示。
    #创建来源主题
kafka-topics.sh--create\--topictestStreams1\--partitions3\--replication-factor1\--zookeeperhadoop01:2181,hadoop02:2181,hadoop03:2181

结果如下图所示:

#创建目标主题

kafka-topics.sh--create\--topic"testStreams1"\--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
欢迎 👍点赞✍评论⭐收藏,欢迎指正

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询