最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spring boot 項目中如何進行kafka Stream app 開發(fā)

 更新時間:2025年09月26日 09:28:54   作者:小落的編程筆記  
文章介紹了Kafka Streams的配置要點與核心方法,強調(diào)應(yīng)用ID唯一性、正確設(shè)置Bootstrap服務(wù)器,區(qū)分KStream與Java Stream的不可變性和多消費特性,本文給大家介紹Springboot項目中如何進行kafka Stream app開發(fā),感興趣的朋友一起看看吧

Kafka Stream

Kafka Stream是Apache Kafka從0.10版本引入的一個新Feature。它是提供了對存儲于Kafka內(nèi)的數(shù)據(jù)進行流式處理和分析的功能。

Kafka Stream的特點

  • Kafka Stream提供了一個非常簡單而輕量的Library,它可以非常方便地嵌入任意Java應(yīng)用中,也可以任意方式打包和部署
  • 除了Kafka外,無任何外部依賴
  • 充分利用Kafka分區(qū)機制實現(xiàn)水平擴展和順序性保證
  • 通過可容錯的state store實現(xiàn)高效的狀態(tài)操作(如windowed join和aggregation)
  • 支持正好一次處理語義
  • 提供記錄級的處理能力,從而實現(xiàn)毫秒級的低延遲
  • 支持基于事件時間的窗口操作,并且可處理晚到的數(shù)據(jù)(late arrival of records)
  • 同時提供底層的處理原語Processor(類似于Storm的spout和bolt),以及高層抽象的DSL(類似于Spark的map/group/reduce)

下面介紹Spring boot 項目中進行kafka Stream app 開發(fā)的詳細過程。

1. 導(dǎo)入依賴

      <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-streams</artifactId>
        <version>3.6.2</version>
      </dependency>

2. 示例代碼(偽代碼)

這段偽代碼只是為了舉例設(shè)置的場景,業(yè)務(wù)場景并不一定合適

@Slf4j
@Component
public class MyKafkaStreamProcessor {
	@Value("${spring.kafka.bootstrap-servers}")
	private String bootstrapServers;
	@PostConstruct
	private void init () {
		String appId = "my-kafka-streams-app";
		myKafkaStreams(appId);
		log.info("? Kafka Streams:{} 初始化完成,開始監(jiān)聽 topic: {}", appId, "source-topic");
	}
	public void myKafkaStreams(String appId) {
		/*=======配置=======*/
		Properties config  = new Properties();
		config.put(StreamsConfig.APPLICATION_ID_CONFIG, appId);
		config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
		StreamsBuilder builder = new StreamsBuilder();
		/*=======構(gòu)建拓撲結(jié)構(gòu)=======*/
		// 數(shù)據(jù)清洗
		KStream<String, Order> stream = builder
			.stream("source-topic", Consumed.with(Serdes.String(), new JsonSerde<>(Order.class)))
			.mapValues(value -> {
				// do something
				return value;
			});
		// 過濾出從app創(chuàng)建的訂單 并進行處理
		stream.filter((k, v) -> Order.getSource.equals("app"))
			.foreach((k, v) -> {
				// do something
			});
		// 發(fā)送到第一個topic
		stream.mapValues(value -> JSON.toJSONString(value), Named.as("to-the-first-target-topic-processor"))
			.to("the-first-target-topic", Produced.with(Serdes.String(), Serdes.String()));
		// 發(fā)送到第二個topic
		stream.filter((k, v) -> {
				// filter something
			})
			.mapValues(value -> {
				// map to another object
			}, Named.as("to-the-second-target-topic-processor"))
			.to("the-second-target-topic", Produced.with(Serdes.String(), Serdes.String()));
		/*=======創(chuàng)建KafkaStreams=======*/
		KafkaStreams streams = new KafkaStreams(builder.build(), config);
		/*=======設(shè)置異常處理器=======*/
		streams.setUncaughtExceptionHandler(new CustomStreamsUncaughtExceptionHandler());
		/*=======啟動streams=======*/
		streams.start();
		/*=======添加jvm hook 確保streams安全退出=======*/
		Runtime.getRuntime().addShutdownHook(new Thread(() -> {
			log.info("關(guān)閉 Kafka Streams:{}...", appId);
			streams.close();
			log.info("Kafka Streams:{}已經(jīng)關(guān)閉!", appId);
		}));
	}
}
@Slf4j
public class CustomStreamsUncaughtExceptionHandler implements StreamsUncaughtExceptionHandler {
	/**
	 * Inspect the exception received in a stream thread and respond with an action.
	 *
	 */
	@Override
	public StreamThreadExceptionResponse handle(Throwable throwable) {
		log.error("Kafka Streams 線程發(fā)生未捕獲異常: {}", ExceptionUtil.stacktraceToString(throwable));
		// 選擇處理策略(以下三選一):
		// 1. 替換線程(繼續(xù)運行)
		return StreamThreadExceptionResponse.REPLACE_THREAD;
		// 2. 關(guān)閉整個 Streams 應(yīng)用
		// return StreamThreadExceptionResponse.SHUTDOWN_CLIENT;
		// 3. 關(guān)閉整個 JVM
		// return StreamThreadExceptionResponse.SHUTDOWN_APPLICATION;
	}
}

