You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何对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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 14:07:50