SseEmitter对比Flux<String>

SseEmitter对比Flux<String>

lx 0 2026-07-24

SseEmitter对比Flux

📊 核心对比

对比维度 SseEmitter Flux
Spring 版本 Spring MVC 传统方式 Spring WebFlux 响应式方式
依赖 spring-webmvc spring-webflux
编程模型 命令式(同步阻塞) 响应式(非阻塞)
线程模型 每个请求占用一个 Servlet 线程 事件驱动,少量线程处理大量请求
错误处理 try-catch + onError() 响应式操作符(onErrorResume等)
超时控制 构造函数设置 全局配置或超时操作符
代码复杂度 中等(需手动管理状态) 低(声明式链式调用)
灵活性 高(可精确控制每次发送) 高(响应式操作符丰富)

一、🎯 选择建议

适合用 SseEmitter 的场景:

  • ✅ 项目本身使用 Spring MVC(非 WebFlux)
  • ✅ 需要精细控制每个事件的发送时机
  • ✅ 需要与遗留的命令式代码集成
  • ✅ 团队成员更熟悉传统 Spring MVC 编程模型

适合用 Flux 的场景:

  • ✅ 项目使用 Spring WebFlux
  • ✅ 希望代码更简洁、声明式
  • ✅ 需要高并发、非阻塞的流式处理
  • ✅ 希望更容易进行流式操作(过滤、转换、合并等)

二、核心原理对比

SseEmitter(命令式)

// 本质:一个可变的容器对象,你主动往里放数据
SseEmitter emitter = new SseEmitter();
emitter.send(data);  // 你控制何时发送
emitter.complete();  // 你控制何时结束
emitter.completeWithError(error); // 你控制错误处理

工作原理:

  • 每个请求分配一个线程
  • 线程阻塞等待数据,通过 send() 主动推送
  • 手动管理生命周期(完成/错误/超时)

Flux(响应式)

// 本质:一个数据流的声明,由框架驱动执行
Flux<String> stream = Flux.just("a", "b", "c")
    .delayElements(Duration.ofMillis(100))
    .map(String::toUpperCase);
// 你只管描述"数据从哪里来、怎么变换"
// 框架负责"何时发送、如何发送"

工作原理:

  • 非阻塞,少量线程处理大量请求
  • 数据流自动驱动(订阅时才开始执行)
  • 通过操作符声明式处理(map/filter/merge等)
  • 框架自动管理背压(Backpressure)

三、使用体验对比

场景1:对接 Ollama 流式 API

SseEmitter 方式:

@PostMapping("/chat/stream")
public SseEmitter chatStream(@RequestBody ChatRequest request) {
    SseEmitter emitter = new SseEmitter(60_000L);
  
    // 必须手动管理线程池
    CompletableFuture.runAsync(() -> {
        try {
            // 调用Ollama
            WebClient.create("http://localhost:11434")
                .post()
                .uri("/api/generate")
                .bodyValue(request)
                .retrieve()
                .bodyToFlux(String.class)
                .subscribe(
                    chunk -> {
                        try {
                            // 手动解析并发送
                            String content = parseContent(chunk);
                            emitter.send(content);
                        } catch (IOException e) {
                            emitter.completeWithError(e);
                        }
                    },
                    error -> emitter.completeWithError(error),
                    () -> emitter.complete()
                );
        } catch (Exception e) {
            emitter.completeWithError(e);
        }
    });
  
    return emitter;
}

Flux 方式:

@PostMapping(value = "/chat/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> chatStream(@RequestBody ChatRequest request) {
    return WebClient.create("http://localhost:11434")
        .post()
        .uri("/api/generate")
        .bodyValue(request)
        .retrieve()
        .bodyToFlux(String.class)
        .map(this::parseContent)  // 流式转换
        .onErrorResume(e -> Flux.just("错误: " + e.getMessage()));
}

体验差异:

  • SseEmitter:代码嵌套深(subscribe里面嵌套try-catch),易出错
  • Flux:链式调用,清晰直观,一行流到底

