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

基于Redis Streams的實(shí)時(shí)消息處理實(shí)戰(zhàn)指南

 更新時(shí)間:2025年07月16日 09:12:39   作者:淺沫云歸  
這篇文章主要為大家詳細(xì)介紹了在生產(chǎn)環(huán)境中基于 Redis Streams 構(gòu)建實(shí)時(shí)消息處理的完整經(jīng)驗(yàn),包括技術(shù)選型、核心代碼示例、踩坑解決和優(yōu)化方案,希望對(duì)大家有所幫助

業(yè)務(wù)場(chǎng)景描述

在我們公司的電商平臺(tái)中,存在大量異步事件需要實(shí)時(shí)處理,例如用戶下單、庫(kù)存更新、支付回調(diào)等。這些事件對(duì)消息的可靠性、順序性和高吞吐量有較高要求。傳統(tǒng)的消息中間件(如Kafka、RabbitMQ)在運(yùn)維成本或部署復(fù)雜度上存在一定挑戰(zhàn),在部分場(chǎng)景下難以滿足“輕量、低延遲、易集成” 的需求。

經(jīng)過(guò)調(diào)研和驗(yàn)證,Redis 6.0+ 提供的 Streams 特性在嵌入式部署、快速上手方面具有顯著優(yōu)勢(shì)。本篇文章將分享我們?cè)谏a(chǎn)環(huán)境中基于 Redis Streams 構(gòu)建實(shí)時(shí)消息處理的完整經(jīng)驗(yàn),包括技術(shù)選型、核心代碼示例、踩坑解決和優(yōu)化方案。

技術(shù)選型過(guò)程

  • 消息可靠性:Redis Streams 支持持久化,且提供 ACK 機(jī)制和 Pending List,能夠有效追蹤消費(fèi)進(jìn)度。
  • 順序消費(fèi):同一消費(fèi)者組內(nèi),可保證分片流(同一 key)中消息按寫入順序被串行消費(fèi)。
  • 橫向擴(kuò)展:可通過(guò) Stream 分片(多個(gè) Stream Key)或消費(fèi)者組內(nèi)多實(shí)例并行消費(fèi)提高吞吐。
  • 運(yùn)營(yíng)成本:Redis 已是團(tuán)隊(duì)基礎(chǔ)設(shè)施,集群部署與監(jiān)控成熟度高,二次成本低。
  • 客戶端生態(tài):Lettuce、Jedis、Redisson 等客戶端均有支持,編碼友好。

基于以上考量,最終選型 Redis Streams,落地于現(xiàn)有 Redis 集群,無(wú)需額外獨(dú)立中間件部署。

實(shí)現(xiàn)方案詳解

環(huán)境與依賴

Maven 依賴(以 Lettuce 客戶端為例):

<dependencies>
    <dependency>
        <groupId>io.lettuce</groupId>
        <artifactId>lettuce-core</artifactId>
        <version>6.1.5.RELEASE</version>
    </dependency>
    <dependency>
        <groupId>org.slf4j</groupId>
        <artifactId>slf4j-api</artifactId>
        <version>1.7.30</version>
    </dependency>
    <dependency>
        <groupId>ch.qos.logback</groupId>
        <artifactId>logback-classic</artifactId>
        <version>1.2.3</version>
    </dependency>
</dependencies>

SpringBoot 配置(application.yml):

spring:
  redis:
    host: redis-cluster-host
    port: 6379
    password: your_password
    timeout: 2000ms

流程設(shè)計(jì)

  • Producer 將事件寫入 Stream:XADD
  • 多消費(fèi)者(Consumer Group)并行讀取:XREADGROUP
  • 消費(fèi)確認(rèn):XACK
  • 異常消息追蹤:Pending-List 與 XCLAIM 回補(bǔ)處理

生產(chǎn)者實(shí)現(xiàn)

import io.lettuce.core.RedisClient;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.sync.RedisCommands;
import java.util.HashMap;
import java.util.Map;

public class RedisStreamProducer {
    private RedisClient client;
    private StatefulRedisConnection<String, String> connection;
    private RedisCommands<String, String> commands;
    private static final String STREAM_KEY = "orderStream";

    public RedisStreamProducer(String uri) {
        client = RedisClient.create(uri);
        connection = client.connect();
        commands = connection.sync();
    }

    public String sendMessage(Map<String, String> message) {
        // XADD key * field value [field value ...]
        return commands.xadd(STREAM_KEY, message);
    }

    public void shutdown() {
        connection.close();
        client.shutdown();
    }

    public static void main(String[] args) {
        RedisStreamProducer producer = new RedisStreamProducer("redis://:your_password@redis-host:6379/0");
        Map<String, String> order = new HashMap<>();
        order.put("orderId", "123456");
        order.put("userId", "u7890");
        order.put("amount", "258.50");
        String messageId = producer.sendMessage(order);
        System.out.println("消息發(fā)送成功, ID=" + messageId);
        producer.shutdown();
    }
}

