在Quarkus中为io.vertx.ext.web.client.WebClient实现请求限流
Quarkus中Vertx WebClient限流解决方案(30秒2请求,等待式限流)
针对你的需求,这里提供两种适配Quarkus环境的实现方案,均支持超出限流阈值的请求排队等待,且可避免触发超时问题:
方案一:基于Vertx原生API自定义令牌桶限流(适配回调式WebClient)
通过实现令牌桶算法,结合Vertx定时器自动补充令牌,请求需先获取令牌再执行,无令牌时进入等待队列。
1. 实现令牌桶限流类
import io.vertx.core.Future; import io.vertx.core.Vertx; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; public class RequestRateLimiter { private final Vertx vertx; private final int maxTokens; private final long refillIntervalMs; private int currentTokens; private final Queue<io.vertx.core.Handler<Void>> waitingQueue = new ConcurrentLinkedQueue<>(); // 构造器:maxTokens=允许的请求数,refillIntervalMs=令牌刷新间隔(毫秒) public RequestRateLimiter(Vertx vertx, int maxTokens, long refillIntervalMs) { this.vertx = vertx; this.maxTokens = maxTokens; this.refillIntervalMs = refillIntervalMs; this.currentTokens = maxTokens; // 定时刷新令牌池 vertx.setPeriodic(refillIntervalMs, id -> { synchronized (this) { currentTokens = maxTokens; // 批量处理等待队列中的请求 while (!waitingQueue.isEmpty() && currentTokens > 0) { io.vertx.core.Handler<Void> handler = waitingQueue.poll(); handler.handle(null); currentTokens--; } } }); } // 获取令牌,有令牌则立即返回成功,无令牌则加入等待队列 public Future<Void> acquire() { synchronized (this) { if (currentTokens > 0) { currentTokens--; return Future.succeededFuture(); } else { return Future.future(promise -> waitingQueue.add(v -> promise.complete())); } } } }
2. 在Quarkus资源中使用限流器
import io.vertx.core.Vertx; import io.vertx.ext.web.client.WebClient; import jakarta.inject.Inject; import jakarta.ws.rs.GET; import jakarta.ws.rs.Path; @Path("/api") public class WebClientResource { @Inject Vertx vertx; @Inject WebClient webClient; // 初始化限流器:30秒内允许2个请求 private final RequestRateLimiter rateLimiter = new RequestRateLimiter(vertx, 2, 30_000); @GET @Path("/fetch") public void fetchExternalData(io.vertx.core.http.HttpServerResponse response) { rateLimiter.acquire() // 获取令牌后执行请求,设置足够长的超时时间避免等待超时 .compose(v -> webClient.get(80, "your-target-host", "/api/data") .send(300_000)) // 设置5分钟超时,覆盖排队等待时间 .onSuccess(resp -> response.end(resp.bodyAsString())) .onFailure(err -> response.setStatusCode(500).end(err.getMessage())); } }
方案二:基于Quarkus Mutiny WebClient + SmallRye Fault Tolerance(响应式风格)
利用Quarkus生态的SmallRye Fault Tolerance组件,通过注解+配置实现限流,支持排队等待,代码更简洁。
1. 添加依赖
在pom.xml中引入SmallRye Fault Tolerance依赖:
<dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-smallrye-fault-tolerance</artifactId> </dependency>
2. 配置限流规则
在application.properties中添加限流配置:
# 定义名为my-webclient-limit的限流规则 smallrye.fault-tolerance.rate-limit.my-webclient-limit.limit=2 smallrye.fault-tolerance.rate-limit.my-webclient-limit.duration=30s # 开启排队等待,避免直接拒绝请求 smallrye.fault-tolerance.rate-limit.my-webclient-limit.queueing.enabled=true # 设置队列最大容量,防止无限排队导致内存溢出 smallrye.fault-tolerance.rate-limit.my-webclient-limit.queueing.capacity=100
3. 编写业务代码
import io.quarkus.vertx.webclient.MutinyWebClient; import jakarta.inject.Inject; import jakarta.ws.rs.GET; import jakarta.ws.rs.Path; import org.eclipse.microprofile.faulttolerance.RateLimit; @Path("/api") public class MutinyWebClientResource { @Inject MutinyWebClient webClient; @GET @Path("/fetch-mutiny") // 绑定限流规则 @RateLimit(name = "my-webclient-limit") public io.smallrye.mutiny.Uni<String> fetchExternalData() { return webClient.get(80, "your-target-host", "/api/data") .send() // 设置请求超时,需覆盖排队等待时间 .ifNoItem().after(300_000).fail() .onItem().transform(resp -> resp.bodyAsString()); } }
关键注意事项
- 超时设置:无论哪种方案,都需要将WebClient的请求超时时间设置得足够长,确保排队等待时间+请求处理时间不会触发超时。
- 队列容量:建议限制等待队列的最大容量,避免大量请求堆积导致内存占用过高。方案一可自行在
acquire方法中添加队列大小判断,方案二通过配置queueing.capacity实现。 - 线程安全:方案一的令牌桶实现已通过同步块和线程安全队列保证线程安全,适配Vertx的事件循环模型。
内容的提问来源于stack exchange,提问作者Kirbylix
相关产品推荐
相关产品推荐

