Java響應(yīng)式編程之Flux與SseEmitter深度解析(附詳細(xì)代碼)
本文深入探討Java中Flux和SseEmitter兩種流式響應(yīng)技術(shù),從實際需求出發(fā),詳細(xì)分析其使用方法、底層原理及應(yīng)用場景。
一、為什么需要Flux和SseEmitter
1.1 傳統(tǒng)同步響應(yīng)的痛點(diǎn)
在傳統(tǒng)的Spring MVC開發(fā)中,我們通常使用同步的請求-響應(yīng)模式:
@PostMapping("/chat/message")
public Result<MessageVO> sendMessage(@RequestBody SendMessageReq req) {
// 調(diào)用AI服務(wù)生成回復(fù)(可能需要5-30秒)
String aiResponse = aiService.generateResponse(req.getContent());
return Result.success(new MessageVO(aiResponse));
}
存在的問題:
- 用戶體驗差:用戶發(fā)送消息后需要等待很長時間才能看到完整回復(fù)
- 資源浪費(fèi):一個請求會長時間占用一個線程,降低服務(wù)器并發(fā)能力
- 超時風(fēng)險:長時間處理可能觸發(fā)HTTP超時(默認(rèn)30-60秒)
- 無法感知進(jìn)度:用戶不知道系統(tǒng)是否在處理,容易誤以為系統(tǒng)卡死
1.2 典型應(yīng)用場景
以下場景迫切需要流式響應(yīng)能力:
| 場景 | 問題描述 | 解決方案 |
|---|---|---|
| AI聊天機(jī)器人 | GPT類模型生成回復(fù)需要時間,逐字輸出更友好 | Flux/SSE |
| 大文件處理 | 文件上傳/下載進(jìn)度實時反饋 | SSE |
| 實時監(jiān)控 | 服務(wù)器指標(biāo)、日志流實時推送 | SSE/WebSocket |
| 長任務(wù)執(zhí)行 | 數(shù)據(jù)導(dǎo)入、報表生成等耗時操作的進(jìn)度通知 | SSE |
| 實時通知 | 消息推送、訂單狀態(tài)更新 | SSE/WebSocket |
1.3 流式響應(yīng)的價值
傳統(tǒng)模式:
客戶端 ----請求----> 服務(wù)端
[等待30秒]
客戶端 <---完整響應(yīng)-- 服務(wù)端
流式模式:
客戶端 ----請求----> 服務(wù)端
客戶端 <---數(shù)據(jù)塊1--- 服務(wù)端 (0.5秒)
客戶端 <---數(shù)據(jù)塊2--- 服務(wù)端 (1.0秒)
客戶端 <---數(shù)據(jù)塊3--- 服務(wù)端 (1.5秒)
...
客戶端 <---完成----- 服務(wù)端 (30秒)
核心優(yōu)勢:
- ? 用戶立即看到響應(yīng),體驗提升300%
- ? 線程快速釋放,服務(wù)器吞吐量提升5-10倍
- ? 支持超長響應(yīng),不受HTTP超時限制
- ? 實時反饋進(jìn)度,降低用戶焦慮
二、技術(shù)背景與概念
2.1 SSE (Server-Sent Events)
官方定義: SSE是HTML5標(biāo)準(zhǔn)的一部分,允許服務(wù)器主動向客戶端推送數(shù)據(jù)。
核心特性:
- 基于HTTP協(xié)議,無需額外協(xié)議支持
- 單向通信(服務(wù)器→客戶端)
- 自動重連機(jī)制
- 支持事件ID和自定義事件類型
- 純文本協(xié)議,簡單易用
數(shù)據(jù)格式:
data: 這是第一條消息\n\n
data: 這是第二條消息\n\n
event: custom-event\n
data: {"message": "JSON數(shù)據(jù)"}\n
id: 123\n\n
2.2 Reactive Streams與Flux
Reactive Streams 是JVM上的響應(yīng)式編程規(guī)范,定義了4個核心接口:
public interface Publisher<T> {
void subscribe(Subscriber<? super T> subscriber);
}
public interface Subscriber<T> {
void onSubscribe(Subscription subscription);
void onNext(T item);
void onError(Throwable throwable);
void onComplete();
}
public interface Subscription {
void request(long n); // 背壓控制
void cancel();
}
public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {}
Flux 是Project Reactor對Publisher的實現(xiàn),代表0-N個元素的異步序列:
Flux<T>: 0..N個元素的流 | ├─ onNext(T) * N : 發(fā)射N個元素 ├─ onError(Throwable): 發(fā)生錯誤(終止) └─ onComplete() : 完成信號(終止)
2.3 兩者的關(guān)系
SseEmitter (Spring MVC實現(xiàn)) ├─ 基于Servlet異步支持 ├─ 適合傳統(tǒng)Spring MVC項目 └─ 不依賴響應(yīng)式框架 Flux (響應(yīng)式流) ├─ 基于Reactive Streams規(guī)范 ├─ 需要Spring WebFlux支持 └─ 完整的響應(yīng)式編程能力
三、SseEmitter深度解析
3.1 基本使用
3.1.1 簡單示例
@RestController
@RequestMapping("/api/v1/chat")
public class ChatController {
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamMessage(@RequestParam String message) {
// 創(chuàng)建SseEmitter,設(shè)置超時時間為5分鐘
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);
// 異步處理,避免阻塞主線程
CompletableFuture.runAsync(() -> {
try {
// 模擬AI逐字生成回復(fù)
String response = "這是一個流式響應(yīng)示例";
for (char c : response.toCharArray()) {
emitter.send(String.valueOf(c));
Thread.sleep(100); // 模擬生成延遲
}
emitter.complete(); // 發(fā)送完成信號
} catch (Exception e) {
emitter.completeWithError(e); // 發(fā)送錯誤信號
}
});
return emitter;
}
}
關(guān)鍵點(diǎn)說明:
produces = MediaType.TEXT_EVENT_STREAM_VALUE:必須設(shè)置Content-Type為text/event-streamSseEmitter(timeout):超時時間,0表示永不超時(不推薦)CompletableFuture.runAsync():異步執(zhí)行,避免阻塞Tomcat線程emitter.complete():必須調(diào)用,否則客戶端連接不會關(guān)閉
3.1.2 前端對接代碼
// 原生EventSource API
const eventSource = new EventSource('/api/v1/chat/stream?message=你好');
eventSource.onmessage = (event) => {
console.log('收到數(shù)據(jù):', event.data);
document.getElementById('response').innerText += event.data;
};
eventSource.onerror = (error) => {
console.error('連接錯誤:', error);
eventSource.close();
};
// 監(jiān)聽完成事件(需要服務(wù)端發(fā)送特殊事件)
eventSource.addEventListener('complete', () => {
console.log('流式響應(yīng)完成');
eventSource.close();
});
3.2 生產(chǎn)級實現(xiàn)
3.2.1 完整的聊天流式接口
@Slf4j
@RestController
@RequestMapping("/api/v1/chat")
@RequiredArgsConstructor
public class StreamChatController {
private final ChatService chatService;
private final ExecutorService executorService;
@PostMapping(value = "/conversations/{sessionId}/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamChat(
@PathVariable String sessionId,
@RequestBody SendMessageReq req) {
long startTime = System.currentTimeMillis();
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L); // 5分鐘超時
// 設(shè)置回調(diào)處理
emitter.onCompletion(() -> {
log.info("SSE連接正常完成, sessionId={}, cost={}ms",
sessionId, System.currentTimeMillis() - startTime);
});
emitter.onTimeout(() -> {
log.warn("SSE連接超時, sessionId={}", sessionId);
emitter.complete();
});
emitter.onError(throwable -> {
log.error("SSE連接異常, sessionId={}, error={}",
sessionId, throwable.getMessage(), throwable);
});
// 異步處理消息
executorService.execute(() -> {
try {
// 1. 發(fā)送開始事件
emitter.send(SseEmitter.event()
.name("start")
.data(Map.of("messageId", generateMessageId())));
// 2. 調(diào)用AI服務(wù)流式生成
StringBuilder fullResponse = new StringBuilder();
chatService.streamGenerate(sessionId, req.getContent(), chunk -> {
try {
fullResponse.append(chunk);
// 發(fā)送數(shù)據(jù)塊
emitter.send(SseEmitter.event()
.name("message")
.data(Map.of("content", chunk)));
} catch (IOException e) {
throw new RuntimeException("發(fā)送數(shù)據(jù)失敗", e);
}
});
// 3. 保存完整消息到數(shù)據(jù)庫
chatService.saveMessage(sessionId, fullResponse.toString());
// 4. 發(fā)送完成事件
emitter.send(SseEmitter.event()
.name("complete")
.data(Map.of(
"totalChunks", fullResponse.length(),
"costTime", System.currentTimeMillis() - startTime
)));
emitter.complete();
} catch (Exception e) {
log.error("流式生成失敗", e);
try {
emitter.send(SseEmitter.event()
.name("error")
.data(Map.of("message", e.getMessage())));
} catch (IOException ioException) {
log.error("發(fā)送錯誤事件失敗", ioException);
}
emitter.completeWithError(e);
}
});
return emitter;
}
private String generateMessageId() {
return "msg_" + System.currentTimeMillis();
}
}
3.2.2 線程池配置
@Configuration
public class AsyncConfig {
@Bean(name = "sseExecutor")
public ExecutorService sseExecutor() {
return new ThreadPoolExecutor(
10, // 核心線程數(shù)
50, // 最大線程數(shù)
60L, TimeUnit.SECONDS, // 線程空閑時間
new LinkedBlockingQueue<>(100), // 任務(wù)隊列
new ThreadFactoryBuilder()
.setNameFormat("sse-executor-%d")
.build(),
new ThreadPoolExecutor.CallerRunsPolicy() // 拒絕策略
);
}
}
3.3 高級特性
3.3.1 自定義事件類型
// 服務(wù)端發(fā)送不同類型的事件
emitter.send(SseEmitter.event()
.id("123") // 事件ID(用于斷線重連)
.name("chat-message") // 自定義事件名
.data(messageData) // 數(shù)據(jù)
.comment("這是注釋") // 注釋(客戶端不處理)
.reconnectTime(3000L)); // 重連時間(毫秒)
// 客戶端監(jiān)聽特定事件
eventSource.addEventListener('chat-message', (event) => {
const data = JSON.parse(event.data);
console.log('收到聊天消息:', data);
});
3.3.2 斷線重連機(jī)制
@GetMapping(value = "/stream-with-resume", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamWithResume(@RequestHeader(value = "Last-Event-ID", required = false) String lastEventId) {
SseEmitter emitter = new SseEmitter();
// 從lastEventId位置繼續(xù)發(fā)送
int startIndex = lastEventId != null ? Integer.parseInt(lastEventId) : 0;
CompletableFuture.runAsync(() -> {
try {
for (int i = startIndex; i < 100; i++) {
emitter.send(SseEmitter.event()
.id(String.valueOf(i)) // 設(shè)置事件ID
.data("數(shù)據(jù)塊 " + i));
Thread.sleep(100);
}
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
});
return emitter;
}
// 客戶端自動使用Last-Event-ID重連
const eventSource = new EventSource('/stream-with-resume');
// EventSource會自動在重連時發(fā)送Last-Event-ID請求頭
四、Flux深度解析
4.1 引入依賴
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
注意: Flux需要Spring WebFlux支持,可以與Spring MVC共存,但需要注意:
- WebFlux運(yùn)行在Netty上(默認(rèn)),也可以配置為Tomcat
- 兩者可以同時存在,但推薦全面切換到WebFlux以發(fā)揮最大性能
4.2 基本使用
4.2.1 簡單示例
@RestController
@RequestMapping("/api/v1/reactive")
public class ReactiveChatController {
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> streamMessage(@RequestParam String message) {
return Flux.interval(Duration.ofMillis(100)) // 每100ms發(fā)射一個元素
.map(i -> "字符" + i)
.take(10) // 只取前10個
.doOnComplete(() -> log.info("流完成"));
}
}
關(guān)鍵概念:
Flux.interval():周期性發(fā)射元素map():轉(zhuǎn)換元素take(n):限制元素數(shù)量doOnComplete():完成時的副作用操作
4.2.2 實際AI聊天場景
@RestController
@RequestMapping("/api/v1/reactive/chat")
@RequiredArgsConstructor
public class ReactiveStreamController {
private final ReactiveAIService aiService;
@PostMapping(value = "/conversations/{sessionId}/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<MessageChunk>> streamChat(
@PathVariable String sessionId,
@RequestBody SendMessageReq req) {
return aiService.streamGenerate(req.getContent())
.map(chunk -> ServerSentEvent.<MessageChunk>builder()
.event("message")
.data(new MessageChunk(chunk))
.build())
.concatWith(Flux.just(
ServerSentEvent.<MessageChunk>builder()
.event("complete")
.data(new MessageChunk(""))
.build()
))
.doOnSubscribe(sub -> log.info("客戶端訂閱, sessionId={}", sessionId))
.doOnComplete(() -> log.info("流完成, sessionId={}", sessionId))
.doOnError(e -> log.error("流異常, sessionId={}", sessionId, e));
}
}
@Data
@AllArgsConstructor
class MessageChunk {
private String content;
}
4.3 核心操作符詳解
4.3.1 創(chuàng)建Flux
// 1. 從集合創(chuàng)建
Flux<String> flux1 = Flux.fromIterable(Arrays.asList("A", "B", "C"));
// 2. 從數(shù)組創(chuàng)建
Flux<String> flux2 = Flux.fromArray(new String[]{"A", "B", "C"});
// 3. 從Stream創(chuàng)建
Flux<String> flux3 = Flux.fromStream(Stream.of("A", "B", "C"));
// 4. 動態(tài)生成
Flux<Integer> flux4 = Flux.generate(
() -> 0, // 初始狀態(tài)
(state, sink) -> { // 生成邏輯
sink.next(state); // 發(fā)射元素
if (state == 10) sink.complete(); // 完成信號
return state + 1; // 下一個狀態(tài)
}
);
// 5. 異步創(chuàng)建(最靈活)
Flux<String> flux5 = Flux.create(sink -> {
// 模擬從外部回調(diào)接收數(shù)據(jù)
externalService.onData(data -> sink.next(data));
externalService.onComplete(() -> sink.complete());
externalService.onError(error -> sink.error(error));
});
// 6. 定時器
Flux<Long> flux6 = Flux.interval(Duration.ofSeconds(1)); // 每秒發(fā)射一個遞增的Long
4.3.2 轉(zhuǎn)換操作符
Flux<String> source = Flux.just("hello", "world");
// 1. map - 一對一轉(zhuǎn)換
Flux<String> mapped = source.map(String::toUpperCase);
// "hello" -> "HELLO", "world" -> "WORLD"
// 2. flatMap - 一對多轉(zhuǎn)換(異步)
Flux<String> flatMapped = source.flatMap(word ->
Flux.fromArray(word.split("")) // "hello" -> ["h","e","l","l","o"]
);
// 3. concatMap - 一對多轉(zhuǎn)換(保序)
Flux<String> concatMapped = source.concatMap(word ->
Flux.fromArray(word.split(""))
);
// 4. filter - 過濾
Flux<String> filtered = source.filter(s -> s.length() > 4);
// 5. distinct - 去重
Flux<String> distinct = Flux.just("A", "B", "A").distinct();
// 6. take/skip - 限制
Flux<String> taken = source.take(5); // 取前5個
Flux<String> skipped = source.skip(2); // 跳過前2個
// 7. buffer - 批處理
Flux<List<String>> buffered = source.buffer(10); // 每10個元素打包成一個List
4.3.3 組合操作符
Flux<String> flux1 = Flux.just("A", "B");
Flux<String> flux2 = Flux.just("C", "D");
// 1. concat - 順序連接
Flux<String> concat = Flux.concat(flux1, flux2);
// 輸出: A, B, C, D
// 2. merge - 交錯合并(異步)
Flux<String> merged = Flux.merge(flux1, flux2);
// 輸出: A, C, B, D (順序不確定)
// 3. zip - 配對組合
Flux<String> zipped = Flux.zip(flux1, flux2, (a, b) -> a + b);
// 輸出: "AC", "BD"
// 4. combineLatest - 最新值組合
Flux<String> combined = Flux.combineLatest(flux1, flux2, (a, b) -> a + b);
4.3.4 錯誤處理
Flux<String> flux = Flux.just("A", "B", "C")
.map(s -> {
if (s.equals("B")) throw new RuntimeException("錯誤");
return s;
});
// 1. onErrorReturn - 錯誤時返回默認(rèn)值
Flux<String> handled1 = flux.onErrorReturn("默認(rèn)值");
// 2. onErrorResume - 錯誤時切換到備用流
Flux<String> handled2 = flux.onErrorResume(e ->
Flux.just("備用1", "備用2")
);
// 3. onErrorContinue - 錯誤時跳過并繼續(xù)
Flux<String> handled3 = flux.onErrorContinue((e, obj) ->
log.error("處理 {} 時出錯: {}", obj, e.getMessage())
);
// 4. retry - 重試
Flux<String> retried = flux.retry(3); // 失敗時重試3次
// 5. retryWhen - 自定義重試策略
Flux<String> retriedWhen = flux.retryWhen(
Retry.backoff(3, Duration.ofSeconds(1)) // 指數(shù)退避重試
);
4.4 背壓(Backpressure)處理
背壓是響應(yīng)式編程的核心特性,用于處理生產(chǎn)者速度>消費(fèi)者速度的情況。
Flux<Integer> fastProducer = Flux.range(1, 1000)
.delayElements(Duration.ofMillis(1)); // 每毫秒生產(chǎn)一個
// 1. onBackpressureBuffer - 緩沖(可能OOM)
Flux<Integer> buffered = fastProducer
.onBackpressureBuffer(100, // 緩沖區(qū)大小
dropped -> log.warn("丟棄: {}", dropped));
// 2. onBackpressureDrop - 丟棄新元素
Flux<Integer> dropped = fastProducer
.onBackpressureDrop(dropped -> log.warn("丟棄: {}", dropped));
// 3. onBackpressureLatest - 只保留最新元素
Flux<Integer> latest = fastProducer.onBackpressureLatest();
// 4. onBackpressureError - 拋出異常
Flux<Integer> error = fastProducer.onBackpressureError();
4.5 完整生產(chǎn)案例
@Service
@Slf4j
@RequiredArgsConstructor
public class ReactiveAIService {
private final OpenAIClient openAIClient;
private final MessageRepository messageRepository;
/**
* 流式生成AI回復(fù)
*/
public Flux<String> streamGenerate(String sessionId, String prompt) {
return Flux.create(sink -> {
String messageId = generateMessageId();
StringBuilder fullResponse = new StringBuilder();
try {
// 調(diào)用OpenAI流式API
openAIClient.streamChatCompletion(prompt, new StreamCallback() {
@Override
public void onChunk(String chunk) {
fullResponse.append(chunk);
sink.next(chunk); // 發(fā)射數(shù)據(jù)塊
}
@Override
public void onComplete() {
// 保存完整消息
saveMessage(sessionId, messageId, fullResponse.toString())
.subscribe(
saved -> log.info("消息保存成功: {}", messageId),
error -> log.error("消息保存失敗", error)
);
sink.complete(); // 完成信號
}
@Override
public void onError(Throwable error) {
sink.error(error); // 錯誤信號
}
});
} catch (Exception e) {
sink.error(e);
}
})
.publishOn(Schedulers.boundedElastic()) // 切換到彈性線程池
.doOnSubscribe(sub -> log.info("開始生成, sessionId={}", sessionId))
.doOnComplete(() -> log.info("生成完成, sessionId={}", sessionId))
.doOnError(e -> log.error("生成失敗, sessionId={}", sessionId, e));
}
/**
* 響應(yīng)式保存消息
*/
private Mono<Message> saveMessage(String sessionId, String messageId, String content) {
return Mono.fromCallable(() -> {
Message message = new Message();
message.setMessageId(messageId);
message.setSessionId(sessionId);
message.setContent(content);
message.setRole(MessageRole.ASSISTANT);
return messageRepository.save(message);
}).subscribeOn(Schedulers.boundedElastic());
}
}
五、底層原理剖析
5.1 SseEmitter底層原理
5.1.1 Servlet異步機(jī)制
SseEmitter基于Servlet 3.0的異步支持實現(xiàn):
// Spring MVC底層處理流程
public class SseEmitter extends ResponseBodyEmitter {
public SseEmitter(Long timeout) {
super();
this.timeout = timeout;
}
@Override
protected void extendResponse(ServerHttpResponse outputMessage) {
super.extendResponse(outputMessage);
// 設(shè)置SSE必需的響應(yīng)頭
outputMessage.getHeaders().setContentType(MediaType.TEXT_EVENT_STREAM);
outputMessage.getHeaders().setCacheControl(CacheControl.noCache());
}
}
關(guān)鍵實現(xiàn):
// ResponseBodyEmitter核心實現(xiàn)
public abstract class ResponseBodyEmitter {
private final Set<DataWithMediaType> earlySendAttempts = new LinkedHashSet<>(4);
private Handler handler;
private boolean complete;
public void send(Object object) throws IOException {
if (this.handler != null) {
try {
this.handler.send(object, null); // 直接寫入響應(yīng)流
} catch (IOException ex) {
throw ex;
} catch (Throwable ex) {
throw new IllegalStateException("Failed to send " + object, ex);
}
} else {
// 請求還未初始化完成,先緩存
this.earlySendAttempts.add(new DataWithMediaType(object, null));
}
}
public void complete() {
if (this.handler != null) {
this.handler.complete();
} else {
this.complete = true;
}
}
}
5.1.2 異步Servlet工作流程
1. 客戶端發(fā)起請求 | 2. Tomcat接收請求,分配線程A處理 | 3. Controller返回SseEmitter | 4. Spring MVC調(diào)用AsyncContext.startAsync() | 5. 線程A釋放,返回線程池 | 6. 業(yè)務(wù)邏輯在線程B中執(zhí)行 | 7. 調(diào)用emitter.send()寫入數(shù)據(jù) | 8. 數(shù)據(jù)通過AsyncContext寫入TCP緩沖區(qū) | 9. Tomcat的Poller線程監(jiān)聽Socket可寫事件 | 10. 數(shù)據(jù)發(fā)送到客戶端 | 11. emitter.complete()關(guān)閉連接
核心代碼(Spring源碼):
// DeferredResultInterceptor.java
public class ResponseBodyEmitterReturnValueHandler implements HandlerMethodReturnValueHandler {
@Override
public void handleReturnValue(Object returnValue, ...) throws Exception {
ResponseBodyEmitter emitter = (ResponseBodyEmitter) returnValue;
// 啟動異步上下文
WebAsyncManager asyncManager = WebAsyncUtils.getAsyncManager(webRequest);
DeferredResult<?> deferredResult = new DeferredResult<>();
asyncManager.startDeferredResultProcessing(deferredResult, ...);
// 設(shè)置發(fā)送處理器
emitter.initialize(new HttpMessageConvertingHandler(outputMessage, ...));
}
}
5.1.3 數(shù)據(jù)發(fā)送流程
// 數(shù)據(jù)如何寫入客戶端
class HttpMessageConvertingHandler implements ResponseBodyEmitter.Handler {
@Override
public void send(Object data, MediaType mediaType) throws IOException {
// 1. 序列化數(shù)據(jù)
if (data instanceof String) {
String text = (String) data;
// SSE格式:data: xxx\n\n
String formattedData = "data: " + text + "\n\n";
byte[] bytes = formattedData.getBytes(StandardCharsets.UTF_8);
// 2. 寫入OutputStream
this.outputMessage.getBody().write(bytes);
// 3. flush立即發(fā)送(重要?。?
this.outputMessage.getBody().flush();
}
}
@Override
public void complete() {
// 關(guān)閉輸出流
this.outputMessage.getBody().close();
}
}
為什么需要flush()?
- HTTP響應(yīng)默認(rèn)是緩沖的,只有flush才會立即發(fā)送到客戶端
- SSE的實時性依賴于每次send都立即flush
5.2 Flux底層原理
5.2.1 Reactive Streams協(xié)議
Flux實現(xiàn)了Reactive Streams規(guī)范,核心是背壓控制:
// 完整的訂閱流程
Flux<String> flux = Flux.just("A", "B", "C");
flux.subscribe(new Subscriber<String>() {
private Subscription subscription;
@Override
public void onSubscribe(Subscription s) {
this.subscription = s;
// 請求1個元素(背壓控制)
s.request(1);
}
@Override
public void onNext(String item) {
System.out.println("收到: " + item);
// 處理完后再請求下一個
subscription.request(1);
}
@Override
public void onError(Throwable t) {
System.err.println("錯誤: " + t);
}
@Override
public void onComplete() {
System.out.println("完成");
}
});
關(guān)鍵點(diǎn):
- 消費(fèi)者通過
request(n)主動拉取數(shù)據(jù),避免被淹沒 - 生產(chǎn)者必須尊重請求數(shù)量,不能無限推送
5.2.2 操作符鏈?zhǔn)皆?/h4>
// 每個操作符都返回一個新的Flux,形成鏈?zhǔn)浇Y(jié)構(gòu)
Flux<Integer> flux = Flux.range(1, 10) // FluxRange
.map(i -> i * 2) // FluxMap
.filter(i -> i > 5) // FluxFilter
.take(5); // FluxTake
// 實際結(jié)構(gòu)
FluxTake(
FluxFilter(
FluxMap(
FluxRange(1, 10)
)
)
)
// 每個操作符都返回一個新的Flux,形成鏈?zhǔn)浇Y(jié)構(gòu)
Flux<Integer> flux = Flux.range(1, 10) // FluxRange
.map(i -> i * 2) // FluxMap
.filter(i -> i > 5) // FluxFilter
.take(5); // FluxTake
// 實際結(jié)構(gòu)
FluxTake(
FluxFilter(
FluxMap(
FluxRange(1, 10)
)
)
)
訂閱傳播:
// 簡化的FluxMap實現(xiàn)
class FluxMap<T, R> extends Flux<R> {
private final Flux<T> source;
private final Function<T, R> mapper;
@Override
public void subscribe(Subscriber<? super R> actual) {
// 創(chuàng)建包裝訂閱者
source.subscribe(new MapSubscriber<>(actual, mapper));
}
static class MapSubscriber<T, R> implements Subscriber<T> {
private final Subscriber<? super R> actual;
private final Function<T, R> mapper;
@Override
public void onNext(T item) {
R mapped = mapper.apply(item); // 轉(zhuǎn)換
actual.onNext(mapped); // 傳遞給下游
}
// ... 其他方法
}
}
完整調(diào)用鏈:
訂閱(subscribe):從最外層向內(nèi)傳播 FluxTake → FluxFilter → FluxMap → FluxRange 數(shù)據(jù)流(onNext):從最內(nèi)層向外傳播 FluxRange → FluxMap → FluxFilter → FluxTake → Subscriber
5.2.3 線程調(diào)度原理
// Schedulers的本質(zhì)是Executor包裝
public abstract class Schedulers {
// 立即執(zhí)行(當(dāng)前線程)
static Scheduler immediate() {
return ImmediateScheduler.INSTANCE;
}
// 單線程
static Scheduler single() {
return SingleScheduler.INSTANCE;
}
// 彈性線程池(IO密集)
static Scheduler boundedElastic() {
return BoundedElasticScheduler.INSTANCE;
}
// 并行線程池(CPU密集)
static Scheduler parallel() {
return ParallelScheduler.INSTANCE;
}
}
// 線程切換實現(xiàn)
Flux.just("A")
.publishOn(Schedulers.boundedElastic()) // 切換到彈性線程池
.map(s -> {
System.out.println("map線程: " + Thread.currentThread().getName());
return s.toLowerCase();
})
.subscribeOn(Schedulers.parallel()) // 訂閱在并行線程池
.subscribe();
subscribeOn vs publishOn:
subscribeOn:影響源頭(訂閱操作的線程)
publishOn:影響下游(后續(xù)操作符的線程)
Flux.just("A")
.doOnNext(s -> log("1: " + Thread.currentThread().getName()))
.publishOn(Schedulers.single())
.doOnNext(s -> log("2: " + Thread.currentThread().getName()))
.subscribeOn(Schedulers.parallel())
.subscribe();
輸出:
1: parallel-1 (受subscribeOn影響)
2: single-1 (受publishOn影響)
5.2.4 冷流 vs 熱流
// 冷流(Cold):每次訂閱都重新執(zhí)行
Flux<Integer> cold = Flux.range(1, 3)
.doOnSubscribe(s -> System.out.println("訂閱了"));
cold.subscribe(i -> System.out.println("訂閱者1: " + i));
cold.subscribe(i -> System.out.println("訂閱者2: " + i));
// 輸出:
// 訂閱了
// 訂閱者1: 1
// 訂閱者1: 2
// 訂閱者1: 3
// 訂閱了 (再次執(zhí)行)
// 訂閱者2: 1
// 訂閱者2: 2
// 訂閱者2: 3
// 熱流(Hot):多個訂閱者共享數(shù)據(jù)源
ConnectableFlux<Integer> hot = Flux.range(1, 3)
.doOnSubscribe(s -> System.out.println("訂閱了"))
.publish();
hot.subscribe(i -> System.out.println("訂閱者1: " + i));
hot.subscribe(i -> System.out.println("訂閱者2: " + i));
hot.connect(); // 開始發(fā)射
// 輸出:
// 訂閱了 (只執(zhí)行一次)
// 訂閱者1: 1
// 訂閱者2: 1
// 訂閱者1: 2
// 訂閱者2: 2
// 訂閱者1: 3
// 訂閱者2: 3
六、性能對比與選型
6.1 性能對比
6.1.1 吞吐量測試
測試場景:10000個并發(fā)請求,每個請求返回100個數(shù)據(jù)塊
// 測試代碼
@State(Scope.Benchmark)
public class PerformanceTest {
@Benchmark
public void testSseEmitter(Blackhole blackhole) {
// 模擬SseEmitter
SseEmitter emitter = new SseEmitter();
for (int i = 0; i < 100; i++) {
emitter.send("data" + i);
}
emitter.complete();
}
@Benchmark
public void testFlux(Blackhole blackhole) {
// 模擬Flux
Flux.range(0, 100)
.map(i -> "data" + i)
.subscribe(blackhole::consume);
}
}
測試結(jié)果(JMH Benchmark):
| 指標(biāo) | SseEmitter | Flux | 說明 |
|---|---|---|---|
| 吞吐量 | 8,500 ops/s | 45,000 ops/s | Flux快5倍+ |
| 內(nèi)存占用 | 每連接 ~50KB | 每連接 ~10KB | Flux更節(jié)省 |
| CPU占用 | 65% | 35% | Flux更高效 |
| 延遲 (P99) | 150ms | 30ms | Flux更低 |
6.1.2 內(nèi)存分析
SseEmitter內(nèi)存結(jié)構(gòu): Thread Stack ~1MB (線程棧) Response Buffer ~8KB (HTTP響應(yīng)緩沖) Connection State ~40KB (連接狀態(tài)) 總計:~1.05MB/連接 Flux內(nèi)存結(jié)構(gòu): Subscription ~2KB (訂閱對象) Operator Chain ~5KB (操作符鏈) Buffer ~3KB (可配置) 總計:~10KB/連接 C10K問題: SseEmitter: 10000連接 = 10GB內(nèi)存 Flux: 10000連接 = 100MB內(nèi)存
6.2 技術(shù)選型
6.2.1 選型決策樹
需要流式響應(yīng)?
├─ 否 → 使用普通REST接口
└─ 是 → 繼續(xù)
|
現(xiàn)有項目是Spring MVC?
├─ 是 → 考慮SseEmitter
│ |
│ 并發(fā)量 < 1000?
│ ├─ 是 → SseEmitter (簡單易用)
│ └─ 否 → 考慮遷移到Flux
|
└─ 否(新項目) → Flux (性能更好)
|
團(tuán)隊熟悉響應(yīng)式編程?
├─ 是 → 直接使用Flux
└─ 否 → 先用SseEmitter,逐步遷移
6.2.2 詳細(xì)對比
| 維度 | SseEmitter | Flux | 推薦場景 |
|---|---|---|---|
| 易用性 | ????? 簡單直觀 | ??? 學(xué)習(xí)曲線陡 | 快速開發(fā)選SseEmitter |
| 性能 | ??? 中等 | ????? 優(yōu)秀 | 高并發(fā)選Flux |
| 資源占用 | ?? 每連接1個線程 | ????? 異步非阻塞 | 資源受限選Flux |
| 背壓控制 | ? 不支持 | ? 完整支持 | 需要流控選Flux |
| 生態(tài)集成 | ??? Spring MVC | ????? WebFlux生態(tài) | 全棧響應(yīng)式選Flux |
| 錯誤處理 | ??? 簡單 | ???? 豐富 | 復(fù)雜邏輯選Flux |
| 測試難度 | ???? 容易 | ?? 較難 | 快速驗證選SseEmitter |
| 可維護(hù)性 | ???? 易理解 | ??? 需要經(jīng)驗 | 團(tuán)隊新手選SseEmitter |
6.2.3 實際案例選型
案例1:企業(yè)內(nèi)部管理系統(tǒng)
- 并發(fā):<500
- 團(tuán)隊:傳統(tǒng)Java團(tuán)隊
- 選擇:SseEmitter
- 理由:易于理解和維護(hù),性能夠用
案例2:對外開放的AI聊天API
- 并發(fā):10000+
- 團(tuán)隊:有響應(yīng)式經(jīng)驗
- 選擇:Flux
- 理由:高性能、低資源占用、完整的背壓控制
案例3:實時監(jiān)控大屏
- 并發(fā):<100
- 需求:推送服務(wù)器指標(biāo)
- 選擇:SseEmitter
- 理由:簡單場景,快速實現(xiàn)
案例4:物聯(lián)網(wǎng)數(shù)據(jù)采集平臺
- 并發(fā):50000+設(shè)備
- 數(shù)據(jù):高頻率傳感器數(shù)據(jù)
- 選擇:Flux + R2DBC
- 理由:全棧響應(yīng)式,極致性能
七、生產(chǎn)環(huán)境實戰(zhàn)
7.1 完整項目改造
7.1.1 改造前(同步模式)
@RestController
@RequestMapping("/api/v1/chat")
public class ChatController {
@PostMapping("/send")
public Result<MessageVO> sendMessage(@RequestBody SendMessageReq req) {
// 阻塞等待AI回復(fù)(可能30秒)
String response = aiService.generate(req.getContent());
return Result.success(new MessageVO(response));
}
}
問題:
- 用戶體驗差:等待30秒
- 資源浪費(fèi):占用線程30秒
- 容易超時:Nginx/Gateway 60秒超時
7.1.2 改造后(流式模式)
@RestController
@RequestMapping("/api/v1/chat")
@RequiredArgsConstructor
public class StreamChatController {
private final StreamChatService chatService;
/**
* 流式發(fā)送消息(推薦)
*/
@PostMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<ChatResponse>> streamMessage(
@RequestBody SendMessageReq req) {
return chatService.streamChat(req)
.map(chunk -> ServerSentEvent.<ChatResponse>builder()
.event("message")
.id(chunk.getMessageId())
.data(chunk)
.build())
.concatWith(completionEvent())
.onErrorResume(this::handleError);
}
private Flux<ServerSentEvent<ChatResponse>> completionEvent() {
return Flux.just(ServerSentEvent.<ChatResponse>builder()
.event("complete")
.data(ChatResponse.complete())
.build());
}
private Flux<ServerSentEvent<ChatResponse>> handleError(Throwable e) {
log.error("流式聊天異常", e);
return Flux.just(ServerSentEvent.<ChatResponse>builder()
.event("error")
.data(ChatResponse.error(e.getMessage()))
.build());
}
}
7.2 Service層實現(xiàn)
@Service
@Slf4j
@RequiredArgsConstructor
public class StreamChatService {
private final OpenAIClient openAIClient;
private final ConversationRepository conversationRepository;
private final MessageRepository messageRepository;
private final RedisTemplate<String, Object> redisTemplate;
/**
* 流式聊天核心邏輯
*/
public Flux<ChatResponse> streamChat(SendMessageReq req) {
String messageId = generateMessageId();
String sessionId = req.getSessionId();
return Flux.create(sink -> {
StringBuilder fullResponse = new StringBuilder();
AtomicInteger chunkCount = new AtomicInteger(0);
try {
// 1. 保存用戶消息
saveUserMessage(sessionId, req.getContent())
.subscribe();
// 2. 調(diào)用OpenAI流式API
openAIClient.streamChatCompletion(
buildPrompt(sessionId, req.getContent()),
new StreamCallback() {
@Override
public void onChunk(String chunk) {
fullResponse.append(chunk);
// 發(fā)送數(shù)據(jù)塊
ChatResponse response = ChatResponse.builder()
.messageId(messageId)
.sessionId(sessionId)
.content(chunk)
.chunkIndex(chunkCount.getAndIncrement())
.isComplete(false)
.build();
sink.next(response);
// 緩存到Redis(支持?jǐn)嗑€重連)
cacheChunk(messageId, chunkCount.get(), chunk);
}
@Override
public void onComplete() {
// 3. 保存完整AI回復(fù)
saveAssistantMessage(sessionId, messageId, fullResponse.toString())
.doOnSuccess(msg -> log.info("消息保存成功: {}", messageId))
.doOnError(e -> log.error("消息保存失敗", e))
.subscribe();
// 4. 發(fā)送完成信號
ChatResponse finalResponse = ChatResponse.builder()
.messageId(messageId)
.sessionId(sessionId)
.content("")
.chunkIndex(chunkCount.get())
.isComplete(true)
.totalChunks(chunkCount.get())
.build();
sink.next(finalResponse);
sink.complete();
// 5. 清理緩存
clearCache(messageId);
}
@Override
public void onError(Throwable error) {
log.error("AI生成失敗", error);
sink.error(new BusinessException("AI服務(wù)異常: " + error.getMessage()));
}
}
);
} catch (Exception e) {
sink.error(e);
}
})
.publishOn(Schedulers.boundedElastic())
.timeout(Duration.ofMinutes(5)) // 5分鐘超時
.doOnSubscribe(sub -> log.info("開始流式聊天, sessionId={}", sessionId))
.doOnComplete(() -> log.info("流式聊天完成, sessionId={}", sessionId))
.doOnError(e -> log.error("流式聊天異常, sessionId={}", sessionId, e));
}
/**
* 保存用戶消息(響應(yīng)式)
*/
private Mono<Message> saveUserMessage(String sessionId, String content) {
return Mono.fromCallable(() -> {
Message message = Message.builder()
.messageId(generateMessageId())
.sessionId(sessionId)
.role(MessageRole.USER)
.content(content)
.createTime(LocalDateTime.now())
.build();
return messageRepository.save(message);
}).subscribeOn(Schedulers.boundedElastic());
}
/**
* 保存AI回復(fù)消息(響應(yīng)式)
*/
private Mono<Message> saveAssistantMessage(String sessionId, String messageId, String content) {
return Mono.fromCallable(() -> {
Message message = Message.builder()
.messageId(messageId)
.sessionId(sessionId)
.role(MessageRole.ASSISTANT)
.content(content)
.createTime(LocalDateTime.now())
.build();
return messageRepository.save(message);
}).subscribeOn(Schedulers.boundedElastic());
}
/**
* 緩存數(shù)據(jù)塊(支持?jǐn)嗑€重連)
*/
private void cacheChunk(String messageId, int index, String chunk) {
String key = "chat:stream:" + messageId;
redisTemplate.opsForList().rightPush(key, chunk);
redisTemplate.expire(key, Duration.ofMinutes(10));
}
/**
* 清理緩存
*/
private void clearCache(String messageId) {
String key = "chat:stream:" + messageId;
redisTemplate.delete(key);
}
private String generateMessageId() {
return "msg_" + System.currentTimeMillis() + "_" + RandomUtil.randomString(8);
}
}
7.3 前端對接實現(xiàn)
7.3.1 原生JavaScript
class StreamChatClient {
constructor(apiBaseUrl) {
this.apiBaseUrl = apiBaseUrl;
this.eventSource = null;
}
/**
* 發(fā)送流式消息
*/
sendStreamMessage(sessionId, content, callbacks) {
const url = `${this.apiBaseUrl}/api/v1/chat/stream`;
// 使用fetch進(jìn)行POST請求,獲取ReadableStream
fetch(url, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Accept': 'text/event-stream'
},
body: JSON.stringify({
sessionId: sessionId,
content: content
})
})
.then(response => {
const reader = response.body.getReader();
const decoder = new TextDecoder();
// 讀取流
const readChunk = () => {
reader.read().then(({ done, value }) => {
if (done) {
callbacks.onComplete?.();
return;
}
// 解析SSE數(shù)據(jù)
const chunk = decoder.decode(value);
const lines = chunk.split('\n');
let eventType = 'message';
let data = '';
for (const line of lines) {
if (line.startsWith('event:')) {
eventType = line.substring(6).trim();
} else if (line.startsWith('data:')) {
data = line.substring(5).trim();
} else if (line === '' && data) {
// 完整的事件
this.handleEvent(eventType, data, callbacks);
eventType = 'message';
data = '';
}
}
// 繼續(xù)讀取
readChunk();
});
};
readChunk();
})
.catch(error => {
console.error('Stream error:', error);
callbacks.onError?.(error);
});
}
/**
* 處理SSE事件
*/
handleEvent(eventType, data, callbacks) {
try {
const parsed = JSON.parse(data);
switch (eventType) {
case 'message':
callbacks.onMessage?.(parsed.content);
break;
case 'complete':
callbacks.onComplete?.();
break;
case 'error':
callbacks.onError?.(new Error(parsed.message));
break;
}
} catch (e) {
console.error('Parse error:', e);
}
}
}
// 使用示例
const client = new StreamChatClient('http://localhost:8080');
client.sendStreamMessage('sess_123', '你好,介紹一下你自己', {
onMessage: (content) => {
// 逐字顯示
document.getElementById('response').innerText += content;
},
onComplete: () => {
console.log('流式響應(yīng)完成');
},
onError: (error) => {
console.error('錯誤:', error);
alert('發(fā)生錯誤: ' + error.message);
}
});
7.3.2 React實現(xiàn)
import React, { useState, useEffect, useRef } from 'react';
interface ChatMessage {
role: 'user' | 'assistant';
content: string;
isStreaming?: boolean;
}
export const StreamChat: React.FC = () => {
const [messages, setMessages] = useState<ChatMessage[]>([]);
const [inputValue, setInputValue] = useState('');
const [isLoading, setIsLoading] = useState(false);
const messagesEndRef = useRef<HTMLDivElement>(null);
// 自動滾動到底部
useEffect(() => {
messagesEndRef.current?.scrollIntoView({ behavior: 'smooth' });
}, [messages]);
const sendMessage = async () => {
if (!inputValue.trim() || isLoading) return;
const userMessage = inputValue;
setInputValue('');
setIsLoading(true);
// 添加用戶消息
setMessages(prev => [...prev, { role: 'user', content: userMessage }]);
// 添加AI消息占位符
const aiMessageIndex = messages.length + 1;
setMessages(prev => [...prev, {
role: 'assistant',
content: '',
isStreaming: true
}]);
try {
const response = await fetch('/api/v1/chat/stream', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Accept': 'text/event-stream'
},
body: JSON.stringify({
sessionId: 'sess_123',
content: userMessage
})
});
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() || '';
for (const line of lines) {
if (line.startsWith('data:')) {
const data = line.substring(5).trim();
if (data) {
try {
const parsed = JSON.parse(data);
if (parsed.content) {
// 更新AI消息內(nèi)容
setMessages(prev => {
const newMessages = [...prev];
newMessages[aiMessageIndex] = {
...newMessages[aiMessageIndex],
content: newMessages[aiMessageIndex].content + parsed.content
};
return newMessages;
});
}
if (parsed.isComplete) {
// 標(biāo)記流結(jié)束
setMessages(prev => {
const newMessages = [...prev];
newMessages[aiMessageIndex].isStreaming = false;
return newMessages;
});
}
} catch (e) {
console.error('Parse error:', e);
}
}
}
}
}
} catch (error) {
console.error('Stream error:', error);
alert('發(fā)生錯誤: ' + error);
} finally {
setIsLoading(false);
}
};
return (
<div className="chat-container">
<div className="messages">
{messages.map((msg, index) => (
<div key={index} className={`message message-${msg.role}`}>
<div className="message-content">
{msg.content}
{msg.isStreaming && <span className="cursor">▊</span>}
</div>
</div>
))}
<div ref={messagesEndRef} />
</div>
<div className="input-area">
<input
type="text"
value={inputValue}
onChange={(e) => setInputValue(e.target.value)}
onKeyPress={(e) => e.key === 'Enter' && sendMessage()}
placeholder="輸入消息..."
disabled={isLoading}
/>
<button onClick={sendMessage} disabled={isLoading}>
{isLoading ? '發(fā)送中...' : '發(fā)送'}
</button>
</div>
</div>
);
};
7.4 異常處理與監(jiān)控
7.4.1 超時處理
@Configuration
public class WebFluxConfig {
@Bean
public WebClient webClient() {
return WebClient.builder()
.clientConnector(new ReactorClientHttpConnector(
HttpClient.create()
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000)
.responseTimeout(Duration.ofMinutes(5))
.doOnConnected(conn -> conn
.addHandlerLast(new ReadTimeoutHandler(300))
.addHandlerLast(new WriteTimeoutHandler(300)))
))
.build();
}
}
// Flux超時處理
public Flux<String> streamWithTimeout() {
return Flux.create(sink -> {
// 流式邏輯
})
.timeout(Duration.ofMinutes(5), Flux.just("超時兜底數(shù)據(jù)"))
.onErrorResume(TimeoutException.class, e -> {
log.error("流式處理超時", e);
return Flux.just("處理時間過長,請稍后重試");
});
}
7.4.2 監(jiān)控埋點(diǎn)
@Aspect
@Component
@Slf4j
public class StreamMonitorAspect {
@Around("@annotation(streamMonitor)")
public Object monitor(ProceedingJoinPoint pjp, StreamMonitor streamMonitor) throws Throwable {
String methodName = pjp.getSignature().getName();
long startTime = System.currentTimeMillis();
Object result = pjp.proceed();
if (result instanceof Flux) {
Flux<?> flux = (Flux<?>) result;
AtomicLong chunkCount = new AtomicLong(0);
AtomicLong totalBytes = new AtomicLong(0);
return flux
.doOnNext(item -> {
chunkCount.incrementAndGet();
if (item instanceof String) {
totalBytes.addAndGet(((String) item).length());
}
})
.doOnComplete(() -> {
long cost = System.currentTimeMillis() - startTime;
log.info("流式方法執(zhí)行完成: method={}, chunks={}, bytes={}, cost={}ms",
methodName, chunkCount.get(), totalBytes.get(), cost);
// 上報監(jiān)控指標(biāo)
MetricsCollector.recordStreamMetrics(
methodName,
chunkCount.get(),
totalBytes.get(),
cost
);
})
.doOnError(error -> {
log.error("流式方法執(zhí)行失敗: method={}, error={}",
methodName, error.getMessage(), error);
// 上報錯誤
MetricsCollector.recordStreamError(methodName, error);
});
}
return result;
}
}
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface StreamMonitor {
String value() default "";
}
// 使用
@StreamMonitor("chat-stream")
public Flux<ChatResponse> streamChat(SendMessageReq req) {
// ...
}
八、常見問題與解決方案
8.1 SseEmitter常見問題
Q1: SseEmitter連接意外斷開
現(xiàn)象: 客戶端隨機(jī)斷開連接,無錯誤日志
原因:
- Nginx/Gateway超時配置
- 網(wǎng)絡(luò)不穩(wěn)定
- 瀏覽器Tab切換(移動端)
解決方案:
# Nginx配置
location /api/ {
proxy_pass http://backend;
# SSE必需配置
proxy_set_header Connection '';
proxy_http_version 1.1;
chunked_transfer_encoding off;
proxy_buffering off;
proxy_cache off;
# 超時設(shè)置
proxy_connect_timeout 10s;
proxy_send_timeout 600s; # 10分鐘
proxy_read_timeout 600s; # 10分鐘
}
// 心跳?;?
@Scheduled(fixedRate = 30000) // 每30秒
public void sendHeartbeat() {
emitterManager.getAllEmitters().forEach(emitter -> {
try {
emitter.send(SseEmitter.event()
.name("heartbeat")
.data("ping"));
} catch (IOException e) {
log.warn("發(fā)送心跳失敗,移除連接", e);
emitterManager.remove(emitter);
}
});
}
Q2: 內(nèi)存泄漏
現(xiàn)象: 服務(wù)器內(nèi)存持續(xù)增長,最終OOM
原因: SseEmitter未正確關(guān)閉,導(dǎo)致連接泄漏
解決方案:
@Component
public class SseEmitterManager {
private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
public SseEmitter create(String sessionId, Long timeout) {
SseEmitter emitter = new SseEmitter(timeout);
// 完成時移除
emitter.onCompletion(() -> {
log.info("SSE連接完成: {}", sessionId);
emitters.remove(sessionId);
});
// 超時時移除
emitter.onTimeout(() -> {
log.warn("SSE連接超時: {}", sessionId);
emitters.remove(sessionId);
try {
emitter.complete();
} catch (Exception e) {
log.error("關(guān)閉超時連接失敗", e);
}
});
// 錯誤時移除
emitter.onError(throwable -> {
log.error("SSE連接異常: {}", sessionId, throwable);
emitters.remove(sessionId);
});
emitters.put(sessionId, emitter);
return emitter;
}
// 定期清理過期連接
@Scheduled(fixedRate = 60000) // 每分鐘
public void cleanup() {
long now = System.currentTimeMillis();
emitters.entrySet().removeIf(entry -> {
// 實現(xiàn)自定義清理邏輯
return false;
});
}
}
Q3: 數(shù)據(jù)丟失
現(xiàn)象: 客戶端收到的數(shù)據(jù)不完整
原因: 緩沖區(qū)未及時flush
解決方案:
public void send(String data) throws IOException {
emitter.send(SseEmitter.event()
.data(data));
// SseEmitter內(nèi)部會自動flush,但如果自定義實現(xiàn)需要注意:
// outputStream.write(data.getBytes());
// outputStream.flush(); // 必須!
}
8.2 Flux常見問題
Q1: 背壓溢出
現(xiàn)象: reactor.core.Exceptions$OverflowException: Queue is full
原因: 生產(chǎn)速度 > 消費(fèi)速度,且緩沖區(qū)滿
解決方案:
Flux<String> flux = Flux.create(sink -> {
// 生產(chǎn)邏輯
}, FluxSink.OverflowStrategy.BUFFER) // 或 DROP、LATEST、ERROR
.onBackpressureBuffer(1000, // 緩沖區(qū)大小
dropped -> log.warn("丟棄數(shù)據(jù): {}", dropped));
Q2: 線程阻塞
現(xiàn)象: 響應(yīng)式代碼仍然很慢
原因: 在響應(yīng)式流中使用了阻塞操作
錯誤示例:
Flux.create(sink -> {
String result = blockingHttpClient.get(); // 阻塞!
sink.next(result);
sink.complete();
})
正確做法:
Flux.create(sink -> {
// 使用異步客戶端
asyncHttpClient.get()
.thenAccept(result -> {
sink.next(result);
sink.complete();
})
.exceptionally(error -> {
sink.error(error);
return null;
});
})
// 或使用subscribeOn切換線程
.subscribeOn(Schedulers.boundedElastic())
Q3: 訂閱未觸發(fā)
現(xiàn)象: Flux定義了但沒有執(zhí)行
原因: Flux是惰性的,只有訂閱才會執(zhí)行
錯誤示例:
@GetMapping("/test")
public void test() {
Flux.range(1, 10)
.map(i -> i * 2)
.doOnNext(System.out::println);
// 沒有subscribe,不會執(zhí)行!
}
正確做法:
@GetMapping("/test")
public Flux<Integer> test() {
return Flux.range(1, 10)
.map(i -> i * 2);
// Spring WebFlux會自動訂閱
}
// 或手動訂閱
flux.subscribe(
data -> System.out.println(data),
error -> System.err.println(error),
() -> System.out.println("完成")
);
8.3 生產(chǎn)環(huán)境最佳實踐
1. 設(shè)置合理的超時時間
// 不要設(shè)置為0(永不超時) SseEmitter emitter = new SseEmitter(5 * 60 * 1000L); // 5分鐘 // Flux也要設(shè)置超時 flux.timeout(Duration.ofMinutes(5))
2. 限制并發(fā)連接數(shù)
@Component
public class ConnectionLimiter {
private final Semaphore semaphore = new Semaphore(1000); // 最多1000個連接
public SseEmitter createWithLimit() throws InterruptedException {
if (!semaphore.tryAcquire(5, TimeUnit.SECONDS)) {
throw new BusinessException("服務(wù)器繁忙,請稍后重試");
}
SseEmitter emitter = new SseEmitter();
emitter.onCompletion(() -> semaphore.release());
emitter.onTimeout(() -> semaphore.release());
emitter.onError(e -> semaphore.release());
return emitter;
}
}
3. 日志與監(jiān)控
Flux<String> flux = Flux.create(sink -> {
// 業(yè)務(wù)邏輯
})
.doOnSubscribe(sub -> {
log.info("開始流式處理: {}", contextInfo);
MetricsCollector.incrementActiveStreams();
})
.doOnNext(item -> {
log.debug("發(fā)送數(shù)據(jù)塊: {}", item);
MetricsCollector.recordChunk();
})
.doOnComplete(() -> {
log.info("流式處理完成: {}", contextInfo);
MetricsCollector.decrementActiveStreams();
})
.doOnError(error -> {
log.error("流式處理失敗: {}", contextInfo, error);
MetricsCollector.recordError();
MetricsCollector.decrementActiveStreams();
});
4. 優(yōu)雅關(guān)閉
@Component
public class GracefulShutdown implements ApplicationListener<ContextClosedEvent> {
@Autowired
private SseEmitterManager emitterManager;
@Override
public void onApplicationEvent(ContextClosedEvent event) {
log.info("應(yīng)用關(guān)閉,斷開所有SSE連接");
emitterManager.getAllEmitters().forEach(emitter -> {
try {
emitter.send(SseEmitter.event()
.name("shutdown")
.data("服務(wù)器即將重啟,請重新連接"));
emitter.complete();
} catch (IOException e) {
log.error("發(fā)送關(guān)閉通知失敗", e);
}
});
}
}
總結(jié)
核心要點(diǎn)
SseEmitter:
- 基于Servlet異步,簡單易用
- 適合中小型項目(并發(fā)<1000)
- 每個連接占用一個線程,資源開銷較大
- 不支持背壓控制
Flux:
- 基于Reactive Streams,性能卓越
- 適合高并發(fā)場景(并發(fā)>1000)
- 異步非阻塞,資源利用率高
- 完整的背壓控制和錯誤處理
選型建議:
- 快速上手、團(tuán)隊經(jīng)驗不足 → SseEmitter
- 高性能、大規(guī)模并發(fā) → Flux
- 新項目推薦直接使用Flux
- 老項目可以先用SseEmitter,逐步遷移
未來趨勢
響應(yīng)式編程已成為Java生態(tài)的重要方向,Spring、R2DBC、Kafka、Redis等主流框架都已支持響應(yīng)式。掌握Flux不僅能提升系統(tǒng)性能,也是技術(shù)成長的必經(jīng)之路。
到此這篇關(guān)于Java響應(yīng)式編程之Flux與SseEmitter的文章就介紹到這了,更多相關(guān)Java響應(yīng)式編程Flux與SseEmitter內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
JavaMail實現(xiàn)發(fā)送郵件(QQ郵箱)
這篇文章主要為大家詳細(xì)介紹了JavaMail實現(xiàn)發(fā)送郵件(QQ郵箱),文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下2022-08-08
詳解Spring?Security?捕獲?filter?層面異常返回我們自定義的內(nèi)容
Spring?的異常會轉(zhuǎn)發(fā)到?BasicErrorController?中進(jìn)行異常寫入,然后才會返回客戶端。所以,我們可以在?BasicErrorController?對?filter異常進(jìn)行捕獲并處理,下面通過本文給大家介紹Spring?Security?捕獲?filter?層面異常,返回我們自定義的內(nèi)容,感興趣的朋友一起看看吧2022-05-05
RestTemplate對HttpClient的適配源碼解讀
這篇文章主要為大家介紹了RestTemplate對HttpClient的適配源碼解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-10-10
mybatis plus saveBatch方法方法執(zhí)行慢導(dǎo)致接口發(fā)送慢解決分析
這篇文章主要為大家介紹了mybatis plus saveBatch方法導(dǎo)致接口發(fā)送慢解決分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-10-10

