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

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. 线程池耦合导致调度阻塞:原实现使用大小为1的固定线程池作为全局ExecutionContext,HTTP请求的Future调度、回调执行与限流逻辑共用同一个线程,RateLimiter的阻塞获取许可、Future的执行调度互相抢占资源,直接堵死请求提交流程。哪怕将RateLimiter阈值设置得远高于目标值,唯一的工作线程被请求执行占有时,下一轮循环的限流判断、请求提交根本无法按时触发,实际发送速率自然上不去。
  2. 限流逻辑与请求执行逻辑未隔离:Guava RateLimiter的acquire()是阻塞方法,放在单线程循环中时,请求执行的耗时会直接打乱限流的固定节拍,导致限流完全失效。
  3. 恒定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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 16:09:20