• 【KafkaStream】流式计算概述&KafkaStream入门


     8420b26844034fab91b6df661ae68671.png

    个人简介: 

    > 📦个人主页:赵四司机
    > 🏆学习方向:JAVA后端开发 
    > 📣种一棵树最好的时间是十年前,其次是现在!
    > ⏰往期文章:SpringBoot项目整合微信支付
    > 🧡喜欢的话麻烦点点关注喔,你们的支持是我的最大动力。

    前言:

    1.前面基于Springboot的单体项目介绍已经完结了,至于项目中的其他功能实现我这里就不打算介绍了,因为涉及的知识点不难,而且都是简单的CRUD操作,假如有兴趣的话可以私信我我再看看要不要写几篇文章做个介绍。

    2.完成上一阶段的学习,我就投入到了微服务的学习当中,所用教程为B站上面黑马的微服务教程。由于我的记性不是很好,所以对于新事物的学习我比较喜欢做笔记以加强理解,在这里我会将笔记的重点内容做个总结发布到“微服务学习”笔记栏目中。我是赵四,一名有追求的程序员,希望大家能多多支持,能给我点个关注就更好了。

    目录

    一: 流式计算

    1.概述

    2.应用场景

     3.技术方案选型

    二:Kafka Stream

    1.概述

     2.Kafka Streams的关键概念

     3.KStream

    3.1:KStream数据结构

     3.2:KStream数据流

    三:入门案例编写

    1.需求分析

    2.引入依赖

    3.创建原生的kafka staream入门案例

    3.1:开启流式计算

    3.2:生产者

    3.3:消费者

    4.测试结果


    一: 流式计算

    1.概述

            一般流式计算会与批量计算相比较。在流式计算模型中,输入是持续的,可以认为在时间上是无界的,也就意味着,永远拿不到全量数据去做计算。同时,计算结果是持续输出的,也即计算结果在时间上也是无界的。流式计算一般对实时性要求较高,同时一般是先定义目标计算,然后数据到来之后将计算逻辑应用于数据。同时为了提高计算效率,往往尽可能采用增量计算代替全量计算。

            可能上面的话看起来不好理解,我举个简单的栗子,批量计算就相当于我们平时搭乘的升降式电梯,运送人是一批一批的,中间会等待人流的积累;而流式计算则相当于我们商场里面搭乘的扶梯,人是源源不断进行运送的,它会不断接收人流,你把人理解成数据就可以理解什么是数据的流式计算了。

    2.应用场景

    流式计算其实在我们生活中十分常见,下面我举几个例子:

    • 日志分析

      网站的用户访问日志进行实时的分析,计算访问量,用户画像,留存率等等,实时的进行数据分析,帮助企业进行决策

    • 大屏看板统计

      可以实时的查看网站注册数量,订单数量,购买数量,金额等。

    • 公交实时数据

      可以随时更新公交车方位,计算多久到达站牌等

    • 实时文章分值计算

      头条类文章的分值计算,通过用户的行为实时文章的分值,分值越高就越被推荐。

     3.技术方案选型

    常用的流式计算技术方案选型有如下几种:

    • Hadoop
    • Apche Storm
    • Flink
    • Kafka Stream

            前面三种都是大数据领域常用的技术方案,因此学习成本会比较高,其相关功能大家可以自行查阅官方文档或者相关教程,这里就不多做介绍了。我的项目是基于Java开发的,需求是实现文章分值实时计算,而Kafka Stream可以轻松地将其嵌入任何Java应用程序中,并与用户为其流应用程序所拥有的任何现有打包,部署和操作工具集成。 所以我这里选用的技术方案的Kafka Stream来实现流式计算。

    二:Kafka Stream

    1.概述

    Kafka Stream是Apache Kafka从0.10版本引入的一个新Feature。它是提供了对存储于Kafka内的数据进行流式处理和分析的功能。

    Kafka Stream的特点如下:

    • Kafka Stream提供了一个非常简单而轻量的Library,它可以非常方便地嵌入任意Java应用中,也可以任意方式打包和部署

    • 除了Kafka外,无任何外部依赖

    • 充分利用Kafka分区机制实现水平扩展和顺序性保证

    • 通过可容错的state store实现高效的状态操作(如windowed join和aggregation)

    • 支持正好一次处理语义

    • 提供记录级的处理能力,从而实现毫秒级的低延迟

    • 支持基于事件时间的窗口操作,并且可处理晚到的数据(late arrival of records)

    • 同时提供底层的处理原语Processor(类似于Storm的spout和bolt),以及高层抽象的DSL(类似于Spark的map/group/reduce)

     2.Kafka Streams的关键概念

    • 源处理器(Source Processor):源处理器是一个没有任何上游处理器的特殊类型的流处理器。它从一个或多个kafka主题生成输入流。通过消费这些主题的消息并将它们转发到下游处理器。

    • Sink处理器:sink处理器是一个没有下游流处理器的特殊类型的流处理器。它接收上游流处理器的消息发送到一个指定的Kafka主题

    下面提供一幅官方的图供大家参考帮助理解:

     3.KStream

    3.1:KStream数据结构

    KStream数据结构有点像Map结构,都是Key-Value键值对,见下图:

     3.2:KStream数据流

    KStream数据流(data stream),即是一段顺序的,可以无限长,不断更新的数据集。 数据流中比较常记录的是事件,这些事件可以是一次鼠标点击(click),一次交易,或是传感器记录的位置数据。

    KStream负责抽象的,就是数据流。与Kafka自身topic中的数据一样,类似日志,每一次操作都是向其中插入(insert)新数据。

    可能上面的话不太好理解,可以看下图:

    你可能会问,消息是按序发送,为什么处理结果不是3而是5呢?这就是上面提到的第二条数据记录不被视为先前记录的更新,而是新增

    三:入门案例编写

    1.需求分析

            我们需要利用Kafka Stream实现对单词数量的统计,假如多个生产者发送多条数据,数据是一段英文,那么消费者将统计这些生产者发送的消息中不同单词的个数进行输出。

    2.引入依赖

    1. <dependency>
    2. <groupId>org.springframework.kafkagroupId>
    3. <artifactId>spring-kafkaartifactId>
    4. <exclusions>
    5. <exclusion>
    6. <groupId>org.apache.kafkagroupId>
    7. <artifactId>kafka-clientsartifactId>
    8. exclusion>
    9. exclusions>
    10. dependency>
    11. <dependency>
    12. <groupId>org.apache.kafkagroupId>
    13. <artifactId>kafka-clientsartifactId>
    14. dependency>
    15. <dependency>
    16. <groupId>org.apache.kafkagroupId>
    17. <artifactId>kafka-streamsartifactId>
    18. <exclusions>
    19. <exclusion>
    20. <groupId>org.apache.kafkagroupId>
    21. <artifactId>connect-jsonartifactId>
    22. exclusion>
    23. <exclusion>
    24. <groupId>org.apache.kafkagroupId>
    25. <artifactId>kafka-clientsartifactId>
    26. exclusion>
    27. exclusions>
    28. dependency>

    3.创建原生的kafka staream入门案例

    3.1:开启流式计算

    1. package com.my.kafka.demo;
    2. import org.apache.kafka.common.serialization.Serdes;
    3. import org.apache.kafka.streams.KafkaStreams;
    4. import org.apache.kafka.streams.KeyValue;
    5. import org.apache.kafka.streams.StreamsBuilder;
    6. import org.apache.kafka.streams.StreamsConfig;
    7. import org.apache.kafka.streams.kstream.KStream;
    8. import org.apache.kafka.streams.kstream.TimeWindows;
    9. import org.apache.kafka.streams.kstream.ValueMapper;
    10. import java.time.Duration;
    11. import java.util.Arrays;
    12. import java.util.Properties;
    13. /**
    14. * 流式处理
    15. */
    16. public class KafkaStreamQuickStart {
    17. public static void main(String[] args) {
    18. //kafka的配置信息
    19. Properties prop = new Properties();
    20. prop.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,"49.234.48.192:9092");
    21. prop.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
    22. prop.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
    23. prop.put(StreamsConfig.APPLICATION_ID_CONFIG,"streams-quickstart");
    24. //stream 构建器
    25. StreamsBuilder streamsBuilder = new StreamsBuilder();
    26. //流式计算
    27. streamProcessor(streamsBuilder);
    28. //创建kafkaStream对象
    29. KafkaStreams kafkaStreams = new KafkaStreams(streamsBuilder.build(),prop);
    30. //开启流式计算
    31. kafkaStreams.start();
    32. }
    33. /**
    34. * 流式计算
    35. * 消息的内容:hello kafka hello tbug
    36. * @param streamsBuilder
    37. */
    38. private static void streamProcessor(StreamsBuilder streamsBuilder) {
    39. //创建kstream对象,同时指定从哪个topic中接收消息
    40. KStream stream = streamsBuilder.stream("tbug-topic-input");
    41. /**
    42. * 处理消息的value
    43. */
    44. stream.flatMapValues(new ValueMapper>() {
    45. @Override
    46. public Iterable apply(String value) {
    47. return Arrays.asList(value.split(" "));
    48. }
    49. })
    50. //按照value进行聚合处理
    51. .groupBy((key,value)->value)
    52. //时间窗口
    53. .windowedBy(TimeWindows.of(Duration.ofSeconds(10)))
    54. //统计单词的个数
    55. .count()
    56. //转换为kStream
    57. .toStream()
    58. .map((key,value)->{
    59. System.out.println("key:"+key+",value:"+value);
    60. return new KeyValue<>(key.key().toString(),value.toString());
    61. })
    62. //发送消息
    63. .to("tbug-topic-out");
    64. }
    65. }

    3.2:生产者

    1. package com.my.kafka.demo;
    2. import org.apache.kafka.clients.producer.KafkaProducer;
    3. import org.apache.kafka.clients.producer.ProducerConfig;
    4. import org.apache.kafka.clients.producer.ProducerRecord;
    5. import java.util.Properties;
    6. /**
    7. * 生产者
    8. */
    9. public class ProducerDemo {
    10. public static void main(String[] args) {
    11. //1.kafka的配置信息
    12. Properties pro = new Properties();
    13. //Kafka的连接地址
    14. pro.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"49.234.48.192:9092");
    15. //发送失败,失败重连次数
    16. pro.put(ProducerConfig.RETRIES_CONFIG,5);
    17. //消息key的序列化器
    18. pro.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer");
    19. //消息value的序列化器
    20. pro.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringSerializer");
    21. //数据压缩
    22. pro.put(ProducerConfig.COMPRESSION_TYPE_CONFIG,"lz4");
    23. //2.生产者对象
    24. KafkaProducer producer = new KafkaProducer(pro);
    25. for (int i = 0; i < 5; i++) {
    26. //3.封装发送消息
    27. ProducerRecord message = new ProducerRecord<>("tbug-topic-input", "ni gan ma ai yuo");
    28. //4.发送消息
    29. producer.send(message);
    30. }
    31. //5.关闭消息通道(必选)
    32. producer.close();
    33. }
    34. }

    3.3:消费者

    1. package com.my.kafka.demo;
    2. import org.apache.kafka.clients.consumer.ConsumerConfig;
    3. import org.apache.kafka.clients.consumer.ConsumerRecord;
    4. import org.apache.kafka.clients.consumer.ConsumerRecords;
    5. import org.apache.kafka.clients.consumer.KafkaConsumer;
    6. import java.time.Duration;
    7. import java.util.Collections;
    8. import java.util.Properties;
    9. /**
    10. * 消费者
    11. */
    12. public class ConsumerDemo {
    13. public static void main(String[] args) {
    14. //1.添加Kafka配置信息
    15. Properties pro = new Properties();
    16. //Kafka的连接地址
    17. pro.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"49.234.48.192:9092");
    18. //消费者组
    19. pro.put(ConsumerConfig.GROUP_ID_CONFIG,"group2");
    20. //消息key的反序列化器
    21. pro.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
    22. //消息value的反序列化器
    23. pro.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
    24. //手动提交偏移量
    25. pro.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
    26. //2.消费者对象
    27. KafkaConsumer consumer = new KafkaConsumer<>(pro);
    28. //3.订阅主题
    29. consumer.subscribe(Collections.singletonList("tbug-topic-out"));
    30. //4.设置线程一种处于监听状态
    31. //同步提交和异步提交偏移量
    32. try {
    33. while(true) {
    34. //5.获取消息
    35. ConsumerRecords messages = consumer.poll(Duration.ofMillis(1000)); //设置每秒钟拉取一次
    36. for (ConsumerRecord message : messages) {
    37. System.out.print(message.key() + ":");
    38. System.out.println(message.value());
    39. }
    40. }
    41. } catch (Exception e) {
    42. e.printStackTrace();
    43. System.out.println("记录错误的信息:"+e);
    44. } finally {
    45. //同步
    46. consumer.commitSync();
    47. }
    48. }
    49. }

    4.测试结果

    先启动消费者和KafkaStream,然后启动消费者发送消息:

    可以看到成功利用流式计算实现了单词数量的统计,当然你也可以设置多个生产者发送消息。

    要注意的是,假如你是第一次启动程序,那么你先启动KafkaStream是会报错的,这时候你需要先启动生产者创建一个topic,然后再启动KafkaStream,等待一会就能接收到消息。 

  • 相关阅读:
    [附源码]java毕业设计超市收银系统
    nacos注册中心AP核心源码
    java毕业设计糖助手服务交流平台mybatis+源码+调试部署+系统+数据库+lw
    【重铸Java根基】理解Java代理模式机制
    OpenHarmony网络编程及多播相关总结
    NLP:使用 SciKit Learn 的文本矢量化方法
    基于RuoYi-Flowable-Plus的若依ruoyi-nbcio支持自定义业务表单流程(三)
    构建LangChain应用程序的示例代码:48、如何使用非文本生成工具创建多模态代理
    3.3 ss-sp寄存器,栈的push和pop指令
    C语言数据结构算法之求出栈序列个数、判断出栈序列是否合法、求所有合法出栈序列、入栈次数、出栈次数、
  • 原文地址:https://blog.csdn.net/weixin_45750572/article/details/126096199