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

java實現(xiàn)請求緩沖合并的示例代碼

 更新時間:2024年04月12日 11:23:56   作者:muguazhi  
我們對外提供了一個rest接口給第三方業(yè)務(wù)進行調(diào)用,但是由于第三方框架限制,導(dǎo)致會發(fā)送大量相似無效請求,這篇文章主要介紹了java實現(xiàn)請求緩沖合并,需要的朋友可以參考下

業(yè)務(wù)背景:

我們對外提供了一個rest接口給第三方業(yè)務(wù)進行調(diào)用,但是由于第三方框架限制,導(dǎo)致會發(fā)送大量相似無效請求,例如:接口入?yún)son包含兩個字段,createBy和receiverList,完整的入?yún)son示例如下:

{
	"createBy": "aa",
	"receiverList": [
		"bb",
		"cc"
	]
}

實際第三方業(yè)務(wù)會進行多次調(diào)用接口,每次傳遞的數(shù)據(jù)可能如下:

{
	"createBy": "aa",
	"receiverList": [
		"bb"
	]
}
或者
{
	"createBy": "aa",
	"receiverList": [
		"cc"
	]
}
或者
{
	"createBy": "bb",
	"receiverList": [
		"cc"
	]
}
或者
{
	"createBy": "aa",
	"receiverList": [
		"bb",
		"cc"
	]
}

所有需要對第三方業(yè)務(wù)傳遞過來的數(shù)據(jù)進行緩沖合并處理,減輕真正的后臺服務(wù)的壓力。

代碼實現(xiàn)

package com.demo;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
import org.springframework.stereotype.Component;
import javax.annotation.PreDestroy;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
/**
 * Description: 請求合并管理類
 */
