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

Jetty Embedded请求串行/可控并行处理方案咨询

嘿,我太懂你这种头疼的情况了——做了个处理高CPU内存任务的Web服务,本来想靠线程池控并发,结果单线程多发请求直接崩了,这确实是QueuedThreadPool的短板,它只能控线程数,管不住单个线程塞进来的大量请求。我给你几个实用的解决方案,亲测靠谱:

方案一:阻塞队列 + 固定线程池,从根源控并发

直接换用LinkedBlockingQueue做任务队列,配合固定大小的线程池,不管多少请求进来,都会先排队,线程池只按你设定的数量(比如2个)来处理,彻底杜绝资源过载。

具体步骤走一遍:

  • 先把每个请求的处理逻辑封装成独立任务:
class ProcessingTask implements Runnable {
    private String[] inputStrings;
    private HttpServletResponse response;

    public ProcessingTask(String[] input, HttpServletResponse resp) {
        this.inputStrings = input;
        this.response = resp;
    }

    @Override
    public void run() {
        // 这里放你的高消耗处理逻辑
        Object result = heavyProcessing(inputStrings);
        // 处理完返回JSON
        try {
            response.setContentType("application/json");
            new ObjectMapper().writeValue(response.getOutputStream(), result);
        } catch (IOException e) {
            response.setStatus(HttpServletResponse.SC_INTERNAL_SERVER_ERROR);
        }
    }
}
  • 初始化队列和线程池,硬限制并发数:
// 队列大小根据服务器内存调,比如设100,避免无限堆积爆内存
LinkedBlockingQueue<Runnable> taskQueue = new LinkedBlockingQueue<>(100);
// 固定2个线程干活,核心和最大线程数设成一样
ExecutorService executor = new ThreadPoolExecutor(2, 2,
        0L, TimeUnit.MILLISECONDS, taskQueue);
  • 在POST接口里把请求丢进队列:
@PostMapping("/process")
public void handleRequest(@RequestBody String[] input, HttpServletResponse response) {
    try {
        executor.submit(new ProcessingTask(input, response));
    } catch (RejectedExecutionException e) {
        // 队列满了就返回繁忙提示
        response.setStatus(HttpServletResponse.SC_SERVICE_UNAVAILABLE);
        response.getWriter().write("{\"error\": \"服务繁忙,请稍后再试\"}");
    }
}

这种方式不管客户端怎么刷请求,哪怕单线程狂发,任务都会乖乖排队,线程池严格控并发,再也不会因为资源耗尽崩服务。

方案二:RxJava Flowable背压,响应式场景更省心

如果你之前试的是RxJava的链式处理,那用Flowable的背压机制就太合适了,它能自动根据处理能力调节请求接收速度,不用手动管队列。

举个简化版例子:

// 创建Flowable,用BUFFER策略缓存请求
Flowable.<String[]>create(emitter -> {
    // 把POST请求的input发射进流
    yourPostHandler.setRequestCallback(input -> emitter.onNext(input));
}, BackpressureStrategy.BUFFER)
        // 限制并发数为2,每个任务在IO线程处理
        .flatMap(input -> Flowable.just(input)
                .subscribeOn(Schedulers.io())
                .map(this::heavyProcessing),
                2)
        // 处理结果或错误
        .subscribe(result -> sendJsonResponse(result),
                   error -> handleProcessingError(error));

背压会帮你处理请求堆积,处理不过来的时候会暂时缓存,或者你可以选丢弃、报错等策略,灵活度很高。

Spring Boot专属简化方案:@Async + 自定义线程池

如果是Spring Boot项目,用@Async配合自定义线程池更简洁,不用自己写队列逻辑:

  • 先配置线程池:
@Configuration
@EnableAsync
public class AsyncConfig {
    @Bean(name = "processingExecutor")
    public Executor processingExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(2); // 核心线程数2
        executor.setMaxPoolSize(2); // 最大线程数2,固定并发
        executor.setQueueCapacity(100); // 队列容量
        executor.setThreadNamePrefix("Processing-Worker-");
        executor.initialize();
        return executor;
    }
}
  • 把处理方法改成异步:
@Async("processingExecutor")
public CompletableFuture<Object> processTask(String[] input) {
    Object result = heavyProcessing(input);
    return CompletableFuture.completedFuture(result);
}
  • 接口用DeferredResult实现非阻塞返回:
@PostMapping("/process")
public DeferredResult<ResponseEntity<Object>> handleRequest(@RequestBody String[] input) {
    DeferredResult<ResponseEntity<Object>> deferredResult = new DeferredResult<>();
    processTask(input)
            .thenAccept(result -> deferredResult.setResult(ResponseEntity.ok(result)))
            .exceptionally(error -> {
                deferredResult.setResult(ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body(error.getMessage()));
                return null;
            });
    return deferredResult;
}

这样客户端发请求后会立刻收到响应等待的信号,服务端后台按设定的并发数处理,完美解决单线程多发请求的问题。

内容的提问来源于stack exchange,提问作者antorqs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 07:54:23