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

springboot基于注解實(shí)現(xiàn)去重表消息防止重復(fù)消費(fèi)

 更新時(shí)間:2025年05月16日 10:20:22   作者:sjsjsbbsbsn  
本文主要介紹了springboot基于注解實(shí)現(xiàn)去重表消息防止重復(fù)消費(fèi),通過記錄消息ID、使用分布式鎖和設(shè)置過期時(shí)間,可以確保消息只會被處理一次,具有一定的參考價(jià)值,感興趣的可以了解一下

1. 背景/問題

在分布式系統(tǒng)中,消息隊(duì)列(如RocketMQ、Kafka)的 消息重復(fù)消費(fèi) 是常見問題,主要原因包括:

  • 網(wǎng)絡(luò)抖動:生產(chǎn)者或消費(fèi)者因網(wǎng)絡(luò)不穩(wěn)定觸發(fā)消息重發(fā)。
  • 消費(fèi)者超時(shí):消費(fèi)者處理時(shí)間過長,消息隊(duì)列誤判為失敗并重新投遞。
  • 集群故障轉(zhuǎn)移:消費(fèi)者宕機(jī)后,未完成的消息會被其他節(jié)點(diǎn)重新拉取。

重復(fù)消費(fèi)帶來的問題

  • 業(yè)務(wù)邏輯多次執(zhí)行(如重復(fù)扣款、重復(fù)生成訂單)。
  • 數(shù)據(jù)一致性被破壞(如庫存超賣、積分累加錯(cuò)誤)。
  • 系統(tǒng)資源浪費(fèi),影響性能和穩(wěn)定性。

為了避免這種情況發(fā)生,需要在客戶端實(shí)現(xiàn)一些機(jī)制來確保消息不會被重復(fù)消費(fèi),例如記錄消費(fèi)者已經(jīng)處理的消息 ID、使用分布式鎖來控制消費(fèi)進(jìn)程的唯一性等。這些機(jī)制能夠保證消息被成功處理,同時(shí)也能夠提高系統(tǒng)的可靠性和穩(wěn)定性。

2. 什么是冪等性

冪等性 是指對同一操作的多次執(zhí)行所產(chǎn)生的影響與一次執(zhí)行的影響相同。

  • 消息消費(fèi)場景:無論消息被消費(fèi)多少次,最終結(jié)果應(yīng)與消費(fèi)一次一致。
  • 實(shí)現(xiàn)目標(biāo):通過冪等設(shè)計(jì),確保業(yè)務(wù)邏輯的重復(fù)執(zhí)行不會產(chǎn)生副作用。

3. 冪等設(shè)計(jì)

核心思路

  • 冪等標(biāo)識:為每條消息生成唯一標(biāo)識(如業(yè)務(wù)ID + 消息ID),記錄其處理狀態(tài)。
  • 狀態(tài)管理:通過數(shù)據(jù)庫或Redis維護(hù)冪等標(biāo)識的狀態(tài)(如“消費(fèi)中”“已消費(fèi)”)。
  • 過期時(shí)間:防止因系統(tǒng)崩潰導(dǎo)致狀態(tài)長期滯留,需設(shè)置合理的超時(shí)時(shí)間(如10分鐘)。
[消費(fèi)者接收消息]  
        │  
        ▼  
[解析消息,生成唯一冪等標(biāo)識]  
        │  
        ▼  
[查詢冪等標(biāo)識狀態(tài)]  
        │  
┌───────┴───────┐  
│ 存在且已消費(fèi)  │           [返回成功,丟棄消息]  
└───────┬───────┘  
        │  
┌───────┴───────┐  
│ 存在且消費(fèi)中  │           [延遲消費(fèi),等待重試]  
└───────┬───────┘  
        │  
┌───────┴───────┐  
│   不存在      │  
└───────┬───────┘  
        │  
[設(shè)置冪等標(biāo)識為“消費(fèi)中”,并設(shè)置過期時(shí)間]  
        │  
        ▼  
[執(zhí)行業(yè)務(wù)邏輯]  
        │  
        ▼  
[業(yè)務(wù)執(zhí)行成功?]  
        │  
┌───────┴───────┐  
│     是        │           [更新標(biāo)識為“已消費(fèi)”]  
│               │           [刪除或保留標(biāo)識]  
└───────┬───────┘  
        │  
┌───────┴───────┐  
│     否        │           [刪除標(biāo)識,允許重試]  
└───────┬───────┘  
        │  
        ▼  
[流程結(jié)束]  

4.抽象通用冪等組件

消息防重復(fù)消費(fèi)冪等組件是通用的通常會提取出來也可供其他模塊/服務(wù) 使用

