KafkaStream基本使用
简介:
kafkaStream:提供了对存储在kafka中的数据进行流式处理和分析的功能
特点:
KafkasSream提供了一个非常简单轻量的Library,它可以非常方便的嵌入到java程序中,也可以任何方式打包部署
入门案例:
1、新建工程kafka-demo
引入kafkaStream依赖
-
<dependencies>
-
<dependency>
-
<groupId>org.springframework.boot</groupId>
-
<artifactId>spring-boot-starter-web</artifactId>
-
</dependency>
-
<!-- kafkfa -->
-
<dependency>
-
<groupId>org.springframework.kafka</groupId>
-
<artifactId>spring-kafka</artifactId>
-
<exclusions>
-
<exclusion>
-
<groupId>org.apache.kafka</groupId>
-
<artifactId>kafka-clients</artifactId>
-
</exclusion>
-
</exclusions>
-
</dependency>
-
<dependency>
-
<groupId>org.apache.kafka</groupId>
-
<artifactId>kafka-clients</artifactId>
-
</dependency>
-
<dependency>
-
<groupId>com.alibaba</groupId>
-
<artifactId>fastjson</artifactId>
-
</dependency>
-
-
<!--kafkaStream-->
-
<dependency>
-
<groupId>org.apache.kafka</groupId>
-
<artifactId>kafka-streams</artifactId>
-
<exclusions>
-
<exclusion>
-
<artifactId>connect-json</artifactId>
-
<groupId>org.apache.kafka</groupId>
-
</exclusion>
-
<exclusion>
-
<groupId>org.apache.kafka</groupId>
-
<artifactId>kafka-clients</artifactId>
-
</exclusion>
-
</exclusions>
-
</dependency>
-
</dependencies>
2、新建流式处理类
代码如下
-
package com.heima.kafkademo.sample;
-
-
import org.apache.kafka.common.serialization.Serdes;
-
import org.apache.kafka.streams.KafkaStreams;
-
import org.apache.kafka.streams.KeyValue;
-
import org.apache.kafka.streams.StreamsBuilder;
-
import org.apache.kafka.streams.StreamsConfig;
-
import org.apache.kafka.streams.kstream.KStream;
-
import org.apache.kafka.streams.kstream.TimeWindows;
-
import org.apache.kafka.streams.kstream.ValueMapper;
-
-
import java.time.Duration;
-
import java.util.Arrays;
-
import java.util.Properties;
-
-
/*
-
* 流式处理
-
* */
-
public class KafkaStreamQuickStart {
-
public static void main(String[] args) {
-
/*创建kafka配置中心并配置参数*/
-
Properties prop = new Properties();
-
//连接地址
-
prop.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,"192.168.200.130:9092");
-
//key序列化
-
prop.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
-
//value序列化
-
prop.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
-
//创建id名称
-
prop.put(StreamsConfig.APPLICATION_ID_CONFIG,"streams-quickstart");
-
-
//stream构造器
-
StreamsBuilder streamsBuilder = new StreamsBuilder();
-
-
//流式计算
-
streamProcessor(streamsBuilder);
-
-
//创建KafkaStream对象
-
KafkaStreams kafkaStreams = new KafkaStreams(streamsBuilder.build(),prop);
-
-
//开启流式计算
-
kafkaStreams.start();
-
}
-
-
//流式计算方法
-
private static void streamProcessor(StreamsBuilder streamsBuilder) {
-
//创建kafka对象,同时指定从哪个topic获取消息
-
KStream<String, String> stream = streamsBuilder.stream("itcast-topic-input");
-
//处理消息的value
-
stream.flatMapValues(new ValueMapper<String, Iterable<?>>() {
-
@Override
-
public Iterable<String> apply(String value) {
-
return Arrays.asList(value.split(" "));
-
}
-
}) //按照value进行聚合
-
.groupBy((key,value)->value)
-
//时间窗口,每隔10秒更新一次
-
.windowedBy(TimeWindows.of(Duration.ofSeconds(10)))
-
//统计单词个数
-
.count()
-
//转换为kStream
-
.toStream()
-
.map((key,value)->{
-
System.out.println("key:" key ",vlaue:" value);
-
return new KeyValue<>(key.key().toString(),value.toString());
-
})
-
//发送消息
-
.to("itcast-topic-out");
-
}
-
}
3、启动消费者类和流式处理类监听消息
使用生产者类发送消息
4、测试
成功接收到消息
这篇好文章是转载于:学新通技术网
- 版权申明: 本站部分内容来自互联网,仅供学习及演示用,请勿用于商业和其他非法用途。如果侵犯了您的权益请与我们联系,请提供相关证据及您的身份证明,我们将在收到邮件后48小时内删除。
- 本站站名: 学新通技术网
- 本文地址: /boutique/detail/tanhgfeaga
系列文章
更多
同类精品
更多
-
photoshop保存的图片太大微信发不了怎么办
PHP中文网 06-15 -
Android 11 保存文件到外部存储,并分享文件
Luke 10-12 -
word里面弄一个表格后上面的标题会跑到下面怎么办
PHP中文网 06-20 -
《学习通》视频自动暂停处理方法
HelloWorld317 07-05 -
photoshop扩展功能面板显示灰色怎么办
PHP中文网 06-14 -
微信公众号没有声音提示怎么办
PHP中文网 03-31 -
excel下划线不显示怎么办
PHP中文网 06-23 -
excel打印预览压线压字怎么办
PHP中文网 06-22 -
怎样阻止微信小程序自动打开
PHP中文网 06-13 -
TikTok加速器哪个好免费的TK加速器推荐
TK小达人 10-01