场景2:数据转换和过滤

SseEmitter 方式:

// 需要手动处理每个chunk
List<String> chunks = new ArrayList<>();
ollamaStream.subscribe(
    chunk -> {
        String content = parseContent(chunk);
        if (content != null && !content.isEmpty()) {
            chunks.add(content.toUpperCase()); // 手动转换
            emitter.send(content);
        }
    },
    error -> emitter.completeWithError(error),
    () -> {
        // 手动计算统计信息
        int totalLength = chunks.stream().mapToInt(String::length).sum();
        // 手动发送完成事件
        emitter.send(Map.of("done", true, "totalLength", totalLength));
        emitter.complete();
    }
);

Flux 方式:

return ollamaStream
    .map(this::parseContent)           // 解析
    .filter(Objects::nonNull)          // 过滤空值
    .filter(s -> !s.isEmpty())         // 过滤空串
    .map(String::toUpperCase)          // 转换
    .doOnNext(System.out::println)     // 打印日志
    .collectList()                     // 收集统计
    .flatMapMany(list -> Flux.concat(
        Flux.fromIterable(list),       // 原数据流
        Flux.just("完成!共" + list.size() + "个片段")  // 额外信息
    ));

体验差异:

  • SseEmitter:需要手动管理集合、状态、完成逻辑
  • Flux:丰富的操作符,链式表达意图,代码即文档

场景3:错误处理和重试

SseEmitter 方式:

// 错误处理分散各处
try {
    // ... 业务逻辑
} catch (IOException e) {
    emitter.completeWithError(e);
} catch (TimeoutException e) {
    emitter.completeWithError(new RuntimeException("超时"));
} catch (Exception e) {
    emitter.completeWithError(e);
}
// 重试逻辑需要手动实现(循环+计数器)

Flux 方式:

return ollamaStream
    .retry(3)  // 自动重试3次
    .timeout(Duration.ofSeconds(30)) // 超时处理
    .onErrorResume(e -> {
        // 统一错误处理
        log.error("流式处理失败", e);
        return Flux.just("系统繁忙,请稍后重试");
    })
    .onErrorReturn("默认回复");

体验差异:

  • SseEmitter:错误处理分散,重试需手动实现
  • Flux:丰富的错误处理操作符,声明式、集中式

四、实战建议

如果选 Flux(推荐方案)

// 1. 添加依赖
// build.gradle
implementation 'org.springframework.boot:spring-boot-starter-webflux'

// 2. Controller
@RestController
public class ChatController {
    @PostMapping(value = "/chat/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<ServerSentEvent<Map<String, Object>>> stream(@RequestBody ChatRequest req) {
        return aiService.streamChat(req);
    }
}

// 3. Service
@Service
public class AiService {
    public Flux<ServerSentEvent<Map<String, Object>>> streamChat(ChatRequest req) {
        return webClient.post()
            .uri("/api/generate")
            .bodyValue(req)
            .retrieve()
            .bodyToFlux(OllamaChunk.class)
            .map(this::toSSE)
            .onErrorResume(e -> Flux.just(errorEvent(e)));
    }
}

如果选 SseEmitter(保守方案)

// 1. 使用标准 Spring MVC
// 2. 封装一个工具类简化使用
@Component
public class SseEmitterHelper {
    public <T> void emitFlux(SseEmitter emitter, Flux<T> flux, 
                             Function<T, Object> mapper) {
        flux.subscribe(
            data -> {
                try {
                    emitter.send(mapper.apply(data));
                } catch (IOException e) {
                    emitter.completeWithError(e);
                }
            },
            emitter::completeWithError,
            emitter::complete
        );
    }
}

// 3. 使用时更简洁
@PostMapping("/chat/stream")
public SseEmitter stream(@RequestBody ChatRequest req) {
    SseEmitter emitter = new SseEmitter();
    Flux<String> stream = aiService.getStream(req);
    sseHelper.emitFlux(emitter, stream, content -> 
        Map.of("content", content, "success", true)
    );
    return emitter;
}