最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Java響應(yīng)式編程之Flux與SseEmitter深度解析(附詳細(xì)代碼)

 更新時間:2026年01月29日 10:59:44   作者:jiayong23  
這篇文章主要介紹了Java響應(yīng)式編程之Flux與SseEmitter深度解析的相關(guān)資料,從需求分析、技術(shù)背景、使用方法、底層原理、性能對比到生產(chǎn)環(huán)境實戰(zhàn),全面解析了這兩種技術(shù)的特點(diǎn)、應(yīng)用場景及優(yōu)缺點(diǎn),需要的朋友可以參考下

本文深入探討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));
}

存在的問題:

  1. 用戶體驗差:用戶發(fā)送消息后需要等待很長時間才能看到完整回復(fù)
  2. 資源浪費(fèi):一個請求會長時間占用一個線程,降低服務(wù)器并發(fā)能力
  3. 超時風(fēng)險:長時間處理可能觸發(fā)HTTP超時(默認(rèn)30-60秒)
  4. 無法感知進(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)說明:

  1. produces = MediaType.TEXT_EVENT_STREAM_VALUE:必須設(shè)置Content-Type為text/event-stream
  2. SseEmitter(timeout):超時時間,0表示永不超時(不推薦)
  3. CompletableFuture.runAsync():異步執(zhí)行,避免阻塞Tomcat線程
  4. 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)
        )
    )
)

訂閱傳播:

// 簡化的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)SseEmitterFlux說明
吞吐量8,500 ops/s45,000 ops/sFlux快5倍+
內(nèi)存占用每連接 ~50KB每連接 ~10KBFlux更節(jié)省
CPU占用65%35%Flux更高效
延遲 (P99)150ms30msFlux更低

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ì)對比

維度SseEmitterFlux推薦場景
易用性????? 簡單直觀??? 學(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ī)斷開連接,無錯誤日志

原因:

  1. Nginx/Gateway超時配置
  2. 網(wǎng)絡(luò)不穩(wěn)定
  3. 瀏覽器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)

  1. SseEmitter

    • 基于Servlet異步,簡單易用
    • 適合中小型項目(并發(fā)<1000)
    • 每個連接占用一個線程,資源開銷較大
    • 不支持背壓控制
  2. Flux

    • 基于Reactive Streams,性能卓越
    • 適合高并發(fā)場景(并發(fā)>1000)
    • 異步非阻塞,資源利用率高
    • 完整的背壓控制和錯誤處理
  3. 選型建議

    • 快速上手、團(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)文章

  • 一次Jvm old過高的排查過程實戰(zhàn)記錄

    一次Jvm old過高的排查過程實戰(zhàn)記錄

    這篇文章主要給大家介紹了一次Jvm old過高的排查過程,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2018-11-11
  • JVM雙親委派模型知識詳細(xì)總結(jié)

    JVM雙親委派模型知識詳細(xì)總結(jié)

    今天帶各位小伙伴學(xué)習(xí)Java虛擬機(jī)的相關(guān)知識,文中對JVM雙親委派模型作了非常詳細(xì)的介紹,對正在學(xué)習(xí)java的小伙伴們有很好的幫助,需要的朋友可以參考下
    2021-05-05
  • JavaMail實現(xiàn)發(fā)送郵件(QQ郵箱)

    JavaMail實現(xiàn)發(fā)送郵件(QQ郵箱)

    這篇文章主要為大家詳細(xì)介紹了JavaMail實現(xiàn)發(fā)送郵件(QQ郵箱),文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2022-08-08
  • 詳解MyBatis日志如何做到兼容所有常用的日志框架

    詳解MyBatis日志如何做到兼容所有常用的日志框架

    這篇文章主要介紹了詳解MyBatis日志如何做到兼容所有常用的日志框架,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-11-11
  • 詳解Spring?Security?捕獲?filter?層面異常返回我們自定義的內(nèi)容

    詳解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的適配源碼解讀

    這篇文章主要為大家介紹了RestTemplate對HttpClient的適配源碼解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-10-10
  • 一文搞懂Java中對象池的實現(xiàn)

    一文搞懂Java中對象池的實現(xiàn)

    池化并不是什么新鮮的技術(shù),它更像一種軟件設(shè)計模式,主要功能是緩存一組已經(jīng)初始化的對象,以供隨時可以使用。本文將為大家詳細(xì)講講Java中對象池的實現(xiàn),需要的可以參考一下
    2022-07-07
  • mybatis plus saveBatch方法方法執(zhí)行慢導(dǎo)致接口發(fā)送慢解決分析

    mybatis plus saveBatch方法方法執(zhí)行慢導(dǎo)致接口發(fā)送慢解決分析

    這篇文章主要為大家介紹了mybatis plus saveBatch方法導(dǎo)致接口發(fā)送慢解決分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-10-10
  • Java?Web關(guān)鍵字填空示例詳解

    Java?Web關(guān)鍵字填空示例詳解

    最近在工作中使用了java?web,發(fā)現(xiàn)有些難度,下面這篇文章主要給大家介紹了關(guān)于Java?Web關(guān)鍵字填空的相關(guān)資料,文中通過實例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-04-04
  • SpringBoot如何使用TraceId日志鏈路追蹤

    SpringBoot如何使用TraceId日志鏈路追蹤

    文章介紹了如何使用TraceId進(jìn)行日志鏈路追蹤,通過在日志中添加TraceId關(guān)鍵字,可以將同一次業(yè)務(wù)調(diào)用鏈上的日志串起來,本文通過實例代碼給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧
    2025-01-01

最新評論

固始县| 高陵县| 东源县| 济源市| 临沭县| 静乐县| 榕江县| 南康市| 时尚| 鹤山市| 宜宾县| 龙海市| 兴仁县| 永登县| 洪泽县| 河南省| 治县。| 巢湖市| 佛冈县| 武平县| 星子县| 集安市| 屯门区| 当雄县| 樟树市| 交口县| 米林县| 嘉黎县| 鹤壁市| 沽源县| 台中县| 大埔县| 佛学| 本溪| 壶关县| 东莞市| 兴和县| 沂源县| 江孜县| 三门县| 陆川县|