@Slf4j
@Component
public class RequestMerger {
    // 線程池核心線程數(shù)
    private final int corePoolSize = 200;
    // 任務(wù)執(zhí)行超時時間,單位:毫秒
    private final int timeout = 5 * 60 * 1000;
    // 隊列,隊列長度為Integer.MAX_VALUE
    private final LinkedBlockingQueue<String> requestQueue = new LinkedBlockingQueue<>();
    // 定時器,所有任務(wù)共用線程池
    private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(corePoolSize,
            new CustomizableThreadFactory("schedule-executor-"));
    // 是否關(guān)閉標志
    private final AtomicBoolean isShutdown = new AtomicBoolean(false);
    /**
     * 構(gòu)造函數(shù),用于初始化請求合并器。
     *
     * @param batchSize   每次合并的最大請求數(shù)量。
     * @param delayMillis 合并請求的周期間隔,單位為毫秒。
     */
    public RequestMerger(int batchSize, long delayMillis) {
        // 啟動定時器,定期合并請求,延遲delayMillis后開始,之后每隔delayMillis執(zhí)行一次
        scheduler.scheduleAtFixedRate(() -> {
            if (!isShutdown.get()) {
                List<String> batch = new ArrayList<>(batchSize);
                int drainedCount = requestQueue.drainTo(batch, batchSize);
                log.info("==>scheduler,drainedCount:{},nowQueueCount:{}", drainedCount, requestQueue.size());
                if (!batch.isEmpty()) {
                    // 異步執(zhí)行任務(wù),防止業(yè)務(wù)執(zhí)行時間過長導(dǎo)致業(yè)務(wù)整體延遲過大
                    scheduler.submit(() -> {
                        sendRequestBatch(batch);
                    });
                }
            }
        }, delayMillis, delayMillis, TimeUnit.MILLISECONDS);
    }
    /**
     * 發(fā)送請求批次的方法。
     *
     * @param batch 請求批次。
     */
    private void sendRequestBatch(List<String> batch) {
        Future<?> future = scheduler.submit(() -> {
            try {
                // 在這里實現(xiàn)你的請求發(fā)送邏輯
                // 可以使用HTTP客戶端庫(如Apache HttpClient或OkHttp)來發(fā)送請求
                // ...
                System.out.println("Sending batch of " + batch.size() + " requests");
            } catch (Exception e) {
                // 異常處理邏輯
                System.err.println("Error sending requests: " + e.getMessage());
            }
        });
        // 嘗試獲取任務(wù)結(jié)果,如果超過超時時間則拋出TimeoutException異常,進行取消任務(wù)
        try {
            // 超時時間,單位:毫秒
            future.get(timeout, TimeUnit.MILLISECONDS);
        } catch (TimeoutException | ExecutionException e) {
            // 超時或執(zhí)行異常時取消任務(wù)
            future.cancel(true);
        } catch (Exception e) {
            log.error("==>任務(wù)執(zhí)行異常", e);
            // 任務(wù)執(zhí)行異常
            future.cancel(true);
        }
    }
    /**
     * 在對象銷毀前執(zhí)行的關(guān)閉操作。
     * 該方法從請求隊列中拉取所有未處理的請求,并將它們批量發(fā)送。
     * 無參數(shù)和返回值。
     */
    @PreDestroy
    public void shutdown() {
        isShutdown.set(true);
        List<String> batch = new ArrayList<>();
        // 獲取請求隊列中的剩余所有請求
        int drainedCount = requestQueue.drainTo(batch);
        log.info("==>shutdown,drainedCount:{},nowQueueCount:{}", drainedCount, requestQueue.size());
        // 批量發(fā)送收集到的剩余請求
        sendRequestBatch(batch);
        // 關(guān)閉定時執(zhí)行器
        scheduler.shutdown();
        try {
            if (!scheduler.awaitTermination(60, TimeUnit.SECONDS)) {
                log.error("Scheduler did not terminate gracefully within 60 seconds, force shutting down.");
                scheduler.shutdownNow();
            }
        } catch (InterruptedException e) {
            log.warn("Interrupted during scheduler termination, force shutting down.");
            scheduler.shutdownNow();
            Thread.currentThread().interrupt();
        }
    }
    /**
     * 向請求隊列中添加一個請求。如果服務(wù)未關(guān)閉,則直接添加到請求隊列中;
     * 如果服務(wù)已關(guān)閉,則將該請求作為一批請求發(fā)送。
     *
     * @param request 要添加的請求字符串。
     */
    public void addRequest(String request) throws InterruptedException {
        // 檢查服務(wù)是否已關(guān)閉
        if (!isShutdown.get()) {
            // 未關(guān)閉,直接添加到請求隊列
            requestQueue.put(request);
        } else {
            // 已關(guān)閉,將當前請求作為一批發(fā)送
            List<String> batch = new ArrayList<>();
            batch.add(request);
            sendRequestBatch(batch);
        }
    }
}

參考資料

https://gitee.com/huangjuncong/mumux-framework/tree/master/merge-request/src/main/java/com/mumux/concurrent

注意:此代碼容易導(dǎo)致數(shù)據(jù)丟失。例如:調(diào)用add方法將10個元素放入隊列,但是真正獲取到9個元素。
造成原因:FlushThread#add()中使用offer方法將數(shù)據(jù)放入隊列,如果此時隊列已滿,返回值為false,實際數(shù)據(jù)未進入隊列,需要額外對數(shù)據(jù)進行處理。
修改建議:調(diào)大隊列長度,并且將offer方法改為put方法,保證數(shù)據(jù)最終進入隊列。

