Advertisement

Kafka(二十四)轻量级流计算 Kafka Streams 实践_轻量级流式计算处理

阅读量:
img
img

丰富的网络学习资源存在;若所学知识缺乏系统性,则在遇到问题时往往停留在表面处理,并且无法深入探究。

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

单兵进度快于团队协作!无论你是资深IT从业者还是Beginner对IT领域充满热情的朋友都欢迎加入我们专业的交流平台(专业交流群+优质的学习素材库+职场吐槽分享区+大厂内推机会+求职面试指导)!在这里我们不仅提供专业的技术讨论还会为大家分享经验传递知识并致力于帮助职场新人快速成长)

流式处理框架属于 Apache Kafka 开源项目的组件。它提供强大的功能和易用性的库工具,并且能够帮助开发人员构建高可用、扩展性强且容错能力突出的Kafka应用程序。

1.2 Kafka Streams 特点

1. 功能强大

  • 高扩展性,弹性,容错

2. 轻量级

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

3. 完全集成

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

4. 实时性

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

目前存在大量流式处理系统,在数据量日益庞大的背景下,Hadoop生态系统的成熟度和稳定性成为企业关注的核心问题之一,其中最为知名的开源流式处理系统之一是Spark Streaming,另一款重要产品则是Apache Storm.经过多年的持续发展,Aphrodite已经形成了广泛的使用场景,并且具备按记录级的处理能力,现支持基于流的SQL查询功能.与此同时,Based on the foundation of Apache Spark构建了强大的数据流平台,Spark Streaming不仅能够方便地集成图计算,SQL查询等功能,而且还提供了高度可扩展和高性能的数据流处理能力.即使是较为熟悉其他Spark组件开发的朋友也能快速上手.此外,the widely-used Hadoop distributions typically include both Storm and Spark modules,从而简化了分布式计算环境下的部署流程

既然 Apache Spark 与 Apache Storm 拥有多项显著优势,那么为什么仍然需要 Kafka Stream 呢?

主要有如下原因。

第一, Spark 和 Storm 均为流处理系统,而 Kafka Streams 则直接提供了基于 Kafka 的流处理功能库...这些架构要求开发者遵循特定的模式进行逻辑实现,由框架负责调用相关组件...由于这些架构的设计较为复杂,在实际操作中难以深入理解其运行机制,导致调试难度较大且限制了灵活性...相比之下, Kafka Streams 则直接提供了完整的流处理功能模块...允许开发者直接调用预设功能模块,从而赋予应用开发者的高度自主权...极大地方便了应用的开发与维护。

尽管 Cloudera 和 Hortonworks 使得 Storm 和 Spark 的部署更加便捷;然而这些框架的部署仍然较为复杂。另一方面 Kafka Streams 作为一个类库 能够非常方便地嵌入到应用程序中 它对应用的打包和部署基本上无需任何特定要求

第三段落指出,在流式处理系统领域中,默认情况下几乎都采用了 Kafka 作为数据源。此外,在具体实现上,Storm 系列提供了专门针对 Kafka 的功能模块(如 kafka-spout),而像 Spark 这样的框架也集成了一套完整的 Kafka 处理机制(spark-streaming-kafka)。实际上,在当前主流的流式处理框架中,默认都会集成 Kafka 作为数据源配置项。换句话说,在大多数流式处理架构中都已经实现了对 Kafka 的集成使用,并且在这样的架构之上构建 Kafka Streams 系统的成本相对较低。

第四,在采用Storm或SparkStreaming时,则必须预先分配好框架内部进程所需的资源配置项。例如Storm的supervisor以及Spark on YARN的应用节点管理器等配置参数都需要特别留意与设置。尽管针对应用实例而言,在运行过程中框架自身也会消耗一定资源。例如SparkStreaming需预先分配shuffle和存储区域的内存空间。然而Kafka作为类库并不占据系统资源。

第五点,在Kafka本身具备数据持久化特性的情况下,在Kafka Streams通过支持流水线部署、逐步升级以及重新计算等特性增强了系统的扩展性和维护性。

第六, 基于 Kafka Consumer Rebalance 机制, Kafka Stream 能够实现对并行处理能力的实时动态调节。

二、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. 具体业务处理
复制代码
    package com.atguigu.kafka.stream;
    import org.apache.kafka.streams.processor.Processor; 
    import org.apache.kafka.streams.processor.ProcessorContext; 
    
    public class LogProcessor implements Processor<byte[], byte[]> { 
    
    	private ProcessorContext context; 
    	@Override 
    	public void init(ProcessorContext context) { 
    		this.context = context; 
    	}
![img]()
![img]()
![img]()
    
    **既有适合小白学习的零基础资料,也有适合3年以上经验的小伙伴深入学习提升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**
    
    升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**

全部评论 (0)

还没有任何评论哟~