消費(fèi)者實(shí)現(xiàn)

import io.lettuce.core.RedisClient;
import io.lettuce.core.StreamMessage;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.sync.RedisCommands;
import io.lettuce.core.models.stream.Consumer;
import io.lettuce.core.models.stream.PendingMessage;

import java.time.Duration;
import java.util.List;
import java.util.Map;

public class RedisStreamConsumer {
    private RedisClient client;
    private StatefulRedisConnection<String, String> connection;
    private RedisCommands<String, String> commands;

    private static final String STREAM_KEY = "orderStream";
    private static final String GROUP_NAME = "orderGroup";
    private static final String CONSUMER_NAME = "consumer-1";

    public RedisStreamConsumer(String uri) {
        client = RedisClient.create(uri);
        connection = client.connect();
        commands = connection.sync();
        // 創(chuàng)建消費(fèi)者組, 如果已創(chuàng)建可 ignore
        try {
            commands.xgroupCreate(STREAM_KEY, GROUP_NAME, "$", true);
        } catch (Exception e) {
            // Group exists
        }
    }

    public void consume() {
        while (true) {
            // 從 Pending List 先處理未 ack 的消息
            List<PendingMessage> pending = commands.xpending(STREAM_KEY, GROUP_NAME, Range.unbounded(), Limit.from(10));
            for (PendingMessage pm : pending) {
                // 重新消費(fèi)
                StreamMessage<String, String> msg = commands.xclaim(
                    STREAM_KEY,
                    GROUP_NAME,
                    CONSUMER_NAME,
                    5000,
                    pm.getId());
                process(msg.getBody());
                commands.xack(STREAM_KEY, GROUP_NAME, pm.getId());
            }

            // 正常讀取新消息
            List<StreamMessage<String, String>> messages = commands.xreadgroup(
                Consumer.from(GROUP_NAME, CONSUMER_NAME),
                XReadArgs.StreamOffset.lastConsumed(STREAM_KEY));
            if (messages != null) {
                for (StreamMessage<String, String> msg : messages) {
                    process(msg.getBody());
                    commands.xack(STREAM_KEY, GROUP_NAME, msg.getId());
                }
            }

            // 輪詢間隔
            try {
                Thread.sleep(200);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
    }

    private void process(Map<String, String> body) {
        // 業(yè)務(wù)處理邏輯
        System.out.println("處理訂單: " + body);
    }

    public void shutdown() {
        connection.close();
        client.shutdown();
    }

    public static void main(String[] args) {
        RedisStreamConsumer consumer = new RedisStreamConsumer("redis://:your_password@redis-host:6379/0");
        consumer.consume();
        consumer.shutdown();
    }
}

踩過(guò)的坑與解決方案

1.消息重復(fù)消費(fèi)

  • 問題:消費(fèi)者處理過(guò)程中拋出異常導(dǎo)致 ack 未發(fā)送,Pending List 中累積大量消息。
  • 解決:定期掃描 Pending List,并結(jié)合 XCLAIM 將“活躍但掛起”消息重新分配給健康消費(fèi)者處理;同時(shí)在業(yè)務(wù)端做好冪等控制。

2.消息積壓與內(nèi)存壓力

  • 問題:Stream 長(zhǎng)度持續(xù)增長(zhǎng),Redis 實(shí)例內(nèi)存壓力上升。
  • 解決:使用 XTRIM MAXLEN ~ N 對(duì)流進(jìn)行修剪,結(jié)合業(yè)務(wù)保留時(shí)間策略,定期分批清理歷史消息。

3.消費(fèi)者實(shí)例重啟后狀態(tài)丟失

  • 問題:未及時(shí)恢復(fù) Pending List 中未處理消息,導(dǎo)致部分消息長(zhǎng)時(shí)間滯留。
  • 解決:消費(fèi)者啟動(dòng)時(shí)優(yōu)先處理 Pending List,再進(jìn)入正常消費(fèi)流程;并通過(guò)定時(shí)任務(wù)對(duì)掛起較久的消息進(jìn)行報(bào)警或二次補(bǔ)償處理。

總結(jié)與最佳實(shí)踐

  • Redis Streams 適合輕量級(jí)、低運(yùn)維成本的實(shí)時(shí)消息場(chǎng)景,結(jié)合 ACK、Pending List 能保證高可靠性。
  • 采用消費(fèi)者組(Consumer Group)可支持橫向擴(kuò)展,讀寫分離與順序消費(fèi)兼得。
  • 業(yè)務(wù)側(cè)必須做好冪等設(shè)計(jì),避免消息重復(fù)帶來(lái)的副作用。
  • 對(duì) Stream 進(jìn)行合理修剪,避免數(shù)據(jù)無(wú)節(jié)制增長(zhǎng)導(dǎo)致內(nèi)存問題。
  • 建議結(jié)合監(jiān)控告警,對(duì) Pending List 長(zhǎng)度、消費(fèi)者積壓情況進(jìn)行實(shí)時(shí)監(jiān)控。

到此這篇關(guān)于基于Redis Streams的實(shí)時(shí)消息處理實(shí)戰(zhàn)指南的文章就介紹到這了,更多相關(guān)Redis Streams消息處理內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 利用Supervisor管理Redis進(jìn)程的方法教程

    利用Supervisor管理Redis進(jìn)程的方法教程

    Supervisor 是可以在類 UNIX 系統(tǒng)中進(jìn)行管理和監(jiān)控各種進(jìn)程的小型系統(tǒng)。它自帶了客戶端和服務(wù)端工具,下面這篇文章主要給大家介紹了關(guān)于利用Supervisor管理Redis進(jìn)程的相關(guān)資料,需要的朋友可以參考借鑒,下面來(lái)一起看看吧。
    2017-08-08
  • Redis?SCAN命令詳解

    Redis?SCAN命令詳解

    SCAN 命令是一個(gè)基于游標(biāo)的迭代器,每次被調(diào)用之后, 都會(huì)向用戶返回一個(gè)新的游標(biāo), 用戶在下次迭代時(shí)需要使用這個(gè)新游標(biāo)作為 SCAN 命令的游標(biāo)參數(shù), 以此來(lái)延續(xù)之前的迭代過(guò)程,這篇文章給大家介紹了Redis?SCAN命令的相關(guān)知識(shí),感興趣的朋友一起看看吧
    2022-07-07
  • Redis如何存儲(chǔ)對(duì)象

    Redis如何存儲(chǔ)對(duì)象

    這篇文章主要介紹了Redis如何存儲(chǔ)對(duì)象,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • redis中bind配置的詳細(xì)步驟

    redis中bind配置的詳細(xì)步驟

    本文主要介紹了redis中bind配置的詳細(xì)步驟,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2023-07-07
  • Redis數(shù)據(jù)結(jié)構(gòu)SortedSet的底層原理解析

    Redis數(shù)據(jù)結(jié)構(gòu)SortedSet的底層原理解析

    這篇文章主要介紹了Redis數(shù)據(jù)結(jié)構(gòu)SortedSet的底層原理解析,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-07-07
  • 虛擬機(jī)下的Redis無(wú)法訪問報(bào)錯(cuò)500解決方法

    虛擬機(jī)下的Redis無(wú)法訪問報(bào)錯(cuò)500解決方法

    這篇文章主要介紹了虛擬機(jī)下的Redis無(wú)法訪問,報(bào)錯(cuò)500解決方法,由于我的redis是在虛擬機(jī)下安裝的,無(wú)法訪問redis的原因是因?yàn)樘摂M機(jī)的ip地址和主機(jī)不同,文中通過(guò)圖文結(jié)合給出了詳細(xì)的解決方法,需要的朋友可以參考下
    2024-02-02
  • Redis中6種緩存更新策略詳解

