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

Spring boot 項(xiàng)目中如何進(jìn)行kafka Stream app 開發(fā)

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

Kafka Stream

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

Kafka Stream的特點(diǎn)

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

下面介紹Spring boot 項(xiàng)目中進(jìn)行kafka Stream app 開發(fā)的詳細(xì)過程。

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)建拓?fù)浣Y(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)建的訂單 并進(jìn)行處理
		stream.filter((k, v) -> Order.getSource.equals("app"))
			.foreach((k, v) -> {
				// do something
			});
		// 發(fā)送到第一個(gè)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ā)送到第二個(gè)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ù)運(yùn)行)
		return StreamThreadExceptionResponse.REPLACE_THREAD;
		// 2. 關(guān)閉整個(gè) Streams 應(yīng)用
		// return StreamThreadExceptionResponse.SHUTDOWN_CLIENT;
		// 3. 關(guān)閉整個(gè) JVM
		// return StreamThreadExceptionResponse.SHUTDOWN_APPLICATION;
	}
}

3. 一些注意事項(xiàng)和說明

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

4. kafka Stream 的一些方法說明

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

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

相關(guān)文章

  • Java反射與注解原理解析

    Java反射與注解原理解析

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

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

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

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

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

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

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

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

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

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

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

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

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

    Java算法練習(xí)題,每天進(jìn)步一點(diǎn)點(diǎn)(1)

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

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

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

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

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

最新評論

景东| 嘉峪关市| 枞阳县| 平安县| 鄯善县| 宽城| 桃源县| 宜兰县| 永城市| 古交市| 佛冈县| 玉山县| 定襄县| 那曲县| 休宁县| 左云县| 安西县| 丽水市| 长沙市| 邹平县| 虹口区| 银川市| 乌苏市| 莲花县| 泸定县| 唐海县| 海盐县| 侯马市| 恭城| 长海县| 呼图壁县| 永吉县| 苍溪县| 修水县| 墨江| 江津市| 临江市| 翼城县| 芦溪县| 岳西县| 内丘县|