基于ThreadPoolTaskScheduler的定时任务请求限流优化问题
解决Spring定时任务中外部API的RPM限流问题
问题场景
基于Spring ThreadPoolTaskScheduler实现的每日定时任务流程如下:
- 调用数据库检查数据是否可更新
- 向外部服务器发起HTTP请求获取数据列表
- 验证并将列表项转换为领域模型
- 针对子列表调用外部服务器获取附加信息(核心限制:外部服务器仅允许20次/分钟请求)
- 将结果列表保存至数据库
- 保存任务元数据
当前通过Thread.sleep实现限流,导致线程被长时间占用(如100条数据需占用5-6分钟),需优化请求调度方式。
可行解决方案
方案1:使用Guava RateLimiter实现平滑限流
Guava的RateLimiter可以按指定速率平滑控制请求频率,避免手动处理睡眠逻辑,同时保证请求符合RPM限制。
首先引入Guava依赖(Maven):
<dependency> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> <version>32.1.3-jre</version> </dependency>
在调度器中实现限流:
import com.google.common.util.concurrent.RateLimiter; import org.springframework.stereotype.Component; @Component public class YourDataUpdateScheduler extends AbstractUpdateScheduler { // 按20次/分钟设置速率,即每秒1/3次请求 private static final RateLimiter API_RATE_LIMITER = RateLimiter.create(20.0 / 60); @Override public void update() { // 步骤1-3:完成数据检查、获取与转换 List<DomainModel> domainModels = fetchAndTransformData(); // 带限流处理附加信息请求 for (DomainModel model : domainModels) { // 阻塞直到获取请求许可,平滑控制速率 API_RATE_LIMITER.acquire(); AdditionalInfo info = fetchAdditionalInfoFromExternal(model.getId()); model.setAdditionalInfo(info); } // 步骤5-6:保存结果与任务元数据 saveDomainModels(domainModels); saveTaskMetadata(); } // 以下为业务逻辑示例方法 private List<DomainModel> fetchAndTransformData() { /* 实现逻辑 */ } private AdditionalInfo fetchAdditionalInfoFromExternal(String modelId) { /* 外部API调用逻辑 */ } private void saveDomainModels(List<DomainModel> models) { /* 数据库保存逻辑 */ } private void saveTaskMetadata() { /* 元数据保存逻辑 */ } @Override protected void setTaskScheduler(ThreadPoolTaskScheduler taskScheduler) { this.taskScheduler = taskScheduler; } @Override public String getCron() { return "0 0 0 * * ?"; // 每日凌晨执行 } }
方案2:异步线程池+RateLimiter释放主线程
如果不想让定时任务主线程长时间阻塞,可以结合Spring异步线程池与RateLimiter,将API请求放到异步线程中处理,主线程仅负责等待任务完成。
首先配置异步线程池:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.TaskExecutor; @Configuration @EnableAsync public class AsyncApiExecutorConfig { @Bean(name = "apiCallExecutor") public TaskExecutor apiCallExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(20); executor.setThreadNamePrefix("ExternalApi-"); executor.initialize(); return executor; } }
修改调度器实现:
import com.google.common.util.concurrent.RateLimiter; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Component; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TaskExecutor; @Component public class YourDataUpdateScheduler extends AbstractUpdateScheduler { private static final RateLimiter API_RATE_LIMITER = RateLimiter.create(20.0 / 60); private final TaskExecutor apiCallExecutor; public YourDataUpdateScheduler(@Qualifier("apiCallExecutor") TaskExecutor apiCallExecutor) { this.apiCallExecutor = apiCallExecutor; } @Override public void update() { List<DomainModel> domainModels = fetchAndTransformData(); // 异步处理每个API请求,带限流控制 CompletableFuture<?>[] futures = domainModels.stream() .map(model -> CompletableFuture.runAsync(() -> { API_RATE_LIMITER.acquire(); AdditionalInfo info = fetchAdditionalInfoFromExternal(model.getId()); model.setAdditionalInfo(info); }, apiCallExecutor)) .toArray(CompletableFuture[]::new); // 等待所有异步API请求完成 CompletableFuture.allOf(futures).join(); saveDomainModels(domainModels); saveTaskMetadata(); } // 业务逻辑方法同方案1... @Override protected void setTaskScheduler(ThreadPoolTaskScheduler taskScheduler) { this.taskScheduler = taskScheduler; } @Override public String getCron() { return "0 0 0 * * ?"; } }
方案3:自定义时间窗口计数器限流
如果不想引入第三方依赖,可以自己实现基于时间窗口的计数器限流逻辑:
import org.springframework.stereotype.Component; import java.util.List; @Component public class YourDataUpdateScheduler extends AbstractUpdateScheduler { private static final int MAX_REQUESTS_PER_MINUTE = 20; private int requestCount = 0; private long currentWindowStartTime = System.currentTimeMillis(); @Override public void update() { List<DomainModel> domainModels = fetchAndTransformData(); for (DomainModel model : domainModels) { synchronized (this) { long now = System.currentTimeMillis(); // 超过1分钟窗口,重置计数器与窗口时间 if (now - currentWindowStartTime > 60_000) { requestCount = 0; currentWindowStartTime = now; } // 达到限流阈值,等待至当前窗口结束 if (requestCount >= MAX_REQUESTS_PER_MINUTE) { long waitMillis = 60_000 - (now - currentWindowStartTime); try { Thread.sleep(waitMillis); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("限流等待被中断", e); } // 等待后重置窗口 requestCount = 0; currentWindowStartTime = System.currentTimeMillis(); } requestCount++; } AdditionalInfo info = fetchAdditionalInfoFromExternal(model.getId()); model.setAdditionalInfo(info); } saveDomainModels(domainModels); saveTaskMetadata(); } // 业务逻辑方法同方案1... @Override protected void setTaskScheduler(ThreadPoolTaskScheduler taskScheduler) { this.taskScheduler = taskScheduler; } @Override public String getCron() { return "0 0 0 * * ?"; } }
内容的提问来源于stack exchange,提问作者iank
相关产品推荐
相关产品推荐