4.1自定義冪等注解

提供了一種通用的冪等注解,并通過 SpEL 的形式生成去重表全局唯一 Key

@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface NoMQDuplicateConsume {

    /**
     * 設(shè)置防重令牌 Key 前綴
     */
    String keyPrefix() default "";

    /**
     * 通過 SpEL 表達(dá)式生成的唯一 Key
     */
    String key();

    /**
     * 設(shè)置防重令牌 Key 過期時(shí)間,單位秒,默認(rèn) 1 小時(shí)
     */
    long keyTimeout() default 3600L;
}

4.2. 定義冪等枚舉

冪等需要設(shè)置兩個(gè)狀態(tài),消費(fèi)中和已消費(fèi),創(chuàng)建對應(yīng)的枚舉

@RequiredArgsConstructor
public enum IdempotentMQConsumeStatusEnum {

    /**
     * 消費(fèi)中
     */
    CONSUMING("0"),

    /**
     * 已消費(fèi)
     */
    CONSUMED("1");

    @Getter
    private final String code;

    /**
     * 如果消費(fèi)狀態(tài)等于消費(fèi)中,返回失敗
     *
     * @param consumeStatus 消費(fèi)狀態(tài)
     * @return 是否消費(fèi)失敗
     */
    public static boolean isError(String consumeStatus) {
        return Objects.equals(CONSUMING.code, consumeStatus);
    }
}

4.3.通過 AOP 的方式進(jìn)行增強(qiáng)注解

如果說方法上加了注解,會被這段 AOP 代碼以環(huán)繞增強(qiáng)方式執(zhí)行

@Slf4j
@Aspect
@RequiredArgsConstructor
public final class NoMQDuplicateConsumeAspect {

    private final StringRedisTemplate stringRedisTemplate;

    private static final String LUA_SCRIPT = """
            local key = KEYS[1]
            local value = ARGV[1]
            local expire_time_ms = ARGV[2]
            return redis.call('SET', key, value, 'NX', 'GET', 'PX', expire_time_ms)
            """;

    /**
     * 增強(qiáng)方法標(biāo)記 {@link NoMQDuplicateConsume} 注解邏輯
     */
    @Around("@annotation(com.nageoffer.onecoupon.framework.idempotent.NoMQDuplicateConsume)")
    public Object noMQRepeatConsume(ProceedingJoinPoint joinPoint) throws Throwable {
        NoMQDuplicateConsume noMQDuplicateConsume = getNoMQDuplicateConsumeAnnotation(joinPoint);
        String uniqueKey = noMQDuplicateConsume.keyPrefix() + SpELUtil.parseKey(noMQDuplicateConsume.key(), ((MethodSignature) joinPoint.getSignature()).getMethod(), joinPoint.getArgs());

        String absentAndGet = stringRedisTemplate.execute(
                RedisScript.of(LUA_SCRIPT, String.class),
                List.of(uniqueKey),
                IdempotentMQConsumeStatusEnum.CONSUMING.getCode(),
                String.valueOf(TimeUnit.SECONDS.toMillis(noMQDuplicateConsume.keyTimeout()))
        );

        // 如果不為空證明已經(jīng)有
        if (Objects.nonNull(absentAndGet)) {
            boolean errorFlag = IdempotentMQConsumeStatusEnum.isError(absentAndGet);
            log.warn("[{}] MQ repeated consumption, {}.", uniqueKey, errorFlag ? "Wait for the client to delay consumption" : "Status is completed");
            if (errorFlag) {
                throw new ServiceException(String.format("消息消費(fèi)者冪等異常,冪等標(biāo)識:%s", uniqueKey));
            }
            return null;
        }

        Object result;
        try {
            // 執(zhí)行標(biāo)記了消息隊(duì)列防重復(fù)消費(fèi)注解的方法原邏輯
            result = joinPoint.proceed();

            // 設(shè)置防重令牌 Key 過期時(shí)間,單位秒
            stringRedisTemplate.opsForValue().set(uniqueKey, IdempotentMQConsumeStatusEnum.CONSUMED.getCode(), noMQDuplicateConsume.keyTimeout(), TimeUnit.SECONDS);
        } catch (Throwable ex) {
            // 刪除冪等 Key,讓消息隊(duì)列消費(fèi)者重試邏輯進(jìn)行重新消費(fèi)
            stringRedisTemplate.delete(uniqueKey);
            throw ex;
        }
        return result;
    }

