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;
}