如何对Quarkus响应式消息流的吞吐量进行限流?
解决方案:基于令牌桶算法实现Mutiny流限流
针对你的场景,要在Quarkus 1.13.7 + SmallRye响应式消息中实现基于时间窗口的吞吐量限流,同时保持@Blocking(ordered=false)的并发处理能力,推荐用令牌桶算法结合Mutiny算子来实现,无需依赖Akka。
核心思路
令牌桶算法会按固定速率生成令牌,每个请求需要先获取令牌才能执行;当令牌耗尽时,后续请求会等待直到有新令牌生成。这种方式既保证了吞吐量限制,又允许合理的突发请求,完美适配外部组件30请求/秒的限制。
实现步骤
1. 自定义令牌桶限流工具类
实现一个线程安全的令牌桶类,提供异步获取令牌的方法(返回Uni<Void>),完全适配Mutiny的响应式流:
import io.smallrye.mutiny.Uni; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; public class RateLimiter { private final int permitsPerSecond; private final Queue<CompletableFuture<Void>> waitingQueue = new ConcurrentLinkedQueue<>(); private long lastRefillTime; private int availablePermits; public RateLimiter(int permitsPerSecond) { this.permitsPerSecond = permitsPerSecond; this.lastRefillTime = System.currentTimeMillis(); this.availablePermits = permitsPerSecond; // 每秒定时补充令牌,线程设为守护线程避免阻止应用关闭 ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(r -> { Thread t = new Thread(r); t.setDaemon(true); return t; }); scheduler.scheduleAtFixedRate(this::refillPermits, 1, 1, TimeUnit.SECONDS); } private void refillPermits() { long now = System.currentTimeMillis(); long elapsedMs = now - lastRefillTime; // 计算时间窗口内可补充的令牌数 int newPermits = (int) (elapsedMs * permitsPerSecond / 1000L); if (newPermits > 0) { availablePermits = Math.min(availablePermits + newPermits, permitsPerSecond); lastRefillTime = now; // 唤醒等待队列中的请求 while (!waitingQueue.isEmpty() && availablePermits > 0) { CompletableFuture<Void> future = waitingQueue.poll(); if (future != null && !future.isDone()) { future.complete(null); availablePermits--; } } } } // 异步获取令牌,无令牌时会挂起直到有可用令牌 public Uni<Void> acquire() { synchronized (this) { if (availablePermits > 0) { availablePermits--; return Uni.createFrom().voidItem(); } else { CompletableFuture<Void> future = new CompletableFuture<>(); waitingQueue.add(future); return Uni.createFrom().completionStage(future); } } } }
2. 在Kafka消费逻辑中集成限流
在消费方法中,先调用令牌桶的acquire()方法获取令牌,再执行外部组件调用,同时保留@Blocking(ordered=false)以启用并发处理:
import io.smallrye.mutiny.Uni; import io.smallrye.reactive.messaging.annotations.Blocking; import org.eclipse.microprofile.reactive.messaging.Incoming; import org.eclipse.microprofile.reactive.messaging.Message; import javax.enterprise.context.ApplicationScoped; import java.time.Duration; @ApplicationScoped public class KafkaMessageConsumer { private final RateLimiter rateLimiter; public KafkaMessageConsumer() { this.rateLimiter = new RateLimiter(30); // 设置为30请求/秒 } @Incoming("your-kafka-topic") // 替换为你的Kafka主题名称 @Blocking(ordered = false) // 启用无序并发处理,提升吞吐量 public Uni<Void> processMessage(Message<String> message) { return rateLimiter.acquire() // 获取令牌后调用外部组件 .chain(() -> callExternalComponent(message.getPayload())) // 处理成功后确认消息 .invoke(message::ack) // 处理失败后拒绝消息 .onFailure().invoke(message::nack); } // 替换为实际调用外部组件的业务代码 private Uni<String> callExternalComponent(String payload) { return Uni.createFrom().item(payload) .onItem().delayIt().by(Duration.ofMillis(500)); // 模拟请求耗时 } }
3. 优化Quarkus线程池配置
由于使用了@Blocking(ordered=false),需要确保Quarkus的工作线程池大小足够,避免因线程不足限制并发能力。在application.properties中添加:
# 配置工作线程池核心大小和最大大小 quarkus.thread-pool.core-threads=20 quarkus.thread-pool.max-threads=50 # 线程池队列大小,按需调整 quarkus.thread-pool.queue-size=100
关键说明
- 令牌桶的定时补充线程设为守护线程,避免应用关闭时被阻塞。
@Blocking(ordered=false)会将任务提交到Quarkus的工作线程池,实现多线程并发处理;配合令牌桶限流,既能达到外部组件的吞吐量上限,又不会超出限制。- 该方案完全基于Mutiny和Java标准库实现,无需引入第三方依赖,避免额外的线程池切换复杂度。
内容的提问来源于stack exchange,提问作者Maxime
相关产品推荐
相关产品推荐