    /**
     * @return 返回自定義防重復(fù)消費(fèi)注解
     */
    public static NoMQDuplicateConsume getNoMQDuplicateConsumeAnnotation(ProceedingJoinPoint joinPoint) throws NoSuchMethodException {
        MethodSignature methodSignature = (MethodSignature) joinPoint.getSignature();
        Method targetMethod = joinPoint.getTarget().getClass().getDeclaredMethod(methodSignature.getName(), methodSignature.getMethod().getParameterTypes());
        return targetMethod.getAnnotation(NoMQDuplicateConsume.class);
    }

lua腳本解釋

local key = KEYS[1] # 第一個(gè) Key,即冪等唯一標(biāo)識 uniqueKey
local value = ARGV[1] # 第一個(gè)參數(shù),即初始化冪等消費(fèi)狀態(tài),為消費(fèi)中
local expire_time_ms = ARGV[2] # 第二個(gè)參數(shù),即冪等 Key 過期時(shí)間

return redis.call('SET', key, value, 'NX', 'GET', 'PX', expire_time_ms)

該腳本的主要作用是:在 Redis 中嘗試以 NX 方式設(shè)置一個(gè)鍵,即如果鍵不存在,則設(shè)置新值,并返回設(shè)置之前的舊值,同時(shí)為該鍵設(shè)置過期時(shí)間(以毫秒為單位)。

獲取到 Redis 里面的 Key 值后,可能會有三個(gè)流程執(zhí)行:

  • absentAndGet 為空:代表消息是第一次到達(dá),執(zhí)行完 LUA 腳本后,會在 Redis 設(shè)置 Key 的 Value 值為 0,消費(fèi)中狀態(tài)。
  • absentAndGet 為 0:代表已經(jīng)有相同消息到達(dá)并且還沒有處理完,會通過拋異常的形式讓 RocketMQ 重試。
  • absentAndGet 為 1:代表已經(jīng)有相同消息消費(fèi)完成,返回空表示不執(zhí)行任何處理。

4.4.注冊為 Spring Bean

另外可以看看另一篇基于分布式鎖注解防重復(fù)提交

https://blog.csdn.net/sjsjsbbsbsn/article/details/145131305?spm=1001.2014.3001.5501

public class IdempotentConfiguration {
    /**
     * 防止消息隊(duì)列消費(fèi)者重復(fù)消費(fèi)消息切面控制器
     */
    @Bean
    public NoMQDuplicateConsumeAspect noMQDuplicateConsumeAspect(StringRedisTemplate stringRedisTemplate) {
        return new NoMQDuplicateConsumeAspect(stringRedisTemplate);
    }
}

4.5EL工具類

public class SpELUtil {
    /**
     * 校驗(yàn)并返回實(shí)際使用的 spEL 表達(dá)式
     *
     * @param spEl spEL 表達(dá)式
     * @return 實(shí)際使用的 spEL 表達(dá)式
     */
    public static Object parseKey(String spEl, Method method, Object[] contextObj) {
        List<String> spELFlag = ListUtil.of("#", "T(");
        Optional<String> optional = spELFlag.stream().filter(spEl::contains).findFirst();
        if (optional.isPresent()) {
            return parse(spEl, method, contextObj);
        }
        return spEl;
    }

    /**
     * 轉(zhuǎn)換參數(shù)為字符串
     *
     * @param spEl       spEl 表達(dá)式
     * @param contextObj 上下文對象
     * @return 解析的字符串值
     */
    public static Object parse(String spEl, Method method, Object[] contextObj) {
        DefaultParameterNameDiscoverer discoverer = new DefaultParameterNameDiscoverer();
        ExpressionParser parser = new SpelExpressionParser();
        Expression exp = parser.parseExpression(spEl);
        String[] params = discoverer.getParameterNames(method);
        StandardEvaluationContext context = new StandardEvaluationContext();
        if (ArrayUtil.isNotEmpty(params)) {
            for (int len = 0; len < params.length; len++) {
                context.setVariable(params[len], contextObj[len]);
            }
        }
        return exp.getValue(context);
    }
}

5.實(shí)戰(zhàn)使用

使用天機(jī)學(xué)堂項(xiàng)目來進(jìn)行實(shí)戰(zhàn)

5.1寫入common模塊

在這里插入圖片描述

5.2使用

在這里插入圖片描述

直接加上注解就可以

但是實(shí)際上這里不存在冪等問題,因?yàn)閡serId和courseId設(shè)置了唯一索引,所以這里不存在冪等性,不需要加上冪等注解

