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

Spring WebFlux 流式數(shù)據(jù)拉取與推送的實現(xiàn)

 更新時間:2025年09月05日 09:12:46   作者:on the way 123  
本文介紹了使用Spring WebFlux實現(xiàn)流式數(shù)據(jù)拉取與推送的方案,通過配置長超時、連接池和重試機制優(yōu)化性能,實現(xiàn)了阻塞與非阻塞的結(jié)合,感興趣的可以了解一下

本文介紹了使用Spring WebFlux實現(xiàn)流式數(shù)據(jù)拉取與推送的方案。文章首先展示了流式返回數(shù)據(jù)的格式(類似DeepSeek大模型的推送模式),然后詳細(xì)講解了三個核心實現(xiàn)部分:1)通過Flux.create實現(xiàn)流式響應(yīng)數(shù)據(jù)的橋接轉(zhuǎn)發(fā);2)配置OkHttpClient的HTTP客戶端參數(shù)(特別是readTimeout和callTimeout設(shè)為0以支持流式傳輸);3)核心數(shù)據(jù)獲取方法queryDifficultFaultMessage的實現(xiàn),包括異步請求處理、錯誤處理和取消訂閱機制。該方案實現(xiàn)了后端對原始數(shù)

前言

1,流式返回數(shù)據(jù)類型如下,是不斷的推送數(shù)據(jù),類似于主流DeepSeek大模型模式,推送數(shù)據(jù)一點點推送,直至推送結(jié)束或者主動點擊停止

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":""}

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":""}

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":"故障"}

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":"根"}

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":"因"}
....
data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"agent_message","answer":":\n\n"}

