SpringBoot實(shí)現(xiàn)SSE(Server-Sent?Events)的完整指南
引言
在 Spring Boot 中實(shí)現(xiàn) SSE (Server-Sent Events) 非常簡單,SSE 是一種服務(wù)器向客戶端推送事件的技術(shù)。以下是完整的實(shí)現(xiàn)步驟:
基礎(chǔ)實(shí)現(xiàn)
1. 添加依賴
確保 pom.xml 中包含 Spring Web 依賴:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
2. 創(chuàng)建 SSE 控制器
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@RestController
public class SseController {
// 用于異步發(fā)送事件的線程池
private final ExecutorService executor = Executors.newCachedThreadPool();
@GetMapping(path = "/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter handleSse() {
SseEmitter emitter = new SseEmitter(60_000L); // 設(shè)置超時(shí)時(shí)間(毫秒)
// 在單獨(dú)的線程中發(fā)送事件
executor.execute(() -> {
try {
for (int i = 0; i < 10; i++) {
// 構(gòu)建事件
SseEmitter.SseEventBuilder event = SseEmitter.event()
.id(String.valueOf(i)) // 事件ID
.name("message") // 事件名稱
.data("Event #" + i); // 事件數(shù)據(jù)
// 發(fā)送事件
emitter.send(event);
// 模擬延遲
Thread.sleep(1000);
}
// 完成發(fā)送
emitter.complete();
} catch (IOException | InterruptedException e) {
// 發(fā)生錯(cuò)誤時(shí)關(guān)閉連接
emitter.completeWithError(e);
}
});
// 處理完成和超時(shí)事件
emitter.onCompletion(() -> System.out.println("SSE completed"));
emitter.onTimeout(() -> {
System.out.println("SSE timeout");
emitter.complete();
});
return emitter;
}
}
3. 前端監(jiān)聽 SSE
<!DOCTYPE html>
<html>
<head>
<title>SSE Demo</title>
</head>
<body>
<div id="events"></div>
<script>
const eventSource = new EventSource('/sse');
// 監(jiān)聽消息事件
eventSource.onmessage = function(event) {
const data = event.data;
const element = document.createElement('p');
element.textContent = 'Received: ' + data;
document.getElementById('events').appendChild(element);
};
// 監(jiān)聽自定義事件
eventSource.addEventListener('message', function(event) {
console.log('Custom event:', event.data);
});
// 錯(cuò)誤處理
eventSource.onerror = function(error) {
console.error('EventSource error:', error);
eventSource.close();
};
</script>
</body>
</html>
添加鑒權(quán)支持
1. 基于 Token 的鑒權(quán)
@GetMapping(path = "/secure-sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter handleSecureSse(@RequestHeader("Authorization") String token) {
// 驗(yàn)證Token
if (!isValidToken(token)) {
throw new ResponseStatusException(HttpStatus.UNAUTHORIZED, "Invalid token");
}
SseEmitter emitter = new SseEmitter();
// 獲取用戶信息
User user = getUserFromToken(token);
executor.execute(() -> {
try {
// 發(fā)送個(gè)性化事件
emitter.send(SseEmitter.event()
.data("Welcome, " + user.getName())
.name("greeting"));
// 繼續(xù)發(fā)送其他事件...
} catch (IOException e) {
emitter.completeWithError(e);
}
});
return emitter;
}
private boolean isValidToken(String token) {
// 實(shí)現(xiàn)Token驗(yàn)證邏輯
return token != null && token.startsWith("Bearer ");
}
private User getUserFromToken(String token) {
// 從Token中提取用戶信息
return new User("John Doe"); // 示例
}
2. 前端發(fā)送鑒權(quán)信息
const token = "Bearer your_jwt_token_here";
const eventSource = new EventSource('/secure-sse', {
headers: {
Authorization: token
}
});
高級(jí)功能實(shí)現(xiàn)
1. 廣播事件給多個(gè)客戶端
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
@Service
public class SseService {
private final List<SseEmitter> emitters = new CopyOnWriteArrayList<>();
public SseEmitter subscribe() {
SseEmitter emitter = new SseEmitter(60_000L);
emitters.add(emitter);
emitter.onCompletion(() -> emitters.remove(emitter));
emitter.onTimeout(() -> emitters.remove(emitter));
return emitter;
}
public void broadcast(String eventName, Object data) {
for (SseEmitter emitter : emitters) {
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
} catch (IOException e) {
emitter.completeWithError(e);
}
}
}
}
2. 在控制器中使用廣播服務(wù)
@RestController
public class SseController {
private final SseService sseService;
public SseController(SseService sseService) {
this.sseService = sseService;
}
@GetMapping("/subscribe")
public SseEmitter subscribe() {
return sseService.subscribe();
}
@PostMapping("/broadcast")
public ResponseEntity<String> broadcastMessage(@RequestBody String message) {
sseService.broadcast("message", message);
return ResponseEntity.ok("Message broadcasted");
}
}
3. 發(fā)送 JSON 數(shù)據(jù)
emitter.send(SseEmitter.event()
.name("userUpdate")
.data(new User("Alice", "alice@example.com"), MediaType.APPLICATION_JSON));
4. 重連機(jī)制
let eventSource;
function connectSSE() {
eventSource = new EventSource('/sse');
eventSource.onmessage = event => {
console.log('Received:', event.data);
};
eventSource.onerror = () => {
console.log('Connection lost. Reconnecting...');
eventSource.close();
setTimeout(connectSSE, 3000); // 3秒后重連
};
}
connectSSE(); // 初始連接
生產(chǎn)環(huán)境最佳實(shí)踐
1. 配置超時(shí)和心跳
@Bean
public SseEmitter createSseEmitter() {
SseEmitter emitter = new SseEmitter(120_000L); // 2分鐘超時(shí)
// 心跳機(jī)制
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(() -> {
try {
emitter.send(SseEmitter.event()
.name("heartbeat")
.data("ping"));
} catch (IOException e) {
scheduler.shutdown();
}
}, 0, 30, TimeUnit.SECONDS); // 每30秒發(fā)送心跳
return emitter;
}
2. 異常處理
@RestControllerAdvice
public class SseExceptionHandler {
@ExceptionHandler(SseException.class)
public ResponseEntity<String> handleSseException(SseException ex) {
return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR)
.body("SSE Error: " + ex.getMessage());
}
}
3. CORS 配置
@Configuration
public class WebConfig implements WebMvcConfigurer {
@Override
public void addCorsMappings(CorsRegistry registry) {
registry.addMapping("/sse/**")
.allowedOrigins("https://your-frontend.com")
.allowedMethods("GET")
.allowCredentials(true);
}
}
4. 性能優(yōu)化
@Configuration
public class AsyncConfig implements AsyncConfigurer {
@Override
public Executor getAsyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("SSE-Executor-");
executor.initialize();
return executor;
}
}
完整示例:實(shí)時(shí)股票報(bào)價(jià)
后端控制器
@RestController
public class StockController {
private final SseService sseService;
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
public StockController(SseService sseService) {
this.sseService = sseService;
startStockUpdates();
}
@GetMapping("/stocks")
public SseEmitter getStockUpdates() {
return sseService.subscribe();
}
private void startStockUpdates() {
scheduler.scheduleAtFixedRate(() -> {
Map<String, Double> stocks = Map.of(
"AAPL", 150 + Math.random() * 10,
"MSFT", 250 + Math.random() * 15,
"GOOGL", 2800 + Math.random() * 50
);
sseService.broadcast("stockUpdate", stocks);
}, 0, 2, TimeUnit.SECONDS);
}
}
前端實(shí)現(xiàn)
<div id="stock-prices"></div>
<script>
const eventSource = new EventSource('/stocks');
eventSource.addEventListener('stockUpdate', event => {
const stocks = JSON.parse(event.data);
let html = '<h2>Stock Prices</h2><ul>';
for (const [symbol, price] of Object.entries(stocks)) {
html += `<li>${symbol}: $${price.toFixed(2)}</li>`;
}
html += '</ul>';
document.getElementById('stock-prices').innerHTML = html;
});
</script>
部署注意事項(xiàng)
1. 負(fù)載均衡配置
# Nginx 配置
location /sse {
proxy_pass http://backend;
proxy_http_version 1.1;
proxy_set_header Connection '';
proxy_buffering off;
}
2. Spring Boot 配置
# application.properties server.servlet.context-path=/api spring.mvc.async.request-timeout=120000 # 2分鐘超時(shí)
3. 監(jiān)控端點(diǎn)
@Endpoint(id = "sse")
public class SseEndpoint {
private final SseService sseService;
public SseEndpoint(SseService sseService) {
this.sseService = sseService;
}
@ReadOperation
public Map<String, Object> sseMetrics() {
return Map.of(
"activeConnections", sseService.getActiveConnections(),
"lastBroadcast", sseService.getLastBroadcastTime()
);
}
}
最佳實(shí)踐總結(jié)
- 使用專用服務(wù)類:封裝 SSE 邏輯,提高代碼復(fù)用性
- 實(shí)現(xiàn)心跳機(jī)制:防止連接超時(shí)斷開
- 添加鑒權(quán)支持:保護(hù)敏感數(shù)據(jù)
- 優(yōu)雅處理錯(cuò)誤:實(shí)現(xiàn)異常處理和重連機(jī)制
- 監(jiān)控連接狀態(tài):使用 Actuator 端點(diǎn)監(jiān)控 SSE 連接
- 優(yōu)化線程池:合理配置異步處理線程
- 前端重連邏輯:自動(dòng)恢復(fù)斷開連接
通過以上實(shí)現(xiàn),您可以在 Spring Boot 應(yīng)用中輕松創(chuàng)建 SSE 端點(diǎn),實(shí)現(xiàn)服務(wù)器向客戶端的實(shí)時(shí)事件推送,同時(shí)滿足鑒權(quán)需求。
以上就是SpringBoot實(shí)現(xiàn)SSE(Server-Sent Events)完整指南的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot實(shí)現(xiàn)SSE(Server-Sent Events)的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
- Python使用WebSocket和SSE實(shí)現(xiàn)HTTP服務(wù)器消息推送方式
- chatGPT前端流式輸出js實(shí)現(xiàn)三種方法—fetch、SSE、websocket
- Websocket的用法及常見應(yīng)用場景
- Springboot?實(shí)現(xiàn)Server-Sent?Events的項(xiàng)目實(shí)踐
- Spring Boot中使用Server-Sent Events (SSE) 實(shí)現(xiàn)實(shí)時(shí)數(shù)據(jù)推送教程
- C#實(shí)現(xiàn)SSE(Server-Sent Events)服務(wù)端和客戶端的示例代碼
- Spring Boot中使用SSE(Server-Sent Events)實(shí)現(xiàn)聊天功能:替代websocket服務(wù)器推送
相關(guān)文章
借助Maven搭建Hadoop開發(fā)環(huán)境的最詳細(xì)教程分享
在Maven插件的幫助下,VSCode寫Java其實(shí)非常方便,所以本文就來和大家詳細(xì)講講如何借助maven用VScode搭建Hadoop開發(fā)環(huán)境,需要的可以參考下2023-05-05
在Java的Struts框架下進(jìn)行web編程的入門教程
這篇文章主要介紹了在Java的Struts框架下進(jìn)行web編程的入門教程,需要的朋友可以參考下2015-11-11
快速解決 MyBatis-Plus 中 ID 自增問題(推薦)
本文介紹了MyBatis-Plus中自動(dòng)生成ID過長導(dǎo)致的問題及解決方法,結(jié)合示例代碼給大家介紹的非常詳細(xì),感興趣的朋友一起看看吧2025-02-02
springboot整合kaptcha驗(yàn)證碼的示例代碼
kaptcha是一個(gè)很有用的驗(yàn)證碼生成工具,本篇文章主要介紹了springboot整合kaptcha驗(yàn)證碼的示例代碼,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧2018-06-06
Java實(shí)現(xiàn)調(diào)用ElasticSearch?API的示例詳解
這篇文章主要為大家詳細(xì)介紹了Java調(diào)用ElasticSearch?API的效果資料,文中的示例代碼講解詳細(xì),具有一定的參考價(jià)值,感興趣的可以了解一下2023-03-03
SpringBoot使用minio進(jìn)行文件管理的流程步驟
MinIO 是一個(gè)高性能的對(duì)象存儲(chǔ)系統(tǒng),兼容 Amazon S3 API,該軟件設(shè)計(jì)用于處理非結(jié)構(gòu)化數(shù)據(jù),如圖片、視頻、日志文件以及備份數(shù)據(jù)等,本文給大家介紹了SpringBoot使用minio進(jìn)行文件管理的流程步驟,需要的朋友可以參考下2025-01-01
IntelliJ IDEA安裝目錄和設(shè)置目錄的說明(IntelliJ IDEA快速入門)
這篇文章主要介紹了IntelliJ IDEA安裝目錄和設(shè)置目錄的說明(IntelliJ IDEA快速入門),本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-04-04
SpringBoot實(shí)現(xiàn)健康檢查的完整指南
本文介紹了Spring Boot實(shí)現(xiàn)健康檢查的方法,包括使用Actuator進(jìn)行應(yīng)用健康檢查,自定義健康檢查項(xiàng)目,搭建可視化監(jiān)控大屏和配置告警系統(tǒng),此外,還提供了實(shí)戰(zhàn)案例和避坑指南,幫助讀者更好地理解和應(yīng)用Spring Boot健康檢查功能,需要的朋友可以參考下2025-12-12