到此這篇關(guān)于springboot基于注解實(shí)現(xiàn)去重表消息防止重復(fù)消費(fèi)的文章就介紹到這了,更多相關(guān)springboot注解防止重復(fù)消費(fèi)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java在長字符串中查找短字符串的實(shí)現(xiàn)多種方法

    Java在長字符串中查找短字符串的實(shí)現(xiàn)多種方法

    這篇文章主要介紹了Java在長字符串中查找短字符串的實(shí)現(xiàn)多種方法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-12-12
  • Java基礎(chǔ)之FileInputStream和FileOutputStream流詳解

    Java基礎(chǔ)之FileInputStream和FileOutputStream流詳解

    這篇文章主要介紹了Java基礎(chǔ)之FileInputStream和FileOutputStream流詳解,文中有非常詳細(xì)的代碼示例,對正在學(xué)習(xí)java基礎(chǔ)的小伙伴們有非常好的幫助,需要的朋友可以參考下
    2021-04-04
  • Java基于HttpClient實(shí)現(xiàn)RPC的示例

    Java基于HttpClient實(shí)現(xiàn)RPC的示例

    HttpClient可以實(shí)現(xiàn)使用Java代碼完成標(biāo)準(zhǔn)HTTP請求及響應(yīng)。本文主要介紹了Java基于HttpClient實(shí)現(xiàn)RPC,文中通過示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-10-10
  • Java Socket使用加密協(xié)議進(jìn)行傳輸對象的方法

    Java Socket使用加密協(xié)議進(jìn)行傳輸對象的方法

    這篇文章主要介紹了Java Socket使用加密協(xié)議進(jìn)行傳輸對象的方法,結(jié)合實(shí)例形式分析了java socket加密協(xié)議相關(guān)接口與類的調(diào)用方法,以及服務(wù)器、客戶端實(shí)現(xiàn)技巧,需要的朋友可以參考下
    2017-06-06
  • 在Java中解析JSON數(shù)據(jù)代碼示例及說明

    在Java中解析JSON數(shù)據(jù)代碼示例及說明

    這篇文章主要介紹了在Java中解析JSON數(shù)據(jù)的相關(guān)資料,文中講解了如何使用Gson和Jackson庫解析JSON數(shù)據(jù),并展示了如何將日期時(shí)間字符串轉(zhuǎn)換為時(shí)間戳,通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2025-03-03
  • Java枚舉(Enum)從基礎(chǔ)到高級應(yīng)用詳解

    Java枚舉(Enum)從基礎(chǔ)到高級應(yīng)用詳解

    文章詳細(xì)介紹了Java枚舉(enum)的概念、基礎(chǔ)語法、內(nèi)置方法及高級用法,強(qiáng)調(diào)了枚舉在提高代碼類型安全性、可讀性和維護(hù)性方面的優(yōu)勢,文中還列舉了添加字段、方法和實(shí)現(xiàn)接口等高級用法,并提供了命名規(guī)范、不可變性、異常處理等最佳實(shí)踐建議,需要的朋友可以參考下
    2026-05-05
  • 防止SpringMVC攔截器攔截js等靜態(tài)資源文件的解決方法

    防止SpringMVC攔截器攔截js等靜態(tài)資源文件的解決方法

    本篇文章主要介紹了防止SpringMVC攔截器攔截js等靜態(tài)資源文件的解決方法,具有一定的參考價(jià)值,有興趣的同學(xué)可以了解一下
    2017-09-09
  • Springboot jpa @Column命名大小寫問題及解決

    Springboot jpa @Column命名大小寫問題及解決

    這篇文章主要介紹了Springboot jpa @Column命名大小寫問題及解決,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-10-10
  • 詳解Java的Hibernate框架中的Interceptor和Collection

    詳解Java的Hibernate框架中的Interceptor和Collection

    這篇文章主要介紹了Java的Hibernate框架中的Interceptor和Collection,Hibernate是Java的SSH三大web開發(fā)框架之一,需要的朋友可以參考下
    2016-01-01
  • 深入理解java中的拷貝機(jī)制

    深入理解java中的拷貝機(jī)制

    這篇文章主要給大家深入介紹了java中的拷貝機(jī)制,網(wǎng)上關(guān)于java中拷貝的文章也很多,但覺得有必要再深的介紹下java的拷貝機(jī)制,有需要的朋友可以參考學(xué)習(xí),下面來一起看看吧。
    2017-02-02

最新評論

乌恰县| 华容县| 锦屏县| 玉门市| 阜宁县| 加查县| 环江| 图木舒克市| 彭阳县| 合江县| 汉寿县| 东港市| 富宁县| 辽源市| 虹口区| 全椒县| 南靖县| 萍乡市| 平阴县| 江城| 临漳县| 兴安盟| 时尚| 新竹市| 大足县| 武穴市| 定远县| 马山县| 山东省| 荆州市| 六盘水市| 林周县| 酉阳| 曲松县| 周至县| 搜索| 竹山县| 濮阳市| 西乌珠穆沁旗| 马龙县| 宜兴市|