data:{"conversationId":"eef6023a-6ace-420f-ad46-324aa53decf8","event":"message_end","answer":"一"}```

2,流式接口又稱基于響應(yīng)式編程(Reactive Programming)的服務(wù)器端到服務(wù)器端(Service-to-Service)的流式數(shù)據(jù)拉取與推送的實現(xiàn)。

3,它的核心作用是:從一個流式API(SSE)消費數(shù)據(jù),并立即將其作為流式響應(yīng)轉(zhuǎn)發(fā)給客戶端

4,功能要求,后段不做數(shù)據(jù)處理,大模型接口返回什么數(shù)據(jù),直接回傳給前端,分阻塞返回,【當(dāng)時做的效果是:大模型返回數(shù)據(jù),后端自動拼接結(jié)束后返回給前端,這種效果很不友好,導(dǎo)致請求的時間比較長】

功能實現(xiàn)

1.流式響應(yīng)數(shù)據(jù)

public Flux<DifficultFaultMessageVo> streamingQueryDifficultFault(DifficultFaultMessageDto paramDto) {
    // 這個 Flux 代表了從下游服務(wù)獲取的原始數(shù)據(jù)流。
    Flux<DifficultFaultMessageVo> faults = queryDifficultFaultMessage(paramDto);
    return Flux.create(sink -> {
        // faults.subscribeOn(Schedulers.boundedElastic()): 這行代碼至關(guān)重要。它告訴上游的 faults 流在 boundedElastic 調(diào)度器上執(zhí)行其訂閱操作(即執(zhí)行網(wǎng)絡(luò)請求和處理響應(yīng))。
        faults.subscribeOn(Schedulers.boundedElastic())
              // 這里創(chuàng)建了一個新的 Flux。create 方法允許我們手動控制如何向這個流中發(fā)射數(shù)據(jù)。我們傳入一個 Consumer,它接收一個 FluxSink 對象(這里的參數(shù)名為 sink)作為參數(shù)。Sink(匯)就是數(shù)據(jù)流的出口,我們可以通過它發(fā)射數(shù)據(jù) (next)、錯誤 (error) 或完成信號 (complete)。
                .subscribe(sink::next, sink::error, sink::complete);
    });
}

代碼解析:

  • subscribe(sink::next, sink::error, sink::complete): 這里訂閱了從 queryDifficultFaults 返回的原始流。
  • 當(dāng)原始流 (faults) 產(chǎn)生一個數(shù)據(jù) (DifficultFaultMessageVo 對象) 時,就通過 sink::next 將它轉(zhuǎn)發(fā)給我們新創(chuàng)建的流的 Sink。
  • 當(dāng)原始流發(fā)生錯誤時,通過 sink::error 將錯誤轉(zhuǎn)發(fā)給新流的 Sink。
  • 當(dāng)原始流結(jié)束時,通過 sink::complete 結(jié)束新流的 Sink。

總結(jié):這個方法的作用可以理解為 “流的橋接” 。它將在一個彈性線程上執(zhí)行的、可能阻塞的原始數(shù)據(jù)流,橋接成一個適合在WebFlux等響應(yīng)式Web框架中返回的響應(yīng)式流。外部調(diào)用者(如Controller)只需返回這個方法的返回值,框架就會自動處理流的訂閱和HTTP響應(yīng)體的流式寫入。

2.HTTP客戶端配置 okhttpclient

private final OkHttpClient httpClient = new OkHttpClient.Builder()
        .connectTimeout(120, TimeUnit.SECONDS) // 連接超時2分鐘
        .readTimeout(0, TimeUnit.SECONDS)      // 讀取超時:0(無限等待,對于流式響應(yīng)關(guān)鍵?。?
        .writeTimeout(120, TimeUnit.SECONDS)   // 寫入超時2分鐘
        .callTimeout(0, TimeUnit.SECONDS)      // 整個調(diào)用超時:0(無限等待)
        .retryOnConnectionFailure(true)        // 自動重試連接失敗
        .connectionPool(new ConnectionPool(20, 5, TimeUnit.MINUTES)) // 連接池(20個空閑連接,存活5分鐘)
        .build();

代碼解讀:

這個配置是流式HTTP客戶端的靈魂,每一項都針對長連接和流式傳輸進(jìn)行了優(yōu)化:

  • readTimeout(0): 這是最關(guān)鍵的配置。普通的HTTP請求需要設(shè)置讀取超時,但流式響應(yīng)是一個長時間存在的連接,數(shù)據(jù)會分塊持續(xù)發(fā)送。設(shè)置為 0 表示永不超時,客戶端會一直等待服務(wù)器發(fā)送更多數(shù)據(jù),直到連接被服務(wù)器或自己主動關(guān)閉。
  • callTimeout(0): 同理,整個調(diào)用的總時間也不應(yīng)設(shè)限。
  • connectTimeout 和 writeTimeout: 這兩個仍然需要設(shè)置一個合理的值,分別控制建立TCP連接的時間和發(fā)送請求體的時間,這些操作不應(yīng)該無限等待。
  • connectionPool: 使用連接池可以復(fù)用TCP連接,避免為每個請求都進(jìn)行三次握手,極大提升性能。這里配置了最多保持20個空閑連接,每個空閑連接最多存活5分鐘。
  • retryOnConnectionFailure: 網(wǎng)絡(luò)抖動時自動重試,提高魯棒性。

3.核心數(shù)據(jù)獲取方法

private Flux<DifficultFaultMessageVo> queryDifficultFaultMessage(DifficultFaultMessageDto paramDto) {
    return Flux.create(emitter -> { // 這個emitter是內(nèi)部Flux的Sink
        Request request = buildRequest(paramDto);
        Call call = httpClient.newCall(request);
        
        call.enqueue(new Callback() { // 異步執(zhí)行HTTP請求
            @Override
            public void onFailure(...) {
                if (!emitter.isCancelled()) {
                    emitter.error(...); // 網(wǎng)絡(luò)失敗,向上游發(fā)射錯誤
                }
            }

            @Override
            public void onResponse(...) throws IOException {
                if (!response.isSuccessful()) {
                    if (!emitter.isCancelled()) {
                        emitter.error(...); // HTTP狀態(tài)碼非2xx,向上游發(fā)射錯誤
                    }
                    return;
                }
                // 成功響應(yīng),開始處理流式響應(yīng)體
                processResponseStream(response, emitter);
            }
        });

        // 重要:注冊取消回調(diào)
        emitter.onCancel(() -> {
            if (!call.isCanceled()) {
                call.cancel(); // 如果下游取消訂閱(如客戶端斷開),則取消OkHttp請求
            }
        });
    });
}

構(gòu)建請求

    private Request buildRequest(DifficultFaultMessageDto paramDto) {
        // 2.組裝請求頭信息
        MediaType parse = MediaType.parse("application/json;charset=UTF-8");
        // 3.組裝請求體信息
        JSONObject requestBody = new JSONObject();
        requestBody.put("x x x", paramDto.getxxxType());
        requestBody.put("apiKey", paramDto.getApiKey());

		// 省略業(yè)務(wù)代碼
        ... 
        log.info("API請求體: {}", requestBody);
        return new Request.Builder().url("http://xxx.x.xx.xx:8000/servicexxx/r/postApi").post(body)
                .addHeader("X-APP-ID", "xx09xx3xx0d7").addHeader("X-APP-KEY", "9jfksjfjkxxkkssdc")
                .addHeader("Content-Type", "application/json").build();
    }

代碼解讀:

  • 目的:創(chuàng)建一個 Flux,用于封裝對下游服務(wù)的異步HTTP調(diào)用和流式響應(yīng)處理。

執(zhí)行流程:

  • 構(gòu)建請求 (buildRequest) 和調(diào)用對象 (Call)。
  • 異步執(zhí)行 (call.enqueue) HTTP請求。

在回調(diào)中:

  • 失敗 (onFailure): 檢查內(nèi)部的 FluxSink (emitter) 是否還未被取消(即下游是否還在關(guān)心結(jié)果),如果是,則發(fā)射一個錯誤信號。
  • 成功 (onResponse): 檢查HTTP狀態(tài)碼,如果不成功則發(fā)射錯誤;如果成功,則調(diào)用 processResponseStream 開始處理響應(yīng)體流。
  • emitter.onCancel(…): 這是響應(yīng)式編程中資源清理的關(guān)鍵。它注冊了一個回調(diào),當(dāng)這個 Flux 的下游訂閱者取消訂閱時(例如,前端用戶關(guān)閉了瀏覽器標(biāo)簽頁),這個回調(diào)會被觸發(fā)。回調(diào)里會取消底層的OkHttp Call 對象,從而立即關(guān)閉網(wǎng)絡(luò)連接,避免資源泄漏。這是一種“背壓”(Backpressure)傳播,體現(xiàn)了響應(yīng)式的優(yōu)點。

4.流式響應(yīng)體處理 【最核心部分】

private void processResponseStream(Response response, FluxSink<DifficultFaultsVo> emitter) {
    try (ResponseBody responseBody = response.body()) { // 使用try-with-resources確保資源關(guān)閉
        ...
        BufferedSource source = responseBody.source(); // 獲取緩沖數(shù)據(jù)源
        AtomicBoolean isComplete = new AtomicBoolean(false); // 標(biāo)志位,是否收到結(jié)束事件

        try {
            while (!emitter.isCancelled()) { // 循環(huán),只要下游沒有取消就繼續(xù)讀
                String line = source.readUtf8Line(); // 讀取一行UTF-8文本
                if (line == null) break; // 讀到null表示流自然結(jié)束(服務(wù)器關(guān)閉連接)
                if (StringUtils.isBlank(line)) continue; // 忽略空行

                // SSE協(xié)議格式:每段數(shù)據(jù)以"data: "開頭
                if (line.startsWith("data:")) {
                    String jsonData = line.substring(6); // 截取"data: "后面的JSON字符串
                    // 過濾:只處理包含"answer"或"message_end"的數(shù)據(jù)行
                    if (!jsonData.contains("answer") && !jsonData.contains("message_end")) {
                        continue;
                    }

                    DifficultFaultMessageVo vo = processModelResponse(jsonData); // 解析JSON為值對象
                    if (vo != null) {
                        emitter.next(vo); // 解析成功,立即發(fā)射給下游(最終到前端)
                        // 如果遇到結(jié)束事件,標(biāo)記完成并結(jié)束循環(huán)
                        if ("message_end".equals(vo.getEvent())) {
                            isComplete.set(true);
                            emitter.complete();
                            break;
                        }
                    }
                }
            }
        } catch (Exception e) {
            // 處理讀取和解析過程中的異常
            if (!emitter.isCancelled()){
                emitter.error(e);
            }
        } finally {
            // 確保流最終完成
            if (!isComplete.get() && !emitter.isCancelled()){
                emitter.complete();
            }
        }
        ...
    } catch (Exception e) {
        ...
    } finally {
        response.close(); // 最終確保HTTP響應(yīng)被關(guān)閉
    }
}


private DifficultFaultMessageVo processModelResponse(String jsonData) {
        try {
            // 1. 解析JSON
            DifficultFaultMessageVo rawResponse = JSONUtil.toBean(
                    jsonData,
                    DifficultFaultMessageVo.class,
                    false
            );

            // 2. 過濾非消息事件
            String event = rawResponse.getEvent();

            // 3. 創(chuàng)建前端響應(yīng)對象
            DifficultFaultMessageVo result = new DifficultFaultMessageVo();
            result.setConversationId(rawResponse.getConversationId());

            // 4. 處理結(jié)束事件
            if ("message_end".equals(event)) {
                result.setEvent("message_end");
                result.setAnswer("");
                return result;
            }
            result.setAnswer(rawResponse.getAnswer());
            result.setEvent(rawResponse.getEvent());

            return result;

        } catch (Exception e) {
            log.error("JSON解析錯誤: {} - {}", jsonData, e.getMessage());
            return null;
        }
    }

代碼解讀:

  • 核心任務(wù):從 ResponseBody 的流中逐行讀取并解析SSE格式。
  • SSE格式簡介:通常為 data: {json}\n\n。代碼只關(guān)心以 data: 開頭的行。

關(guān)鍵點:

  • 逐行讀取: source.readUtf8Line() 是阻塞方法,這就是為什么必須在 boundedElastic 線程上執(zhí)行的原因。
  • 過濾: 并非所有 data: 行都需要處理,這里通過檢查內(nèi)容來過濾。
  • JSON解析: processModelResponse(jsonData) 方法(代碼未給出)負(fù)責(zé)將JSON字符串解析為 DifficultFaultsVo 對象。
  • 實時發(fā)射: 一旦解析成功,立即通過 emitter.next(vo) 將數(shù)據(jù)推送給下游,實現(xiàn)了數(shù)據(jù)塊的零延遲轉(zhuǎn)發(fā)。

結(jié)束條件:

  • 顯式結(jié)束: 收到 “message_end” 事件,調(diào)用 emitter.complete()。
  • 隱式結(jié)束: 服務(wù)器關(guān)閉連接(readUtf8Line() 返回 null),跳出循環(huán)。
  • 異常結(jié)束: 捕獲到任何異常,調(diào)用 emitter.error(e)。
  • 健壯性保證: 大量的 if (!emitter.isCancelled()) 檢查確保了在下游已經(jīng)不感興趣的情況下,不會進(jìn)行無效的操作(發(fā)射數(shù)據(jù)、錯誤或完成信號)。finally 塊確保了在任何情況下流最終都會被關(guān)閉,防止資源泄漏。

總結(jié):

這段代碼實現(xiàn)了一個高效、健壯的雙重流式處理管道:

  • 下游流 (OkHttp -> 本服務(wù)):使用配置了長超時的OkHttp客戶端,異步調(diào)用外部流式API,并在獨立的彈性線程上阻塞地、逐行讀取SSE響應(yīng)。
  • 上游流 (本服務(wù) -> 客戶端):通過Project Reactor的 Flux 和 Sink,將下游獲取到的數(shù)據(jù)塊立即、實時地轉(zhuǎn)發(fā)給最終的客戶端(如Web瀏覽器)。

關(guān)鍵特性:

  • 非阻塞IO: 通過將阻塞操作卸載到專用線程池,保護(hù)了Web容器的核心線程。
  • 背壓傳播: 下游的取消訂閱會向上傳播,最終取消OkHttp請求,及時釋放資源。
  • 全面錯誤處理: 對網(wǎng)絡(luò)錯誤、HTTP錯誤、解析錯誤、連接意外關(guān)閉等都有處理。
  • 資源安全: 廣泛使用 try-with-resources 和 finally 塊確保網(wǎng)絡(luò)連接和響應(yīng)體被正確關(guān)閉。

這是一種在Spring WebFlux等響應(yīng)式框架中集成傳統(tǒng)阻塞式HTTP客戶端以消費流式服務(wù)的標(biāo)準(zhǔn)且優(yōu)雅的模式。

到此這篇關(guān)于Spring WebFlux 流式數(shù)據(jù)拉取與推送的實現(xiàn)的文章就介紹到這了,更多相關(guān)Spring WebFlux 流式拉取與推送內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • k8s部署的java服務(wù)添加idea調(diào)試參數(shù)的方法

    k8s部署的java服務(wù)添加idea調(diào)試參數(shù)的方法

    文章介紹了如何在K8S容器中的Java服務(wù)上進(jìn)行遠(yuǎn)程調(diào)試,包括配置Deployment、Service以及本地IDEA的調(diào)試設(shè)置,感興趣的朋友跟隨小編一起看看吧
    2025-02-02
  • 使用Gradle打依賴包失敗的問題及解決

    使用Gradle打依賴包失敗的問題及解決

    這篇文章主要介紹了使用Gradle打依賴包失敗的問題及解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • mybatis?resultMap之collection聚集兩種實現(xiàn)方式

    mybatis?resultMap之collection聚集兩種實現(xiàn)方式

    本文主要介紹了mybatis?resultMap之collection聚集兩種實現(xiàn)方式,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2024-09-09
  • idea的終端(Terminal)cmd的命令換成linux的命令詳解

    idea的終端(Terminal)cmd的命令換成linux的命令詳解

    本文介紹IDEA配置Git的步驟:安裝Git、修改終端設(shè)置并重啟IDEA,強調(diào)順序,作為個人經(jīng)驗分享,希望提供參考并支持腳本之家
    2025-07-07
  • Vue?+?springboot實現(xiàn)拼圖人機驗證功能

    Vue?+?springboot實現(xiàn)拼圖人機驗證功能

    本文介紹了如何使用Vue和Spring?Boot實現(xiàn)拼圖人機驗證功能,本文通過實例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友參考下吧
    2026-01-01
  • 2025版Maven安裝與配置終極指南(最全)

    2025版Maven安裝與配置終極指南(最全)

    Maven是Apache軟件基金會的開源工具,主要用于Java項目的構(gòu)建、依賴管理和報告生成,本文將為大家詳細(xì)介紹一下Maven的安裝和環(huán)境,有需要的可以了解下
    2025-10-10
  • java循環(huán)結(jié)構(gòu)、數(shù)組的使用小結(jié)

    java循環(huán)結(jié)構(gòu)、數(shù)組的使用小結(jié)

    這篇文章主要介紹了java循環(huán)結(jié)構(gòu)、數(shù)組的使用小結(jié),本文通過實例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-09-09
  • Java跨庫事務(wù)如何實現(xiàn)保證

    Java跨庫事務(wù)如何實現(xiàn)保證

    這段文章詳細(xì)介紹了Java中實現(xiàn)跨庫事務(wù)管理的方法,包括使用XA協(xié)議、Spring的JTA事務(wù)管理器和第三方解決方案如Seata,文章通過具體步驟和配置示例,解釋了如何在Java應(yīng)用中實現(xiàn)高效、可靠的分布式事務(wù)管理,適用于不同場景和需求
    2026-05-05
  • MyBatis與Hibernate的比較

    MyBatis與Hibernate的比較

    Hibernate 與Mybatis都是流行的持久層開發(fā)框架,但Hibernate開發(fā)社區(qū)相對多熱鬧些,支持的工具也多,更新也快,當(dāng)前最高版本4.1.8。而Mybatis相對平靜,工具較少,當(dāng)前最高版本3.2
    2016-01-01
  • Java線程活鎖的實現(xiàn)與死鎖等的區(qū)別

    Java線程活鎖的實現(xiàn)與死鎖等的區(qū)別

    活鎖是一種遞歸情況,其中兩個或更多線程將繼續(xù)重復(fù)特定的代碼邏輯,本文主要介紹了Java線程活鎖的實現(xiàn)與死鎖等的區(qū)別,具有一定的參考價值,感興趣的可以了解一下
    2024-04-04

最新評論

南通市| 乐东| 东乡族自治县| 阿图什市| 福海县| 章丘市| 灵寿县| 西昌市| 长子县| 青岛市| 从化市| 江华| 淄博市| 陵水| 朝阳县| 宝应县| 康乐县| 五寨县| 普陀区| 鲁山县| 新宁县| 邛崃市| 平江县| 济宁市| 桐乡市| 德惠市| 西乌珠穆沁旗| 临沧市| 思茅市| 黎平县| 屏山县| 咸宁市| 桦南县| 锡林郭勒盟| 达拉特旗| 海兴县| 泾川县| 图木舒克市| 桃园县| 娄底市| 襄垣县|