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

Spring Webflux双请求合并响应的线程安全问题及优化方案问询

Webflux实现双请求字符串合并返回的正确姿势

关于同步锁的性能问题

没错,Webflux里用同步锁确实会严重拖垮性能。Webflux的核心是靠少量非阻塞线程处理高并发请求,加锁会让线程进入阻塞状态,无法再处理其他请求,直接废掉了Reactor非阻塞模型的优势,导致服务吞吐量暴跌,绝对不建议这么做。

正确的非阻塞实现方案

核心是用Reactor原生的并发安全工具,避免手动操作共享变量的竞态条件。这里可以用AtomicReference来原子化管理待配对的请求,结合Mono.create()创建可控制的响应流。

代码实现

首先定义一个封装待配对请求的类:

private static class PendingRequest {
    private final String content;
    private final Sink<String> sink;

    public PendingRequest(String content, Sink<String> sink) {
        this.content = content;
        this.sink = sink;
    }

    public String getContent() {
        return content;
    }

    public Sink<String> getSink() {
        return sink;
    }
}

然后在控制器中实现核心逻辑:

import reactor.core.publisher.Mono;
import reactor.core.publisher.Sinks;
import java.time.Duration;
import java.util.concurrent.atomic.AtomicReference;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class MergeController {
    // 原子引用,用来保存等待配对的第一个请求
    private final AtomicReference<PendingRequest> pendingRequestRef = new AtomicReference<>();

    @PostMapping(value = "/merge", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Mono<String> mergeStrings(@RequestBody String input) {
        return Mono.create(sink -> {
            // 1. 尝试获取并清除已有的待配对请求
            PendingRequest existingRequest = pendingRequestRef.getAndSet(null);
            if (existingRequest != null) {
                // 2. 找到配对请求,合并字符串并返回给双方
                String mergedResult = existingRequest.getContent() + input;
                existingRequest.getSink().success(mergedResult);
                sink.success(mergedResult);
            } else {
                // 3. 没有待配对请求,尝试将当前请求存入等待队列
                PendingRequest currentRequest = new PendingRequest(input, sink);
                if (!pendingRequestRef.compareAndSet(null, currentRequest)) {
                    // 4. 并发冲突:刚有另一个请求存入,重新尝试配对
                    PendingRequest raceRequest = pendingRequestRef.getAndSet(null);
                    if (raceRequest != null) {
                        String mergedResult = raceRequest.getContent() + input;
                        raceRequest.getSink().success(mergedResult);
                        sink.success(mergedResult);
                    } else {
                        // 极端兜底情况,理论上不会触发
                        sink.error(new IllegalStateException("Failed to pair request due to race condition"));
                    }
                } else {
                    // 5. 成功存入,添加超时清理逻辑,避免内存泄漏
                    sink.onTimeout(Duration.ofSeconds(10), () -> {
                        // 超时后尝试清除自己的待配对记录
                        pendingRequestRef.compareAndSet(currentRequest, null);
                        sink.error(new java.util.concurrent.TimeoutException("No paired request received within 10 seconds"));
                    });
                }
            }
        });
    }
}

逻辑说明

  • 原子操作保证线程安全:用AtomicReference的getAndSet和compareAndSet方法,这些都是JVM级别的非阻塞原子操作,不会引发竞态条件,同时保持Webflux的非阻塞特性。
  • 请求配对逻辑:第一个请求到来时会被存入原子引用,等待第二个请求;第二个请求到来时会取出第一个请求的信息,合并后同时完成两个请求的响应流,让两个客户端都拿到结果。
  • 超时清理:给等待的请求添加超时逻辑,避免因客户端断开或长时间无配对导致的内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 10:47:41