Kafka轻量级流计算实践


丰富的网络学习资源浩如烟海(原文:网上学习资料一大堆),但若所学知识缺乏系统性(原文:学到的知识不成体系),遇到问题往往停留在表面(原文:遇到问题时只是浅尝辄止),不进行深入探讨(原文:不再深入研究),那么难以实现真正的技术进步(原文:技术提升)。
单兵作战效率很高, 但集体的力量更能推动事情向前发展! 不论你是已经是IT行业的资深从业者, 还是抱着对这一领域浓厚兴趣的新手, 都欢迎加入我们的圈子(技术交流群组, 提供丰富的学习资源, 解决工作中的趣事分享, 向大厂输送优质人才, 以及专业的面试辅导服务), 让我们携手共同进步!
1.2 Kafka Streams 特点
1. 功能强大
- 高扩展性,弹性,容错
2. 轻量级
- 无需专门的集群
- 一个库,而不是框架
3. 完全集成
- 100%的 Kafka 0.10.0 版本兼容
- 易于集成到现有的应用程序
4. 实时性
- 毫秒级延迟
- 并非微批处理
- 窗口允许乱序数据
- 允许迟到数据
1.3 为什么要有 Kafka Streams
目前已有众多流式处理系统以Spark Streaming和Apache Storm为代表的开源流式处理系统因其卓越的应用价值而备受关注。经过多年的演进与发展Apache Storm现也具备广泛的适用性能够实现按记录级别进行高效处理。不仅可轻松整合图计算功能及基于流的SQL查询处理Storm现也支持基于流的SQL查询(SQL on Stream)。对于熟悉其他Spark生态系统开发者而言其学习成本相对较低。同时主流Hadoop发行版如Cloudera和Hortonworks都集成安装了Storm和Spark组件简化了部署流程。
由于 Apache Spark 和 Apache Storm 拥有多方面的优势,为何仍然需要 Kafka Stream 呢?
主要有如下原因。
流式处理框架包括Spark和Storm等技术。然而,Kafka Streams则提供了一个基于Kafka的流式处理类库,其特点在于直接提供了具体的类供开发者调用,从而使得整个应用的运行流程更加灵活可控。具体而言,开发人员需要按照特定的方式编写逻辑代码,并由框架负责调用相关组件,这一过程虽然简化了一定程度,但也带来了两个主要挑战:一是由于缺乏对框架内部机制的深入理解,导致调试难度加大;二是应用范围受到了一定的限制

尽管 Cloudera 和 Hortonworks 简化了 Storm 和 Spark 的部署流程,
然而这些框架本身仍具有较高的配置复杂度。
相比之下,
Kafka Streams 相对于其他流式处理框架而言,
能够轻松集成到应用程序中,
无需对应用进行打包或部署方面的特殊处理。
在流式处理系统的应用中(flow processing systems),大多数系统均支持将Kafka用作其数据源(data source)。具体而言,在Apache Storm中提供了专门的kafka-spout组件(module),而在Apache Spark中则提供了相应的spark-streaming-kafka组件(module)。值得注意的是,在主流流式处理框架中,默认配置通常会采用Kafka作为其数据源(data source)。这表明,在实际应用中已有大量框架默认选择了Kafka作为其数据源(data source)。这种架构下运行的成本显著降低(lower)与传统批处理技术相比
第四,在使用Storm或Spark Streaming时必须向框架自身的进程分配资源。例如,在Storm中需指定supervisor,在Spark on YARN中需指定node manager。尽管针对具体应用实例而言框架自身也会占用一定资源:例如,在Spark Streaming中需为数据整理(shuffle)和存储(storage)阶段预留内存空间。然而作为类库Kafka无需占用系统资源
第五,在Kafka的基础上具备数据持久化的特性后,Kafka Streams实现了基于此的灵活扩展能力。具体而言,在系统架构上它支持通过滚动部署实现扩展性,并通过滚动升级确保系统稳定性的同时支持增量式重新计算。
第六,在基于 Kafka 消费者重平衡机制的情况下,Kafka 流能够在线动态调节平行度。
二、Kafka Streams 数据清洗案例
0)需求
实时处理单词带有”>>>”前缀的内容。例如输入”aaa>>>bbb”,最终处理成 “bbb”
1)需求分析

2)案例实操
- 创建一个工程,并添加 jar 包
- 创建主类
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();
}
}
- 具体业务处理
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;
}
@Override
public void process(byte[] key, byte[] value) {
String input = new String(value);
![img]()
![img]()
![img]()
**既有适合小白学习的零基础资料,也有适合3年以上经验的小伙伴深入学习提升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**
升的进阶课程,涵盖了95%以上大数据知识点,真正体系化!** **由于文件比较多,这里只是将部分目录截图出来,全套包含大厂面经、学习笔记、源码讲义、实战项目、大纲路线、讲解视频,并且后续会持续更新** **[需要这份系统化资料的朋友,可以戳这里获取]()**