    Redis中6種緩存更新策略詳解

    Redis作為一款高性能的內(nèi)存數(shù)據(jù)庫(kù),已經(jīng)成為緩存層的首選解決方案,然而,使用緩存時(shí)最大的挑戰(zhàn)在于保證緩存數(shù)據(jù)與底層數(shù)據(jù)源的一致性,本文將介紹Redis中6種緩存更新策略,需要的朋友可以參考下
    2025-05-05
  • Redis中五種數(shù)據(jù)類型簡(jiǎn)單操作

    Redis中五種數(shù)據(jù)類型簡(jiǎn)單操作

    這篇文章主要介紹了Redis中五種數(shù)據(jù)類型簡(jiǎn)單操作的相關(guān)資料,需要的朋友可以參考下
    2017-04-04
  • Redisson 加鎖解鎖的實(shí)現(xiàn)

    Redisson 加鎖解鎖的實(shí)現(xiàn)

    本文主要介紹了Redisson 加鎖解鎖的實(shí)現(xiàn),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2022-08-08
  • 控制Redis的hash的field中的過(guò)期時(shí)間

    控制Redis的hash的field中的過(guò)期時(shí)間

    這篇文章主要介紹了控制Redis的hash的field中的過(guò)期時(shí)間問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-01-01

最新評(píng)論

白玉县| 陵川县| 福清市| 福建省| 临潭县| 沂水县| 昌吉市| 盖州市| 花莲市| 西安市| 呼玛县| 施秉县| 特克斯县| 南川市| 迁西县| 茂名市| 九寨沟县| 汽车| 新竹县| 乌拉特中旗| 辽阳市| 沈阳市| 宣武区| 会泽县| 萍乡市| 平顶山市| 宜川县| 云阳县| 漠河县| 客服| 大新县| 富阳市| 大竹县| 靖宇县| 台山市| 鹤壁市| 大港区| 义马市| 西昌市| 奇台县| 东光县|