Redis + Lua 實(shí)現(xiàn)高性能分布式限流的詳細(xì)過程
在高并發(fā)分布式系統(tǒng)中,限流是保障系統(tǒng)穩(wěn)定性的最后一道防線。當(dāng)電商秒殺、大促活動(dòng)、API 網(wǎng)關(guān)等場景遭遇流量洪峰時(shí),若沒有合理的限流機(jī)制,系統(tǒng)會(huì)因 CPU、內(nèi)存、數(shù)據(jù)庫連接等資源耗盡而發(fā)生雪崩式崩潰。傳統(tǒng)的單機(jī)限流方案在微服務(wù)架構(gòu)下已顯乏力,而基于 Redis+Lua 的分布式限流方案憑借其高性能、原子性和全局一致性,成為業(yè)界主流選擇。本文將從原理出發(fā),深入解析令牌桶算法的核心邏輯,并結(jié)合生產(chǎn)級代碼實(shí)現(xiàn),帶你從零搭建一個(gè)支持多維度組合的分布式限流系統(tǒng)。
一、為什么必須用分布式限流?
1.1 單機(jī)限流的致命缺陷
傳統(tǒng)的單機(jī)限流工具(如 Google Guava 的 RateLimiter)在單體應(yīng)用時(shí)代表現(xiàn)尚可,但在微服務(wù)架構(gòu)下存在無法克服的局限性:
- 全局狀態(tài)不一致:每個(gè)服務(wù)實(shí)例獨(dú)立維護(hù)自己的限流計(jì)數(shù)器,無法跨節(jié)點(diǎn)共享狀態(tài)。例如,若設(shè)置全局 QPS 為 1000,部署 10 個(gè)實(shí)例,理論上每個(gè)實(shí)例應(yīng)限 100QPS,但實(shí)際中負(fù)載均衡可能導(dǎo)致某個(gè)實(shí)例被集中訪問,提前觸發(fā)限流,而其他實(shí)例仍處于空閑狀態(tài),整體系統(tǒng)的實(shí)際承載能力遠(yuǎn)低于預(yù)期。
- 擴(kuò)容困難:當(dāng)業(yè)務(wù)增長需要新增服務(wù)實(shí)例時(shí),必須手動(dòng)重新計(jì)算每個(gè)實(shí)例的限流閾值,運(yùn)維成本極高,且容易出現(xiàn)配置錯(cuò)誤。
- 無法應(yīng)對分布式攻擊:針對 IP、用戶維度的惡意刷接口行為,單機(jī)限流無法識別跨節(jié)點(diǎn)的請求,導(dǎo)致攻擊者可以通過分散請求到不同實(shí)例來繞過限流。
1.2 分布式限流的核心優(yōu)勢
分布式限流通過中心化存儲(chǔ)(Redis)統(tǒng)一管理所有服務(wù)實(shí)例的限流狀態(tài),從根本上解決了單機(jī)限流的問題:
- 全局一致性:所有服務(wù)實(shí)例共享同一份限流數(shù)據(jù),無論部署多少個(gè)實(shí)例,都能嚴(yán)格遵守全局限流閾值。
- 靈活擴(kuò)展:支持動(dòng)態(tài)調(diào)整限流策略,無需重啟服務(wù),可通過配置中心實(shí)時(shí)修改閾值、時(shí)間窗口等參數(shù)。
- 多維度精細(xì)化控制:可以按接口、IP、用戶 ID、設(shè)備號等任意維度進(jìn)行限流,甚至支持多維度組合(如同時(shí)限制某個(gè)用戶在某個(gè) IP 上的請求頻率)。
- 高可用:Redis 本身支持主從復(fù)制、哨兵模式和集群部署,可保證限流服務(wù)的高可用性。
二、技術(shù)選型與核心算法詳解
2.1 技術(shù)棧選型:AOP + Redisson + Lua
本方案采用 Spring AOP + Redisson + Lua 的技術(shù)組合,各組件的職責(zé)如下:
- Spring AOP:通過切面攔截帶自定義注解的方法,實(shí)現(xiàn)無侵入式限流,無需修改業(yè)務(wù)代碼。
- Redisson:Java 生態(tài)中最成熟的 Redis 客戶端,原生支持 Redis 集群模式、Lua 腳本執(zhí)行和分布式鎖,比 Jedis 更適合復(fù)雜的分布式場景。
- Lua 腳本:將限流邏輯封裝在 Lua 腳本中,由 Redis 原子執(zhí)行,避免多線程并發(fā)下的競態(tài)條件,保證限流邏輯的正確性。
2.2 限流算法對比
常見的限流算法有四種:固定窗口、滑動(dòng)窗口、漏桶和令牌桶。每種算法都有其適用場景,我們需要根據(jù)業(yè)務(wù)需求選擇最合適的算法。
(1)固定窗口算法
原理:將時(shí)間劃分為固定大小的窗口(如 1 分鐘),每個(gè)窗口內(nèi)維護(hù)一個(gè)計(jì)數(shù)器,請求到達(dá)時(shí)計(jì)數(shù)器加 1,當(dāng)計(jì)數(shù)器達(dá)到閾值時(shí),拒絕后續(xù)請求,窗口結(jié)束后計(jì)數(shù)器清零。
優(yōu)點(diǎn):實(shí)現(xiàn)簡單,性能高。
缺點(diǎn):存在嚴(yán)重的 "臨界問題"。例如,設(shè)置 1 分鐘限 100QPS,在第 59 秒發(fā)送 100 個(gè)請求,第 1 分 01 秒再發(fā)送 100 個(gè)請求,實(shí)際上在 2 秒內(nèi)處理了 200 個(gè)請求,遠(yuǎn)超限流閾值。
適用場景:對限流精度要求不高的簡單場景。
代碼實(shí)現(xiàn):
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
/**
* 固定窗口計(jì)數(shù)器限流算法
* 優(yōu)點(diǎn):實(shí)現(xiàn)簡單、性能極高
* 缺點(diǎn):存在"臨界問題"(窗口交界處可能出現(xiàn)2倍閾值的突發(fā)流量)
*/
public class FixedWindowRateLimiter {
// 時(shí)間窗口大?。ê撩耄?
private final long windowSizeMs;
// 窗口內(nèi)允許的最大請求數(shù)
private final int maxRequests;
// 當(dāng)前窗口的請求計(jì)數(shù)器(原子類保證線程安全)
private final AtomicInteger counter = new AtomicInteger(0);
// 窗口開始時(shí)間(原子類保證線程安全)
private final AtomicLong windowStartTime = new AtomicLong(System.currentTimeMillis());
public FixedWindowRateLimiter(long windowSizeMs, int maxRequests) {
this.windowSizeMs = windowSizeMs;
this.maxRequests = maxRequests;
}
/**
* 嘗試獲取令牌
* @return true-獲取成功(允許請求),false-獲取失?。ㄏ蘖鳎?
*/
public boolean tryAcquire() {
long currentTime = System.currentTimeMillis();
// 1. 判斷是否進(jìn)入新窗口
if (currentTime - windowStartTime.get() > windowSizeMs) {
// 重置窗口:CAS操作保證只有一個(gè)線程能重置成功
if (windowStartTime.compareAndSet(windowStartTime.get(), currentTime)) {
counter.set(0);
}
}
// 2. 原子性增加計(jì)數(shù)并判斷是否超過閾值
return counter.incrementAndGet() <= maxRequests;
}
// 測試用例
public static void main(String[] args) throws InterruptedException {
// 1秒內(nèi)最多允許5個(gè)請求
FixedWindowRateLimiter limiter = new FixedWindowRateLimiter(1000, 5);
// 模擬10個(gè)并發(fā)請求
for (int i = 0; i < 10; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":成功");
} else {
System.out.println("請求" + requestId + ":被限流");
}
}).start();
}
// 等待1秒后,窗口重置,再發(fā)5個(gè)請求
Thread.sleep(1000);
System.out.println("===== 窗口重置 =====");
for (int i = 10; i < 15; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":成功");
} else {
System.out.println("請求" + requestId + ":被限流");
}
}).start();
}
}
}(2)滑動(dòng)窗口算法
原理:將固定窗口進(jìn)一步劃分為多個(gè)小的時(shí)間片(如將 1 分鐘劃分為 6 個(gè) 10 秒的時(shí)間片),每個(gè)時(shí)間片維護(hù)獨(dú)立的計(jì)數(shù)器。當(dāng)請求到達(dá)時(shí),計(jì)算當(dāng)前時(shí)間所在的窗口(包含最近 6 個(gè)時(shí)間片),將所有時(shí)間片的計(jì)數(shù)器相加,若超過閾值則拒絕請求。窗口會(huì)隨著時(shí)間滑動(dòng),不斷淘汰過期的時(shí)間片。
優(yōu)點(diǎn):解決了固定窗口的臨界問題,限流精度更高。
缺點(diǎn):實(shí)現(xiàn)復(fù)雜,需要維護(hù)多個(gè)時(shí)間片的計(jì)數(shù)器,性能稍差。
適用場景:對限流精度要求較高的場景。
代碼實(shí)現(xiàn):
import java.util.concurrent.atomic.AtomicIntegerArray;
/**
* 滑動(dòng)窗口計(jì)數(shù)器限流算法
* 優(yōu)點(diǎn):解決了固定窗口的"臨界問題",限流精度更高
* 缺點(diǎn):實(shí)現(xiàn)稍復(fù)雜,需要維護(hù)多個(gè)時(shí)間片的計(jì)數(shù)器
*/
public class SlidingWindowRateLimiter {
// 總窗口大?。ê撩耄?
private final long windowSizeMs;
// 每個(gè)小時(shí)間片的大小(毫秒)
private final long sliceSizeMs;
// 時(shí)間片數(shù)量
private final int sliceCount;
// 每個(gè)時(shí)間片的請求計(jì)數(shù)器(原子數(shù)組保證線程安全)
private final AtomicIntegerArray counters;
// 上一次請求的時(shí)間戳
private volatile long lastRequestTime;
// 上一次請求所在的時(shí)間片索引
private volatile int lastSliceIndex;
public SlidingWindowRateLimiter(long windowSizeMs, int sliceCount, int maxRequests) {
this.windowSizeMs = windowSizeMs;
this.sliceCount = sliceCount;
this.sliceSizeMs = windowSizeMs / sliceCount;
this.counters = new AtomicIntegerArray(sliceCount);
this.lastRequestTime = System.currentTimeMillis();
this.lastSliceIndex = getCurrentSliceIndex();
}
/**
* 獲取當(dāng)前時(shí)間對應(yīng)的時(shí)間片索引
*/
private int getCurrentSliceIndex() {
return (int) ((System.currentTimeMillis() / sliceSizeMs) % sliceCount);
}
/**
* 嘗試獲取令牌
* @return true-獲取成功,false-被限流
*/
public synchronized boolean tryAcquire() {
long currentTime = System.currentTimeMillis();
int currentSliceIndex = getCurrentSliceIndex();
// 1. 清除所有過期的時(shí)間片(滑動(dòng)窗口)
long timePassed = currentTime - lastRequestTime;
if (timePassed > windowSizeMs) {
// 超過整個(gè)窗口大小,全部清零
for (int i = 0; i < sliceCount; i++) {
counters.set(i, 0);
}
} else {
// 清除從上次請求到現(xiàn)在之間過期的時(shí)間片
int slicesToClear = (int) (timePassed / sliceSizeMs);
for (int i = 1; i <= slicesToClear; i++) {
int index = (lastSliceIndex + i) % sliceCount;
counters.set(index, 0);
}
}
// 2. 計(jì)算當(dāng)前窗口內(nèi)的總請求數(shù)
int totalRequests = 0;
for (int i = 0; i < sliceCount; i++) {
totalRequests += counters.get(i);
}
// 3. 判斷是否超過閾值
if (totalRequests >= maxRequests) {
return false;
}
// 4. 當(dāng)前時(shí)間片計(jì)數(shù)加1
counters.incrementAndGet(currentSliceIndex);
lastRequestTime = currentTime;
lastSliceIndex = currentSliceIndex;
return true;
}
// 測試用例
public static void main(String[] args) throws InterruptedException {
// 1秒窗口,劃分為10個(gè)100毫秒的時(shí)間片,最多允許5個(gè)請求
SlidingWindowRateLimiter limiter = new SlidingWindowRateLimiter(1000, 10, 5);
// 模擬10個(gè)請求,間隔100毫秒
for (int i = 0; i < 10; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":成功");
} else {
System.out.println("請求" + requestId + ":被限流");
}
}).start();
Thread.sleep(100);
}
}
}(3)漏桶算法
原理:將請求比作水,漏桶比作隊(duì)列,水以恒定的速度從漏桶流出。當(dāng)水流入速度超過流出速度時(shí),多余的水會(huì)溢出(拒絕請求)。
優(yōu)點(diǎn):可以強(qiáng)制限制請求的處理速度,實(shí)現(xiàn)削峰填谷,使系統(tǒng)輸出流量保持平穩(wěn)。
缺點(diǎn):不允許突發(fā)流量,即使系統(tǒng)資源空閑,也只能以固定速度處理請求,無法充分利用系統(tǒng)資源。
適用場景:需要嚴(yán)格控制請求處理速度的場景,如消息隊(duì)列消費(fèi)。
代碼實(shí)現(xiàn):
/**
* 漏桶限流算法
* 優(yōu)點(diǎn):可以強(qiáng)制限制請求的處理速度,輸出流量非常平穩(wěn)
* 缺點(diǎn):不允許突發(fā)流量,即使系統(tǒng)資源空閑也只能以固定速度處理
*/
public class LeakyBucketRateLimiter {
// 桶的容量(最大排隊(duì)請求數(shù))
private final int capacity;
// 漏水速度(每秒處理的請求數(shù))
private final double leakRate;
// 當(dāng)前桶中的水量(排隊(duì)的請求數(shù))
private double currentWater;
// 上次漏水的時(shí)間戳
private long lastLeakTime;
public LeakyBucketRateLimiter(int capacity, double leakRate) {
this.capacity = capacity;
this.leakRate = leakRate;
this.currentWater = 0;
this.lastLeakTime = System.currentTimeMillis();
}
/**
* 嘗試獲取令牌
* @return true-獲取成功(請求進(jìn)入桶中等待處理),false-被限流(桶滿)
*/
public synchronized boolean tryAcquire() {
long currentTime = System.currentTimeMillis();
// 1. 計(jì)算從上次漏水到現(xiàn)在應(yīng)該漏出的水量
double leakedWater = (currentTime - lastLeakTime) / 1000.0 * leakRate;
// 2. 更新當(dāng)前水量(不能小于0)
currentWater = Math.max(0, currentWater - leakedWater);
lastLeakTime = currentTime;
// 3. 判斷桶是否已滿
if (currentWater < capacity) {
currentWater += 1;
return true;
}
return false;
}
// 測試用例
public static void main(String[] args) throws InterruptedException {
// 桶容量10,每秒漏水5個(gè)(即每秒處理5個(gè)請求)
LeakyBucketRateLimiter limiter = new LeakyBucketRateLimiter(10, 5);
// 模擬15個(gè)突發(fā)請求
for (int i = 0; i < 15; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":進(jìn)入桶中");
} else {
System.out.println("請求" + requestId + ":被限流(桶滿)");
}
}).start();
}
// 等待2秒,觀察漏水效果
Thread.sleep(2000);
System.out.println("===== 2秒后 =====");
// 再發(fā)5個(gè)請求
for (int i = 15; i < 20; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":進(jìn)入桶中");
} else {
System.out.println("請求" + requestId + ":被限流(桶滿)");
}
}).start();
}
}
}(4)令牌桶算法(本方案選擇)
原理:系統(tǒng)以恒定的速度向令牌桶中放入令牌,當(dāng)請求到達(dá)時(shí),需要從桶中獲取一個(gè)令牌才能被處理。如果桶中沒有令牌,則拒絕請求。令牌桶的容量是固定的,當(dāng)令牌放滿時(shí),多余的令牌會(huì)被丟棄。
核心優(yōu)勢:
- 允許突發(fā)流量:令牌桶可以積累令牌,當(dāng)系統(tǒng)空閑時(shí),令牌會(huì)逐漸填滿桶,此時(shí)如果有突發(fā)流量到來,可以一次性獲取多個(gè)令牌進(jìn)行處理,充分利用系統(tǒng)資源。
- 平滑限流:令牌以恒定速度放入桶中,避免了固定窗口的臨界問題,使請求處理速度更加平滑。
- 易于實(shí)現(xiàn)多維度限流:每個(gè)限流維度(如接口、IP、用戶)可以維護(hù)獨(dú)立的令牌桶,互不干擾。
令牌桶算法的數(shù)學(xué)模型:
- 設(shè)令牌桶容量為
max_tokens(即最大突發(fā)請求數(shù)) - 令牌生成速率為
rate(即每秒生成的令牌數(shù),等于限流 QPS) - 當(dāng)前令牌數(shù)為
current_tokens - 當(dāng)請求到達(dá)時(shí),若
current_tokens >= 1,則current_tokens -= 1,請求被處理;否則拒絕請求。 - 每隔
1/rate秒,current_tokens += 1,但不超過max_tokens。
代碼實(shí)現(xiàn):
/**
* 令牌桶限流算法(工業(yè)界標(biāo)準(zhǔn))
* 優(yōu)點(diǎn):允許突發(fā)流量,同時(shí)能平滑限流,兼顧性能和靈活性
* 缺點(diǎn):實(shí)現(xiàn)稍復(fù)雜
*/
public class TokenBucketRateLimiter {
// 令牌桶容量(最大突發(fā)請求數(shù))
private final int capacity;
// 令牌生成速度(每秒生成的令牌數(shù),即限流QPS)
private final double tokenRate;
// 當(dāng)前令牌數(shù)
private double currentTokens;
// 上次生成令牌的時(shí)間戳
private long lastTokenTime;
public TokenBucketRateLimiter(int capacity, double tokenRate) {
this.capacity = capacity;
this.tokenRate = tokenRate;
// 初始時(shí)桶是滿的
this.currentTokens = capacity;
this.lastTokenTime = System.currentTimeMillis();
}
/**
* 嘗試獲取令牌
* @return true-獲取成功(允許請求),false-獲取失?。ㄏ蘖鳎?
*/
public synchronized boolean tryAcquire() {
return tryAcquire(1);
}
/**
* 嘗試獲取指定數(shù)量的令牌
* @param permits 需要獲取的令牌數(shù)
* @return true-獲取成功,false-獲取失敗
*/
public synchronized boolean tryAcquire(int permits) {
if (permits <= 0 || permits > capacity) {
return false;
}
long currentTime = System.currentTimeMillis();
// 1. 計(jì)算從上次生成令牌到現(xiàn)在應(yīng)該生成的令牌數(shù)
double generatedTokens = (currentTime - lastTokenTime) / 1000.0 * tokenRate;
// 2. 更新當(dāng)前令牌數(shù)(不能超過桶的容量)
currentTokens = Math.min(capacity, currentTokens + generatedTokens);
lastTokenTime = currentTime;
// 3. 判斷是否有足夠的令牌
if (currentTokens >= permits) {
currentTokens -= permits;
return true;
}
return false;
}
// 測試用例
public static void main(String[] args) throws InterruptedException {
// 令牌桶容量10,每秒生成5個(gè)令牌(即限流QPS=5,最大突發(fā)10個(gè)請求)
TokenBucketRateLimiter limiter = new TokenBucketRateLimiter(10, 5);
// 模擬15個(gè)突發(fā)請求
System.out.println("===== 突發(fā)15個(gè)請求 =====");
for (int i = 0; i < 15; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":成功");
} else {
System.out.println("請求" + requestId + ":被限流");
}
}).start();
}
// 等待1秒,令牌桶會(huì)補(bǔ)充5個(gè)令牌
Thread.sleep(1000);
System.out.println("===== 1秒后 =====");
// 再發(fā)10個(gè)請求
for (int i = 15; i < 25; i++) {
final int requestId = i;
new Thread(() -> {
if (limiter.tryAcquire()) {
System.out.println("請求" + requestId + ":成功");
} else {
System.out.println("請求" + requestId + ":被限流");
}
}).start();
}
}
}三、生產(chǎn)級分布式限流系統(tǒng)實(shí)現(xiàn)
3.1 自定義限流注解:靈活配置多維度策略
首先定義一個(gè) @RateLimit 注解,用于標(biāo)記需要限流的方法,并配置限流參數(shù)。注解支持多維度組合限流、可配置的時(shí)間窗口和降級方法,滿足不同業(yè)務(wù)場景的需求。
import java.lang.annotation.*;
import java.util.concurrent.TimeUnit;
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface RateLimit {
/**
* 限流維度枚舉
*/
enum Dimension {
GLOBAL, // 全局限流(所有請求共享一個(gè)令牌桶)
IP, // IP維度限流(每個(gè)IP一個(gè)令牌桶)
USER // 用戶維度限流(每個(gè)用戶一個(gè)令牌桶)
}
/**
* 限流維度(支持組合,如同時(shí)限制全局和用戶)
*/
Dimension[] dimensions() default {
Dimension.GLOBAL
};
/**
* 時(shí)間窗口內(nèi)允許的最大請求數(shù)(即令牌桶容量)
*/
double count();
/**
* 時(shí)間窗口大小
*/
long interval() default 1;
/**
* 時(shí)間單位
*/
TimeUnit timeUnit() default TimeUnit.SECONDS;
/**
* 降級方法名(限流時(shí)調(diào)用該方法返回結(jié)果)
*/
String fallback() default "";
}- 支持多維度組合,例如
@RateLimit(dimensions = {Dimension.GLOBAL, Dimension.USER}, count = 1000, interval = 1)表示同時(shí)限制全局 QPS 為 1000,且每個(gè)用戶的 QPS 不超過 1000。 - 可配置降級方法,限流時(shí)優(yōu)雅返回自定義結(jié)果,而非直接拋出異常,提升用戶體驗(yàn)。
- 時(shí)間單位靈活,支持秒、分鐘、小時(shí)等多種時(shí)間窗口。
3.2 AOP 切面:無侵入式攔截與邏輯處理
通過 Spring AOP 攔截所有標(biāo)記了 @RateLimit 注解的方法,在方法執(zhí)行前執(zhí)行限流邏輯。切面負(fù)責(zé)生成限流 Key、調(diào)用 Lua 腳本執(zhí)行限流判斷、處理限流結(jié)果和降級邏輯。
import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.reflect.MethodSignature;
import org.redisson.api.RScript;
import org.redisson.api.RedissonClient;
import org.redisson.client.codec.StringCodec;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.web.context.request.RequestContextHolder;
import org.springframework.web.context.request.ServletRequestAttributes;
import javax.annotation.PostConstruct;
import javax.servlet.http.HttpServletRequest;
import java.lang.reflect.Method;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
@Aspect
@Component
public class RateLimitAspect {
@Autowired
private RedissonClient redissonClient;
private String luaScriptSha;
// Lua腳本內(nèi)容(見下文)
private static final String LUA_SCRIPT = "..."
@PostConstruct
public void init() {
// 預(yù)加載 Lua 腳本到 Redis,獲取 SHA 1值,減少網(wǎng)絡(luò)傳輸開銷
this.luaScriptSha =
redissonClient.getScript(StringCodec.INSTANCE).scriptLoad(LUA_SCRIPT);
}
@Around("@annotation(rateLimit)")
public Object around(ProceedingJoinPoint joinPoint, RateLimit rateLimit) throws Throwable {
// 1. 計(jì)算時(shí)間窗口(轉(zhuǎn)換為毫秒)
long intervalMs = rateLimit.timeUnit().toMillis(rateLimit.interval());
// 2. 獲取目標(biāo)類和方法名
String className = joinPoint.getTarget().getClass().getName();
String methodName = joinPoint.getSignature().getName();
// 3. 生成限流Key(使用Hash Tag適配Redis集群)
List<String> keys = generateKeys(className, methodName, rateLimit.dimensions());
// 4. 構(gòu)造Lua腳本參數(shù)
List<Object> args = new ArrayList<>();
args.add(rateLimit.count()); // 最大令牌數(shù)
args.add(intervalMs); // 時(shí)間窗口(毫秒)
args.add(System.currentTimeMillis()); // 當(dāng)前時(shí)間戳
args.add(1); // 每次請求消耗的令牌數(shù)
args.add(intervalMs * 2 / 1000); // Key過期時(shí)間(窗口的2倍)
// 5. 執(zhí)行Lua腳本
RScript script = redissonClient.getScript(StringCodec.INSTANCE);
Long result = script.evalSha(RScript.Mode.READ_WRITE, luaScriptSha,
RScript.ReturnType.VALUE, keys, args.toArray());
// 6. 處理限流結(jié)果
if (result == null || result == 0) {
// 觸發(fā)限流,執(zhí)行降級方法
return handleFallback(joinPoint, rateLimit);
}
// 7. 執(zhí)行業(yè)務(wù)方法
return joinPoint.proceed();
}
/**
* 生成限流Key,使用Hash Tag確保同一方法的所有Key落在同一個(gè)Redis Slot
*/
private List<String> generateKeys(String className, String methodName, RateLimit.Dimension[] dimensions) {
List<String> keys = new ArrayList<>();
// Hash Tag:用{}包裹類名和方法名,確保所有Key落在同一個(gè)Slot
String hashTag = "{" + className + ":" + methodName + "}";
String keyPrefix = "ratelimit:" + hashTag;
for (RateLimit.Dimension dimension : dimensions) {
switch (dimension) {
case GLOBAL:
keys.add(keyPrefix + ":global");
break;
case IP:
String ip = getClientIp();
keys.add(keyPrefix + ":ip:" + ip);
break;
case USER:
// 從請求上下文獲取當(dāng)前用戶ID(需根據(jù)實(shí)際項(xiàng)目調(diào)整)
String userId = getCurrentUserId();
keys.add(keyPrefix + ":user:" + userId);
break;
}
}
return keys;
}
/**
* 執(zhí)行降級方法
*/
private Object handleFallback(ProceedingJoinPoint joinPoint, RateLimit rateLimit) throws Throwable {
String fallbackName = rateLimit.fallback();
if (fallbackName.isEmpty()) {
// 未配置降級方法,拋出默認(rèn)異常
throw new RuntimeException("請求過于頻繁,請稍后再試");
}
// 查找降級方法(優(yōu)先匹配同參數(shù)列表,其次匹配無參方法)
MethodSignature signature = (MethodSignature) joinPoint.getSignature();
Class<?> targetClass = joinPoint.getTarget().getClass();
Class<?>[] parameterTypes = signature.getParameterTypes();
Method fallbackMethod;
try {
fallbackMethod = targetClass.getDeclaredMethod(fallbackName, parameterTypes);
} catch (NoSuchMethodException e) {
fallbackMethod = targetClass.getDeclaredMethod(fallbackName);
}
fallbackMethod.setAccessible(true);
// 調(diào)用降級方法
return fallbackMethod.invoke(joinPoint.getTarget(), joinPoint.getArgs());
}
// 輔助方法:獲取客戶端IP
private String getClientIp() {
HttpServletRequest request = ((ServletRequestAttributes) Objects.requireNonNull(RequestContextHolder.getRequestAttributes())).getRequest();
String ip = request.getHeader("X-Forwarded-For");
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getHeader("Proxy-Client-IP");
}
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getHeader("WL-Proxy-Client-IP");
}
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getRemoteAddr();
}
return ip;
}
// 輔助方法:獲取當(dāng)前用戶ID(需根據(jù)實(shí)際項(xiàng)目實(shí)現(xiàn))
private String getCurrentUserId() {
// 示例:從Token中解析用戶ID
return "123456";
}
}- Redis Cluster 兼容性:使用 Hash Tag(
{className:methodName})將同一方法的所有限流 Key 映射到同一個(gè) Redis Slot。在 Redis Cluster 模式下,Lua 腳本只能操作同一個(gè) Slot 內(nèi)的 Key,否則會(huì)報(bào)錯(cuò)。Hash Tag 通過將 Key 中{}內(nèi)的部分作為分片依據(jù),確保相關(guān) Key 落在同一個(gè)節(jié)點(diǎn)。 - Lua 腳本預(yù)加載:在
@PostConstruct方法中將 Lua 腳本加載到 Redis 并獲取 SHA1 值,后續(xù)通過evalSha調(diào)用腳本,避免每次傳輸完整的腳本內(nèi)容,大幅減少網(wǎng)絡(luò)開銷。 - 智能降級處理:支持兩種降級方法簽名 —— 與原方法同參數(shù)列表的方法和無參方法,提升了降級邏輯的靈活性。
3.3 Lua 腳本:原子性限流的核心
Lua 腳本是整個(gè)限流系統(tǒng)的靈魂,它將所有限流邏輯封裝在一個(gè)腳本中,由 Redis 原子執(zhí)行,確保在高并發(fā)下不會(huì)出現(xiàn)競態(tài)條件。本方案采用兩階段提交的設(shè)計(jì)思想:先檢查所有維度的令牌是否充足,全部通過后再統(tǒng)一扣減令牌,避免部分扣減導(dǎo)致的數(shù)據(jù)不一致。
-- 令牌桶限流Lua腳本
-- KEYS: 限流Key列表(每個(gè)維度一個(gè)Key)
-- ARGV: [1]max_tokens(最大令牌數(shù)), [2]interval_ms(時(shí)間窗口毫秒), [3]now_ms(當(dāng)前時(shí)間戳), [4]permits(每次請求消耗的令牌數(shù)), [5]expire_time(Key過期時(shí)間秒)
local max_tokens = tonumber(ARGV[1])
local interval_ms = tonumber(ARGV[2])
local now_ms = tonumber(ARGV[3])
local permits = tonumber(ARGV[4])
local expire_time = tonumber(ARGV[5])
-- 第一階段:預(yù)檢查所有維度的令牌是否充足
for i, key in ipairs(KEYS) do
local value_key = key .. ":value" -- 存儲(chǔ)當(dāng)前令牌數(shù)的Key
local permits_key = key .. ":permits" -- 存儲(chǔ)請求記錄的ZSet Key
-- 初始化令牌桶(如果不存在)
if redis.call("exists", value_key) == 0 then
redis.call("set", value_key, max_tokens)
end
-- 回收過期令牌:刪除interval_ms之前的請求記錄,并將對應(yīng)的令牌放回桶中
local expired_values = redis.call("zrangebyscore", permits_key, 0, now_ms - interval_ms)
if #expired_values > 0 then
local expired_count = 0
for _, v in ipairs(expired_values) do
-- 解析請求記錄中的令牌數(shù)(格式:request_id:permits)
local _, p = string.match(v, "(.*):(.*)")
expired_count = expired_count + tonumber(p)
end
-- 刪除過期記錄
redis.call("zremrangebyscore", permits_key, 0, now_ms - interval_ms)
-- 回收令牌
local curr_v = tonumber(redis.call("get", value_key))
redis.call("set", value_key, math.min(max_tokens, curr_v + expired_count))
end
-- 檢查當(dāng)前令牌是否足夠
local current_val = tonumber(redis.call("get", value_key))
if current_val < permits then
return 0 -- 任一維度令牌不足,直接返回失敗
end
end
-- 第二階段:所有維度檢查通過,統(tǒng)一扣減令牌
for i, key in ipairs(KEYS) do
local value_key = key .. ":value"
local permits_key = key .. ":permits"
-- 生成唯一請求ID(時(shí)間戳+隨機(jī)數(shù))
local request_id = now_ms .. ":" .. math.random(1000000)
-- 記錄本次請求到ZSet(分?jǐn)?shù)為時(shí)間戳,值為request_id:permits)
redis.call("zadd", permits_key, now_ms, request_id .. ":" .. permits)
-- 扣減令牌
local current_v = tonumber(redis.call("get", value_key))
redis.call("set", value_key, current_v - permits)
-- 設(shè)置Key過期時(shí)間,防止內(nèi)存泄漏
redis.call("expire", value_key, expire_time)
redis.call("expire", permits_key, expire_time)
end
return 1 -- 限流通過- 兩階段提交:先檢查所有維度的令牌,全部通過后再扣減,避免出現(xiàn) "部分維度扣減成功,部分失敗" 的不一致情況。
- ZSET 記錄請求歷史:使用有序集合(ZSET)存儲(chǔ)每次請求的時(shí)間戳和消耗的令牌數(shù),便于精確回收過期令牌。ZSET 的分?jǐn)?shù)為請求時(shí)間戳,值為
request_id:permits格式的字符串。 - 自動(dòng)內(nèi)存管理:通過
zremrangebyscore定期清理過期的請求記錄,并設(shè)置 Key 的過期時(shí)間為時(shí)間窗口的 2 倍,確保過期令牌能被正?;厥?,同時(shí)避免長期占用 Redis 內(nèi)存。 - 原子性保證:整個(gè)腳本在 Redis 中作為一個(gè)原子操作執(zhí)行,即使有多個(gè)請求同時(shí)到達(dá),也不會(huì)出現(xiàn)競態(tài)條件。
四、性能優(yōu)化與生產(chǎn)環(huán)境注意事項(xiàng)
4.1 性能優(yōu)化最佳實(shí)踐
- 減少網(wǎng)絡(luò)往返:
- 預(yù)加載 Lua 腳本:使用
scriptLoad+evalSha替代eval,每次調(diào)用僅傳輸 40 字節(jié)的 SHA1 值,而非完整的腳本內(nèi)容。 - 批量操作:在 Lua 腳本內(nèi)部完成所有 Redis 命令,避免多次網(wǎng)絡(luò)往返。
- 預(yù)加載 Lua 腳本:使用
- 高效數(shù)據(jù)結(jié)構(gòu):
- 使用 String 存儲(chǔ)當(dāng)前令牌數(shù),O (1) 復(fù)雜度的讀取和更新。
- 使用 ZSET 存儲(chǔ)請求記錄,支持按時(shí)間戳范圍查詢和刪除,效率遠(yuǎn)高于 List 或 Hash。
- Redis 配置優(yōu)化:
- 關(guān)閉不必要的持久化:限流數(shù)據(jù)屬于臨時(shí)數(shù)據(jù),丟失后不會(huì)影響業(yè)務(wù),可關(guān)閉 RDB 和 AOF 持久化,提高 Redis 性能。
- 合理設(shè)置連接池:根據(jù)服務(wù)實(shí)例數(shù)量和并發(fā)量調(diào)整 Redisson 的連接池大小,避免連接瓶頸。
- 本地緩存兜底:當(dāng) Redis 出現(xiàn)故障時(shí),使用本地 Guava RateLimiter 作為兜底方案,保證服務(wù)的可用性。
4.2 生產(chǎn)環(huán)境注意事項(xiàng)
- 限流閾值設(shè)置:
- 通過壓測確定系統(tǒng)的最大承載能力,限流閾值應(yīng)設(shè)置為最大承載能力的 80% 左右,預(yù)留一定的緩沖空間。
- 不同接口的限流閾值應(yīng)根據(jù)業(yè)務(wù)重要性和訪問頻率單獨(dú)設(shè)置,核心接口可適當(dāng)提高閾值。
- 監(jiān)控與告警:
- 監(jiān)控限流觸發(fā)次數(shù)、Redis 的 QPS、內(nèi)存使用率等指標(biāo),及時(shí)發(fā)現(xiàn)異常流量。
- 當(dāng)限流觸發(fā)次數(shù)超過閾值時(shí),發(fā)送告警通知,便于運(yùn)維人員及時(shí)處理。
- 降級策略設(shè)計(jì):
- 降級邏輯應(yīng)盡可能簡單,避免降級方法本身成為性能瓶頸。
- 對于非核心接口,可直接返回默認(rèn)值或緩存數(shù)據(jù);對于核心接口,可采用排隊(duì)等待或降級到備用服務(wù)的方式。
- 防刷機(jī)制:
- 對頻繁觸發(fā)限流的 IP 或用戶進(jìn)行臨時(shí)拉黑,防止惡意攻擊。
- 結(jié)合驗(yàn)證碼、滑塊驗(yàn)證等手段,進(jìn)一步提高系統(tǒng)的安全性。
五、總結(jié)與擴(kuò)展
本文詳細(xì)介紹了基于 Redis + Lua 的分布式限流方案,從原理到實(shí)現(xiàn),深入解析了令牌桶算法的核心邏輯和生產(chǎn)級代碼的設(shè)計(jì)要點(diǎn)。該方案具有高性能、原子性、全局一致性和多維度支持等優(yōu)點(diǎn),能夠有效應(yīng)對高并發(fā)場景下的流量沖擊。
未來,我們可以在此基礎(chǔ)上進(jìn)行進(jìn)一步擴(kuò)展:
- 支持動(dòng)態(tài)調(diào)整限流策略:通過配置中心(如 Nacos、Apollo)實(shí)時(shí)修改限流閾值和時(shí)間窗口,無需重啟服務(wù)。
- 支持更復(fù)雜的限流算法:如基于令牌桶的動(dòng)態(tài)限流,根據(jù)系統(tǒng)負(fù)載自動(dòng)調(diào)整限流閾值。
- 支持分布式限流集群:當(dāng)單 Redis 實(shí)例性能不足時(shí),可采用 Redis Cluster 或分片模式,將限流數(shù)據(jù)分散到多個(gè)節(jié)點(diǎn)。
分布式限流是高并發(fā)系統(tǒng)中不可或缺的組件,只有深入理解其原理并結(jié)合實(shí)際業(yè)務(wù)場景進(jìn)行優(yōu)化,才能構(gòu)建出穩(wěn)定可靠的系統(tǒng)。
到此這篇關(guān)于Redis + Lua 實(shí)現(xiàn)高性能分布式限流的文章就介紹到這了,更多相關(guān)Redis Lua高性能分布式限流內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Redis所實(shí)現(xiàn)的Reactor模型設(shè)計(jì)方案
這篇文章主要介紹了Redis所實(shí)現(xiàn)的Reactor模型,本文將帶領(lǐng)讀者從源碼的角度來查看redis關(guān)于reactor模型的設(shè)計(jì),需要的朋友可以參考下2024-06-06
Windows中Redis安裝配置流程并實(shí)現(xiàn)遠(yuǎn)程訪問功能
很多在windows環(huán)境中安裝Redis總是出錯(cuò),今天小編抽空給大家分享在Windows中Redis安裝配置流程并實(shí)現(xiàn)遠(yuǎn)程訪問功能,本文通過圖文并茂的形式給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧2021-06-06
Redis主從配置和底層實(shí)現(xiàn)原理解析(實(shí)戰(zhàn)記錄)
今天給大家分享Redis主從配置和底層實(shí)現(xiàn)原理解析,本文通過實(shí)戰(zhàn)項(xiàng)目給大家源碼解析,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧2021-06-06
詳解redis大幅性能提升之使用管道(PipeLine)和批量(Batch)操作
這篇文章主要介紹了詳解redis大幅性能提升之使用管道(PipeLine)和批量(Batch)操作 ,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下。2016-12-12
RedisDesktopManager?連接redis的方法
這篇文章主要介紹了RedisDesktopManager?連接redis,需要的朋友可以參考下2023-08-08