到此這篇關(guān)于java實現(xiàn)請求緩沖合并的文章就介紹到這了,更多相關(guān)java請求緩沖合并內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 簡介Java編程中的Object類

    簡介Java編程中的Object類

    這篇文章主要介紹了簡介Java編程中的Object類,是Java入門學(xué)習(xí)中的基礎(chǔ)知識,需要的朋友可以參考下
    2015-09-09
  • 詳解SpringBoot 快速整合Mybatis(去XML化+注解進階)

    詳解SpringBoot 快速整合Mybatis(去XML化+注解進階)

    本篇文章主要介紹了詳解SpringBoot 快速整合Mybatis(去XML化+注解進階),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-11-11
  • Java結(jié)合Kotlin實現(xiàn)寶寶年齡計算

    Java結(jié)合Kotlin實現(xiàn)寶寶年齡計算

    這篇文章主要為大家介紹了Java結(jié)合Kotlin實現(xiàn)寶寶年齡計算示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-06-06
  • SpringBoot項目中HTTP請求體只能讀一次的解決方案

    SpringBoot項目中HTTP請求體只能讀一次的解決方案

    在基于Spring開發(fā)Java項目時,可能需要重復(fù)讀取HTTP請求體中的數(shù)據(jù),例如使用攔截器打印入?yún)⑿畔⒌?但當我們重復(fù)調(diào)用getInputStream()或者getReader()時,通常會遇到SpringBoot HTTP請求只讀一次的問題,本文給出了幾種解決方案,需要的朋友可以參考下
    2024-08-08
  • JVM鉤子函數(shù)的使用場景詳解

    JVM鉤子函數(shù)的使用場景詳解

    當jvm進程退出的時候,或者受到了系統(tǒng)的中斷信號,hook線程就會啟動,一個線程可以注入多個鉤,下面這篇文章主要給大家介紹了關(guān)于JVM鉤子函數(shù)使用的相關(guān)資料,需要的朋友可以參考下
    2021-08-08
  • Java跨環(huán)境部署的完整指南(開發(fā)/測試/生產(chǎn)配置隔離)

    Java跨環(huán)境部署的完整指南(開發(fā)/測試/生產(chǎn)配置隔離)

    在現(xiàn)代軟件開發(fā)中,一次編寫,到處運行的 Java 理念雖然廣為人知,但真正實現(xiàn) 跨環(huán)境無縫部署 卻遠非易事,本文將深入探討如何在 Java 項目中實現(xiàn) 開發(fā)(dev)、測試(test)、生產(chǎn)(prod) 等多環(huán)境的配置隔離與部署策略,需要的朋友可以參考下
    2026-03-03
  • Java?Timer單線程下的定時任務(wù)舉例詳解

    Java?Timer單線程下的定時任務(wù)舉例詳解

    在日常的項目開發(fā)中,多多少少都會涉及到一些定時任務(wù)的需求,下面這篇文章主要介紹了Java?Timer單線程下定時任務(wù)的相關(guān)資料,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2025-10-10
  • 詳解Java Socket通信封裝MIna框架

    詳解Java Socket通信封裝MIna框架

    Mina異步IO使用的Java底層JNI框架,Mina提供服務(wù)端和客戶端,將我們的業(yè)務(wù)解耦開發(fā),真正做到高內(nèi)聚低耦合的思想。
    2021-06-06
  • Java基礎(chǔ)教程之String深度分析

    Java基礎(chǔ)教程之String深度分析

    這篇文章主要給大家介紹了關(guān)于Java基礎(chǔ)教程之String的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • Java?迭代器Iterator完整示例解析

    Java?迭代器Iterator完整示例解析

    迭代器(Iterator)是Java集合框架中的一個核心接口,位于java.util包下,本文給大家講解Java迭代器Iterator完整示例,感興趣的朋友跟隨小編一起看看吧
    2025-09-09

最新評論

彭州市| 城口县| 莱芜市| 道真| 永安市| 柳州市| 平南县| 南靖县| 铜陵市| 东丰县| 云南省| 哈尔滨市| 宁国市| 璧山县| 翁牛特旗| 康马县| 特克斯县| 澄迈县| 呼伦贝尔市| 旅游| 罗定市| 陈巴尔虎旗| 嵩明县| 大理市| 陆良县| 民乐县| 杭锦旗| 磐安县| 大安市| 阿巴嘎旗| 监利县| 孟州市| 宝兴县| 汾西县| 门源| 项城市| 澄迈县| 乌兰察布市| 海盐县| 黔西县| 哈尔滨市|