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

基于WebFlux实现流中继的正确方式及相关技术问题咨询

WebFlux流中继实现方案与常见问题解答

一、正确的流中继实现方式

核心是用Sinks.Many<byte[]>结合有界背压策略,替代嵌套的Sinks.One<Flux<byte[]>>——后者会导致背压无法从订阅者(Get请求)传递到生产者(Post请求),而Sinks.Many可以直接转发单个字节数组元素,天然支持背压传递。

1. 核心组件:流管理器

用ConcurrentHashMap管理每个id对应的Sink,确保单生产者/多生产者与订阅者的绑定:

@Component
public class StreamRelayManager {
    private final ConcurrentHashMap<String, Sinks.Many<byte[]>> streamMap = new ConcurrentHashMap<>();

    public Sinks.Many<byte[]> getOrCreateSink(String id) {
        return streamMap.computeIfAbsent(id, key -> 
            // 单订阅者场景用unicast(更高效)
            Sinks.many().unicast().onBackpressureBuffer()
            // 多订阅者场景用multicast,指定缓冲区容量与溢出策略
            // Sinks.many().multicast().onBackpressureBuffer(1024, Sinks.EmitFailureHandler.FAIL_FAST)
        );
    }

    // 可选:定时清理无订阅者的Sink,避免内存泄漏
    @Scheduled(fixedRate = 30000)
    public void cleanIdleSinks() {
        streamMap.entrySet().removeIf(entry -> entry.getValue().currentSubscriberCount() == 0);
    }
}

2. Post控制器(生产者端)

接收字节流并转发到对应Sink,处理背压失败场景:

@RestController
@RequestMapping("/stream")
public class StreamProducerController {
    private final StreamRelayManager relayManager;

    public StreamProducerController(StreamRelayManager relayManager) {
        this.relayManager = relayManager;
    }

    @PostMapping("{id}")
    public Mono<Void> pushStream(@PathVariable String id, @RequestBody Flux<byte[]> dataStream) {
        Sinks.Many<byte[]> sink = relayManager.getOrCreateSink(id);
        return dataStream.doOnNext(data -> {
            Sinks.EmitResult result = sink.tryEmitNext(data);
            // 处理emit失败逻辑
            if (result.isFailure()) {
                if (result == Sinks.EmitResult.FAIL_OVERFLOW) {
                    throw new RuntimeException("Stream buffer overflow for id: " + id);
                } else if (result == Sinks.EmitResult.FAIL_CANCELLED) {
                    relayManager.getStreamMap().remove(id);
                }
            }
        }).then();
    }
}

3. Get控制器(订阅者端)

返回Sink对应的Flux,订阅者取消时自动清理资源:

@RestController
@RequestMapping("/stream")
public class StreamConsumerController {
    private final StreamRelayManager relayManager;

    public StreamConsumerController(StreamRelayManager relayManager) {
        this.relayManager = relayManager;
    }

    @GetMapping("{id}")
    public Flux<byte[]> pullStream(@PathVariable String id) {
        Sinks.Many<byte[]> sink = relayManager.getOrCreateSink(id);
        return sink.asFlux()
                .doOnCancel(() -> relayManager.getStreamMap().remove(id));
    }
}

二、发送端(Post控制器)返回值选择

  • 优先返回Mono<Void>:这是Spring WebFlux的标准做法,能正确向客户端反馈处理状态——数据流发送完毕时完成,出现错误时返回错误信号。
  • 不推荐返回Disposable:框架会将其当作普通对象序列化,不符合响应式规范。
  • 谨慎使用Mono.never():仅适用于生产者持续发送数据且不需要向客户端反馈完成状态的场景,会长期占用连接资源。

三、DataBuffer的正确使用与Double Free问题规避

Double Free通常是因为池化DataBuffer(如Netty的PooledByteBuf)被多次释放导致的,遵循以下规则即可避免:

1. 核心原则

  • 不要手动释放框架传入的DataBuffer,除非你明确创建了它。
  • 多路复用场景下,转发DataBuffer前必须调用DataBufferUtils.retain()增加引用计数,确保每个订阅者释放时不会影响其他订阅者。
  • 消费完DataBuffer后,调用DataBufferUtils.release()释放资源(框架在多数场景下会自动处理,但手动释放更稳妥)。

2. 多路复用场景示例

// Post控制器接收DataBuffer流
@PostMapping(value = "{id}", consumes = MediaType.APPLICATION_OCTET_STREAM_VALUE)
public Mono<Void> pushDataBufferStream(@PathVariable String id, @RequestBody Flux<DataBuffer> dataStream) {
    Sinks.Many<DataBuffer> sink = relayManager.getOrCreateDataBufferSink(id);
    return dataStream
            .doOnNext(buffer -> {
                // 多路复用必须retain,增加引用计数
                DataBufferUtils.retain(buffer);
                Sinks.EmitResult result = sink.tryEmitNext(buffer);
                if (result.isFailure()) {
                    // emit失败时释放缓冲区
                    DataBufferUtils.release(buffer);
                    if (result == Sinks.EmitResult.FAIL_CANCELLED) {
                        relayManager.getStreamMap().remove(id);
                    }
                }
            })
            .then()
            .doFinally(signal -> sink.tryEmitComplete());
}

// Get控制器返回DataBuffer流
@GetMapping(value = "{id}", produces = MediaType.APPLICATION_OCTET_STREAM_VALUE)
public Flux<DataBuffer> pullDataBufferStream(@PathVariable String id) {
    Sinks.Many<DataBuffer> sink = relayManager.getOrCreateDataBufferSink(id);
    return sink.asFlux()
            .doOnCancel(() -> {
                relayManager.getStreamMap().remove(id);
                // 取消时释放剩余缓冲区
                sink.asFlux().subscribe(DataBufferUtils::release);
            })
            .doOnNext(DataBufferUtils::release); // 消费后释放
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 04:15:44