Advertisement

Kafka(二十四)实践总结:轻量级流式计算处理流程

阅读量:
img
img
img

本课程既提供了专为小白设计的基础学习材料(零基础资料),也专门针对经验丰富的资深开发者设置了深入研究与能力提升的进阶课程(3年以上经验)。这些课程覆盖了95%以上的核心知识点,并且系统性强。

因为文件数量较多的原因,在这里我们仅对其中一部分目录进行了截图展示。整套资源则包含了大厂面经、学习笔记、源码讲义以及实战项目等关键内容,并且后续也会不断更新和完善。

需要这份系统化资料的朋友,可以戳这里获取

复制代码
    	- [2)案例实操](#2_54)
    + [三、总结](#_159)

一、概述

1.1 Kafka Streams

流处理组件。Apache Kafka开源项目的重要组成部分是一套功能丰富且易于使用的库。它被广泛应用于构建高可靠性的扩展型分布式流处理平台。

1.2 Kafka Streams 特点

1. 功能强大

  • 高扩展性,弹性,容错

2. 轻量级

  • 无需专门的集群
  • 一个库,而不是框架

3. 完全集成

  • 100%的 Kafka 0.10.0 版本兼容
  • 易于集成到现有的应用程序

4. 实时性

  • 毫秒级延迟
  • 并非微批处理
  • 窗口允许乱序数据
  • 允许迟到数据
1.3 为什么要有 Kafka Streams

现有的流式处理系统已非常丰富。主要应用的开源流式处理系统包括Spark Streaming和Apache Storm。经过长期发展,Apache Storm已在多个领域得到广泛应用,并按记录级别进行数据处理的能力现也支持基于流式的SQL运算。而Spark Streaming作为以Apache Spark为基础的技术,在功能上不仅支持与图计算、SQL处理等集成操作,并可帮助开发者快速搭建复杂的数据分析平台。功能全面且强大,在掌握其他相关技术的前提下使用门槛相对较低。另外一项重要的是主流Hadoop发行版如Cloudera和Hortonworks都已经集成了Storm和Spark技术,并大大简化了部署流程。

既然 Apache Spark 与 Apache Storm 拥有多项显著优势,则为什么仍有必要采用 Kafka Stream 作为一种数据流处理技术呢?

主要有如下原因。

第一部分:
Spark 和 Storm 均是流式处理框架,在 Kafka 上构建流式处理架构时可选择 Kafka Streams 这一基于 Kafka 的流式处理类库工具包。在框架设计过程中要求开发者遵循特定的开发逻辑进行系统模块化构建,并由框架负责后续操作流程的具体实现细节等环节;然而由于这种封闭式的开发模式导致开发者难以深入理解其运行机制从而增加了调试难度并限制了其灵活性;相比之下 Kafka Streams 作为一个开放式的流式处理类库工具包直接提供了相应的功能组件供开发者调用使得整个应用架构更加灵活开放并有利于提高系统的扩展性与可维护性等特性;这种设计使开发者的开发效率得到了显著提升同时降低了系统的维护成本。

第二, 尽管 Cloudera 與 Hortonworks 的出现使得 Storm 與 Spark 的部署更加便捷, 但這些筆記工具的部署過程依旧较为繁雜. 前者作為 Kafka 分布式文檔管理工具結合使用時, 只需簡單配置即可完成, 完全無需對程式進行Any packing or deployment optimization.

从流式处理系统的角度来看,几乎所有的系统都默认将 Kafka 设为数据源.此外,像 Storm 这样的平台还提供了 kafka-spout 这样的特定组件,而像 Spark 这样的框架则配备了 spark-streaming-kafka 模块.换句话说,在大多数流式处理系统中,Kafka 已经被集成为其核心组件之一.因此,Kafka Streams 的使用成本相对较低.

在采用Storm或Spark Streaming时,则需预先为该框架自身的进程预留必要的资源支持。具体包括Storm的supervisor以及Spark on YARN中的node manager等组件。即便针对应用实例而言,则仍需认识到框架本身的运行会对系统产生一定影响;例如,在Spark Streaming的情况下,则需预先为数据shuffle和存储区域预留内存空间。相比之下,Kafka作为一个类库,并不需要占用系统资源。

第五部分中指出由于Kafka本身具备数据持久化特性因此Kafka Streams支持滚动部署滚动升级以及再生计算的能力

第六部分基于 Kafka Consumer Rebalance 机制的实施措施... 能够实现实时动态调节处理能力。

二、Kafka Streams 数据清洗案例

0)需求

实时处理单词带有”>>>”前缀的内容。例如输入”aaa>>>bbb”,最终处理成 “bbb”

1)需求分析
在这里插入图片描述
2)案例实操
  1. 创建一个工程,并添加 jar 包
  2. 创建主类
复制代码
    package com.atguigu.kafka.stream; 
    import java.util.Properties; 
    import org.apache.kafka.streams.KafkaStreams; 
    import org.apache.kafka.streams.StreamsConfig; 
    import org.apache.kafka.streams.processor.Processor; 
    import org.apache.kafka.streams.processor.ProcessorSupplier; 
    import org.apache.kafka.streams.processor.TopologyBuilder; 
    
    public class Application { 
    
    	public static void main(String[] args) { 
    
    	// 定义输入的 
    	topic String from = "first"; 
    	// 定义输出的 
    	topic String to = "second"; 
    	// 设置参数 
    	Properties settings = new Properties();
    	settings.put(StreamsConfig.APPLICATION_ID_CONFIG, "logFilter");
    	settings.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "hadoop102:9092");
    	StreamsConfig config = new StreamsConfig(settings); 
    
    	// 构建拓扑 
    	TopologyBuilder builder = new TopologyBuilder(); 
    
    	builder.addSource("SOURCE", from) 
    	.addProcessor("PROCESS", new ProcessorSupplier<byte[], byte[]>() { 
    			@Override 
    			public Processor<byte[], byte[]> get() { 
    		// 具体分析处理
    			return new LogProcessor();
    			 } 
    		 }, "SOURCE")
    	.addSink("SINK", to, "PROCESS"); 
    	// 创建 kafka streams 
    	KafkaStreams streams = new KafkaStreams(builder, config); 
    	streams.start(); 
    	} 
    }
  1. 具体业务处理
复制代码
![img]()
![img]()
![img]()
    
    **既有适合小白学习的零基础资料,也有适合3年以上经验的小伙伴深入学习提升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**
    
    升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**

全部评论 (0)

还没有任何评论哟~