3. 一些注意事項和說明

  • kafka stream 的處理部分集中在構(gòu)建的拓撲中,其他部分大同小異
  • 在配置部分 StreamsConfig.APPLICATION_ID_CONFIG 這個參數(shù)是必須的,且不能重復(fù),否則會啟動失敗,StreamsConfig.BOOTSTRAP_SERVERS_CONFIG 是kafka的IP與端口
  • 需要注意的是kafka的 KStream 與 java 中的 stream并不相同,在java中 stream只能被消費一次,但是kstream 可以被消費多次,在上面的demo中可以看到,同一個 kstream 被多次消費,且kstream中的數(shù)據(jù)是不可變的,也就是無論在上一個處理器(processor)對數(shù)據(jù)進行了何種處理,下一個處理器從kstream 中獲取的數(shù)據(jù)依舊是原來的數(shù)據(jù)
  • 在kafka stream app中應(yīng)該對可能會拋出的異常進行處理,而不是全部交給UncaughtExceptionHandler,UncaughtExceptionHandler應(yīng)該只處理哪些無法預(yù)料的異常
  • 如果kafka stream app 捕獲未處理異常之后的處理策略也是替換線程,那么kafka stream app 中如果拋出未捕獲異常,那么這個消費者組就會進入再平衡狀態(tài)(PreparingRebalance),老的消費者從消費者組中剔除,新的消費者加入消費者組,然后再開始消費,注意這種替換線程的處理策略可能導(dǎo)致消息重放,也就是原本的線程消費的offset沒有提交導(dǎo)致新的線程會重復(fù)消費之前已經(jīng)被消費的數(shù)據(jù),如果業(yè)務(wù)會因為消息重放出現(xiàn)異常,建議做冪等

4. kafka Stream 的一些方法說明

  • stream()
    • stream()方法是從源topic獲取數(shù)據(jù)的方法,示例中第一個參數(shù)是字符串,也就是源topic的名稱,
    • 第二個參數(shù)是 Consumed ,用來定義對于消息的key與value反序列化的規(guī)則,
    • 示例中將key序列化為string, value 序列化為order對象,需要注意的是,如果在配置的config中沒有設(shè)置適用于整個stream app的序列化與反序列化規(guī)則,那么后續(xù)的 to()中必須要指定序列化規(guī)則
  • filter()
    • filter()的用法與java stream 中的filter()一致,這里不做說明
  • mapValues()
    • mapValues()的用法,是只對消息的value進行操作,比如將value轉(zhuǎn)換為其他對象。這個方法不會對key進行修改
  • map()
    • 與mapValues()類似,但是可以修改消息的key
  • foreach()
    • 與java stream 的 foreach()類似,也是一個終結(jié)方法
  • to()
    • 終結(jié)方法,用于將數(shù)據(jù)發(fā)送到另外的topic,示例中第一個參數(shù)是目標topic, 第二個參數(shù)是Produced,用于定義key和value的序列化規(guī)則
    • 除了stream()和to()之外,其他方法基本都可以傳一個參數(shù) Named,這個參數(shù)是為每個處理器節(jié)點命名,如果不傳則自動生成,但是在同一個kstream中 命名不能重復(fù)。這個名稱不會影響功能,但是如果有一個名稱可以在后續(xù)調(diào)試和監(jiān)控中提供一點幫助

