Scala中如何实现每秒发送N个请求 精准控制100TPS稳定请求速率
Scala实现恒定100TPS外部服务请求发送方案
问题场景
- 基于Scala开发批处理请求应用,流程为从存储层拉取数据、构造请求后调用外部服务ServiceB API
- ServiceB分配的吞吐量配额为110TPS,需将请求发送速率稳定控制在100TPS以最大化配额利用率
- 已知单次API调用平均耗时约500ms,单线程串行模式下每秒最多发起2次请求,理论上需要50个并行线程才能达到100TPS,但Scala Future的并行机制下难以精准估算所需线程数
- 已尝试基于Guava RateLimiter的实现,将限流阈值设置为300QPS时,实际监控TPS仅为83,3分钟总请求量仅16.4k,未达预期
- 参考负载测试框架的恒定TPS实现逻辑,落地精准速率控制方案
原有问题实现代码:
def main(args: Array[String]): Unit = { val executorService = Executors.newFixedThreadPool(1) // threadpoolsize 1 implicit val executionContext: ExecutionContextExecutor = ExecutionContextFactory.get(executorService) .... val request = getRequest(..) val responseList = ListBuffer[Future[Response]](); val stTime = System.currentTimeMillis() val rateLimiter = RateLimiter.create(300) //guava-ratelimiter while(System.currentTimeMillis() - stTime <= (1000*60*3)) { // Running for 3 mins rateLimiter.acquire(1) responseList += http(request) // using dispatch } //TODO : use different threadpool to process the futures in responseList. }
监控指标:
TPS is 83 Total no of calls to tat api is 16.4k
根因分析
- 线程池耦合导致调度阻塞:原实现使用大小为1的固定线程池作为全局ExecutionContext,HTTP请求的Future调度、回调执行与限流逻辑共用同一个线程,RateLimiter的阻塞获取许可、Future的执行调度互相抢占资源,直接堵死请求提交流程。哪怕将RateLimiter阈值设置得远高于目标值,唯一的工作线程被请求执行占有时,下一轮循环的限流判断、请求提交根本无法按时触发,实际发送速率自然上不去。
- 限流逻辑与请求执行逻辑未隔离:Guava RateLimiter的
acquire()是阻塞方法,放在单线程循环中时,请求执行的耗时会直接打乱限流的固定节拍,导致限流完全失效。 - 恒定TPS控制的核心原则未落地:所有成熟负载测试框架实现固定TPS的核心是彻底解耦节拍控制和请求执行,节拍生成逻辑绝对不能受请求耗时、响应处理的影响。
落地实现方案
核心设计
拆分两个完全独立的线程池,互不抢占资源:
- 节拍调度线程池:固定1个核心线程的调度线程池,仅负责按精确时间间隔生成请求提交任务,不执行任何业务逻辑、不等待请求响应,保证节拍精度
- 请求执行线程池:核心线程数按照
目标TPS * 平均接口耗时(秒) * 冗余系数(1.2~1.5)计算,比如100TPS、500ms平均耗时的场景,设置60~70个核心线程即可,专门负责执行HTTP请求、处理响应逻辑
参考实现代码
import java.util.concurrent.{Executors, ScheduledThreadPoolExecutor, TimeUnit} import scala.collection.mutable.ListBuffer import scala.concurrent.{ExecutionContext, Future} // 自行引入HTTP客户端、请求响应模型依赖 def main(args: Array[String]): Unit = { // 独立的请求执行线程池,预留冗余量避免线程耗尽导致提交阻塞 val requestExecutor = Executors.newFixedThreadPool(70) implicit val requestEC: ExecutionContext = ExecutionContext.fromExecutor(requestExecutor) // 独立的单线程节拍调度池,仅做定时任务触发,无任何阻塞操作 val scheduler = new ScheduledThreadPoolExecutor(1) scheduler.setRemoveOnCancelPolicy(true) val responseList = ListBuffer[Future[Response]]() val runDurationMs = 3 * 60 * 1000 // 总运行时长3分钟 val targetTps = 100 val submitIntervalMs = 1000 / targetTps // 每10ms提交1个请求,折算为100TPS val startTime = System.currentTimeMillis() // 固定间隔调度提交请求,调度逻辑不等待请求返回 scheduler.scheduleAtFixedRate( new Runnable { override def run(): Unit = { if (System.currentTimeMillis() - startTime <= runDurationMs) { val request = getRequest(/* 传入请求构造所需参数 */) // 直接将请求提交到执行池,不阻塞等待结果 responseList += http(request) } else { // 达到运行时长,关闭调度池 scheduler.shutdown() // 等待所有请求完成后执行后续响应处理、关闭资源 Future.sequence(responseList).onComplete { _ => requestExecutor.shutdown() // 自定义响应处理逻辑 } } } }, 0, // 初始无延迟 submitIntervalMs, TimeUnit.MILLISECONDS ) }
优化点
- 若仍需使用Guava RateLimiter,可将
rateLimiter.acquire()放在调度线程run()方法的最开头,此时调度线程独立不受请求执行阻塞,直接设置RateLimiter.create(100)即可精准控速,无需将阈值设置为远高于目标TPS的值,避免突发流量触发ServiceB的限流。 - 给请求执行线程池配置有界队列,队列满时记录异常、做降级处理,避免队列无限堆积反向影响调度线程。
- TPS埋点要放在HTTP请求实际发起的位置,不要放在任务提交位置,避免线程池队列堆积导致的计数偏差。
负载测试框架恒定TPS实现原理
JMeter、Gatling等工具能实现精准恒定TPS,核心逻辑和上述方案一致:
- 调度线程与工作线程完全隔离:独立的高精度调度线程按固定间隔生成请求事件,绝对不会被工作线程的请求执行、响应等待阻塞
- 控速逻辑与业务逻辑完全隔离:限流仅控制请求事件的生成速率,请求实际发送由独立工作线程池承载,工作线程的耗时波动不会反向影响调度节拍
- 自动按目标TPS、接口平均响应时间调整工作线程池大小,预留足够冗余避免线程不足导致请求排队,保证每个生成的请求事件能立刻被发送。
内容的提问来源于stack exchange,提问作者Jeevi
相关产品推荐
相关产品推荐