到此這篇關(guān)于Spring boot 項目中如何進行kafka Stream app 開發(fā)的文章就介紹到這了,更多相關(guān)Spring boot kafka Stream app 開發(fā)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java反射與注解原理解析

    Java反射與注解原理解析

    本文詳細介紹了Java反射的基礎(chǔ)知識,包括概念、示例代碼和進階應(yīng)用,如框架設(shè)計、動態(tài)代理和模板方法,同時,講解了Java注解的概念、基本語法、自定義注解以及如何通過反射獲取注解信息,感興趣的朋友跟隨小編一起看看吧
    2026-02-02
  • spring-cloud入門之eureka-client(服務(wù)注冊)

    spring-cloud入門之eureka-client(服務(wù)注冊)

    本篇文章主要介紹了spring-cloud入門之eureka-client(服務(wù)注冊),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-01-01
  • Servlet會話技術(shù)基礎(chǔ)解析

    Servlet會話技術(shù)基礎(chǔ)解析

    這篇文章主要介紹了Servlet會話技術(shù)基礎(chǔ)解析,具有一定借鑒價值,需要的朋友可以參考下。
    2017-12-12
  • springboot?jpa?實現(xiàn)返回結(jié)果自定義查詢

    springboot?jpa?實現(xiàn)返回結(jié)果自定義查詢

    這篇文章主要介紹了springboot?jpa?實現(xiàn)返回結(jié)果自定義查詢方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • 使用JPA自定義id策略避免主鍵自增

    使用JPA自定義id策略避免主鍵自增

    這篇文章主要介紹了使用JPA自定義id策略避免主鍵自增問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • JAVA使用ElasticSearch查詢in和not in的實現(xiàn)方式

    JAVA使用ElasticSearch查詢in和not in的實現(xiàn)方式

    今天小編就為大家分享一篇關(guān)于JAVA使用Elasticsearch查詢in和not in的實現(xiàn)方式,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2018-12-12
  • springboot多數(shù)據(jù)源配置及切換的示例代碼詳解

    springboot多數(shù)據(jù)源配置及切換的示例代碼詳解

    這篇文章主要介紹了springboot多數(shù)據(jù)源配置及切換,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-09-09
  • Java算法練習題,每天進步一點點(1)

    Java算法練習題,每天進步一點點(1)

    方法下面小編就為大家?guī)硪黄狫ava算法的一道練習題(分享)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧,希望可以幫到你
    2021-07-07
  • springBoot Maven 剔除無用的jar引用問題記錄

    springBoot Maven 剔除無用的jar引用問題記錄

    這篇文章主要介紹了springBoot Maven 剔除無用的jar引用問題記錄,本文給大家介紹的非常詳細,感興趣的朋友跟隨小編一起看看吧
    2024-12-12
  • Java String、StringBuffer與StringBuilder的區(qū)別

    Java String、StringBuffer與StringBuilder的區(qū)別

    本文主要介紹Java String、StringBuffer與StringBuilder的區(qū)別的資料,這里整理了相關(guān)資料及詳細說明其作用和利弊點,有需要的小伙伴可以參考下
    2016-09-09

最新評論

无极县| 西乌珠穆沁旗| 全州县| 木里| 茶陵县| 嘉祥县| 丰镇市| 鹿泉市| 万宁市| 涿州市| 麻江县| 长岛县| 通化市| 开江县| 富锦市| 郯城县| 青铜峡市| 普陀区| 马山县| 思茅市| 新巴尔虎左旗| 策勒县| 靖边县| 宿州市| 会同县| 天峨县| 怀宁县| 保靖县| 大田县| 和龙市| 平昌县| 松桃| 昌都县| 遂宁市| 宣武区| 天柱县| 措勤县| 长兴县| 荃湾区| 广丰县| 西安市|