Server-Sent Events (SSE)
Giriş — Sunucudan Tek Yönlü Push
Bir e-ticaret dashboard'u düşün. Sipariş geldiğinde ekranda anlık bildirim görmek istiyorsun. Kullanıcı dashboard'dan sunucuya sürekli veri göndermez — sadece sunucudan gelen güncellemeleri izler. Bu senaryoda WebSocket overkill. İki yönlü iletişime ihtiyacın yok. İşte SSE tam burada devreye girer.
Server-Sent Events (SSE), sunucudan client'a tek yönlü, sürekli veri akışı sağlayan bir HTTP standardıdır. WebSocket'in kardeşi gibi düşün ama daha basit, daha hafif ve HTTP üzerinde çalışır — extra protokol yok. Bu derste SSE'nin ne olduğunu, Spring Boot'ta nasıl kullanıldığını ve gerçek dünya örneklerini göreceğiz.
1. SSE vs WebSocket: Hangisini Seçmeli?
Karşılaştırma Tablosu
| Özellik | SSE | WebSocket |
|---|---|---|
| İletişim yönü | Tek yön (server → client) | Çift yön (full-duplex) |
| Protokol | HTTP | ws:// / wss:// |
| Browser desteği | Tüm modern tarayıcılar | Tüm modern tarayıcılar |
| Otomatik reconnect | ✅ Built-in | ❌ Manuel implementasyon |
| HTTP/2 uyumu | ✅ Mükemmel (multiplexing) | ❌ Ayrı bağlantı |
| Firewall/Proxy | ✅ Normal HTTP, sorun çıkmaz | ⚠️ Upgrade engellenebilir |
| Data formatı | Text (UTF-8) | Text + Binary |
| Max connection | 6 per domain (HTTP/1.1) | Limit yok |
| Complexity | Düşük | Orta-Yüksek |
| CORS | Standard HTTP CORS | Handshake'te CORS |
Ne Zaman SSE?
Sunucudan client'a tek yönlü push yeterli
Dashboard'lar, bildirimler, canlı feed'ler
Borsa fiyatları, hava durumu güncellemeleri
İşlem durumu takibi (deployment, dosya işleme)
Log streaming
HTTP/2 kullanıyorsan (connection limit sorunu çözülür)
Ne Zaman WebSocket?
İki yönlü iletişim gerekli (chat, oyun)
Binary data transfer gerekli
Çok düşük latency kritik
Client da sık sık sunucuya veri gönderiyor
Karar Ağacı
Client sunucuya veri gönderiyor mu?
├── Evet, sık sık → WebSocket
├── Evet, nadiren → SSE + REST combo (push SSE, gönderim REST)
└── Hayır → SSE ✓
Binary data var mı?
├── Evet → WebSocket
└── Hayır → SSE ✓
HTTP/2 kullanılıyor mu?
├── Evet → SSE mükemmel (multiplexing sayesinde connection sınırı yok)
└── Hayır → SSE hâlâ iyi ama 6 connection limiti var💡 İpucu: Birçok gerçek dünya senaryosunda SSE + REST kombinasyonu WebSocket'ten daha pratiktir. SSE ile güncellemeleri al, REST ile aksiyon gönder. İki ayrı concern, iki ayrı mekanizma. Test etmesi, debug etmesi daha kolay.
2. SSE'nin Avantajları
HTTP Uyumluluğu
SSE standart bir HTTP response'tur. Content-Type: text/event-stream. Bu şu demek:
Mevcut HTTP altyapın (load balancer, CDN, proxy) olduğu gibi çalışır
Ayrı bir port veya protokol gerekmez
HTTP authentication, cookie, header'lar sorunsuz çalışır
CORS standart HTTP CORS kurallarıyla yönetilir
HTTP/2 Multiplexing
HTTP/1.1'de tarayıcı bir domain'e max 6 connection açar. 6 SSE bağlantısı açarsan, o domain'e başka HTTP isteği atamazsın. HTTP/2'de bu sorun yok — tek bir TCP bağlantısı üzerinden yüzlerce stream multiplexing yapılır.
Otomatik Reconnect
SSE'nin en güzel özelliği: bağlantı koparsa tarayıcı otomatik olarak yeniden bağlanır. EventSource API'si bunu built-in destekler. WebSocket'te bunu kendin yazarsın.
Firewall-Friendly
Kurumsal ağlarda WebSocket upgrade'i engellenebilir. SSE? Normal HTTP. Kimse engellemez.
3. Spring Boot'ta SseEmitter
Temel Kullanım
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@RestController
@RequestMapping("/api/events")
public class SseController {
@GetMapping("/stream")
public SseEmitter streamEvents() {
// 30 dakika timeout (milisaniye)
SseEmitter emitter = new SseEmitter(30 * 60 * 1000L);
// Asenkron olarak event gönder
CompletableFuture.runAsync(() -> {
try {
for (int i = 1; i <= 10; i++) {
// Event gönder
emitter.send(
SseEmitter.event()
.id(String.valueOf(i)) // Event ID
.name("update") // Event type
.data("Güncelleme #" + i) // Data
.reconnectTime(5000) // Reconnect süresi (ms)
);
Thread.sleep(2000); // 2 saniyede bir
}
emitter.complete(); // Stream bitti
} catch (IOException | InterruptedException e) {
emitter.completeWithError(e);
}
});
return emitter;
}
}SSE Event Formatı
Sunucudan gönderilen event şu formatta olur:
id:1
event:update
retry:5000
data:Güncelleme #1
id:2
event:update
data:Güncelleme #2
Her event boş bir satırla (\n\n) ayrılır. Alanlar:
| Alan | Açıklama | Zorunlu? |
|---|---|---|
data | Event verisi (string veya JSON) | Evet |
id | Event kimliği (reconnect'te kullanılır) | Hayır |
event | Event tipi (client'ta hangi listener tetiklenir) | Hayır |
retry | Reconnect süresi (milisaniye) | Hayır |
JSON Data Gönderme
@GetMapping("/orders")
public SseEmitter streamOrders() {
SseEmitter emitter = new SseEmitter(0L); // 0 = timeout yok
CompletableFuture.runAsync(() -> {
try {
// Her yeni sipariş geldiğinde
OrderEvent event = new OrderEvent("ORD-123", "CREATED",
"Yeni sipariş: Laptop x1");
emitter.send(
SseEmitter.event()
.id("ord-123")
.name("new-order")
.data(event, MediaType.APPLICATION_JSON) // JSON olarak serialize
);
} catch (IOException e) {
emitter.completeWithError(e);
}
});
return emitter;
}// OrderEvent model
public class OrderEvent {
private String orderId;
private String status;
private String description;
private LocalDateTime timestamp;
public OrderEvent(String orderId, String status, String description) {
this.orderId = orderId;
this.status = status;
this.description = description;
this.timestamp = LocalDateTime.now();
}
// Getter'lar
public String getOrderId() { return orderId; }
public String getStatus() { return status; }
public String getDescription() { return description; }
public LocalDateTime getTimestamp() { return timestamp; }
}4. Dashboard Real-Time Güncelleme Örneği
Gerçek dünyada SSE'nin en yaygın kullanımı: dashboard'da canlı veri gösterimi. Sipariş durumu, stok takibi, sistem metrikleri...
Sipariş Takip Dashboard'u
@RestController
@RequestMapping("/api/dashboard")
public class DashboardController {
// Tüm aktif emitter'ları tut
private final List<SseEmitter> emitters = new CopyOnWriteArrayList<>();
@GetMapping("/stream")
public SseEmitter streamDashboard() {
SseEmitter emitter = new SseEmitter(0L); // Timeout yok
// Callback'leri ayarla
emitter.onCompletion(() -> {
emitters.remove(emitter);
System.out.println("Client ayrıldı. Aktif: " + emitters.size());
});
emitter.onTimeout(() -> {
emitter.complete();
emitters.remove(emitter);
System.out.println("Timeout! Aktif: " + emitters.size());
});
emitter.onError(ex -> {
emitters.remove(emitter);
System.err.println("Hata: " + ex.getMessage());
});
emitters.add(emitter);
System.out.println("Yeni client bağlandı. Aktif: " + emitters.size());
// İlk bağlantıda mevcut durumu gönder
try {
emitter.send(SseEmitter.event()
.name("init")
.data(getCurrentDashboardState()));
} catch (IOException e) {
emitter.completeWithError(e);
}
return emitter;
}
// Tüm client'lara event gönder
public void broadcastEvent(String eventType, Object data) {
List<SseEmitter> deadEmitters = new ArrayList<>();
emitters.forEach(emitter -> {
try {
emitter.send(SseEmitter.event()
.name(eventType)
.data(data, MediaType.APPLICATION_JSON));
} catch (IOException e) {
deadEmitters.add(emitter);
}
});
// Başarısız emitter'ları temizle
emitters.removeAll(deadEmitters);
}
private DashboardState getCurrentDashboardState() {
return new DashboardState(
42, // Bekleyen sipariş
156, // Bugünkü sipariş
12450 // Bugünkü ciro
);
}
}Sipariş Servisi — Event Tetikleme
@Service
public class OrderService {
private final DashboardController dashboardController;
public OrderService(DashboardController dashboardController) {
this.dashboardController = dashboardController;
}
public Order createOrder(OrderRequest request) {
Order order = new Order();
order.setId(UUID.randomUUID().toString());
order.setStatus("CREATED");
order.setProduct(request.getProduct());
order.setAmount(request.getAmount());
// ... veritabanına kaydet
// Dashboard'a bildir
dashboardController.broadcastEvent("new-order",
Map.of(
"orderId", order.getId(),
"product", order.getProduct(),
"amount", order.getAmount(),
"status", "CREATED"
));
return order;
}
public void updateOrderStatus(String orderId, String newStatus) {
// ... veritabanında güncelle
dashboardController.broadcastEvent("order-status",
Map.of(
"orderId", orderId,
"status", newStatus,
"updatedAt", LocalDateTime.now().toString()
));
}
}Stok Takip Örneği
@Component
public class StockMonitor {
private final DashboardController dashboardController;
public StockMonitor(DashboardController dashboardController) {
this.dashboardController = dashboardController;
}
// Her 10 saniyede stok durumunu kontrol et
@Scheduled(fixedRate = 10_000)
public void checkStockLevels() {
List<StockAlert> alerts = findLowStockProducts();
if (!alerts.isEmpty()) {
dashboardController.broadcastEvent("stock-alert", alerts);
}
}
private List<StockAlert> findLowStockProducts() {
// Veritabanından düşük stoklu ürünleri getir
// Örnek:
return List.of(
new StockAlert("SKU-001", "Laptop", 3, 10), // Stok: 3, Min: 10
new StockAlert("SKU-045", "Mouse", 5, 20) // Stok: 5, Min: 20
);
}
}⚠️ Dikkat:
broadcastEvent()metodundaIOExceptionyakalayıp dead emitter'ları temizlemeyi unutma. Client bağlantıyı kapattığındaemitter.send()exception fırlatır. Bu temizlenmezse emitter listesi şişer ve memory leak oluşur.
5. Timeout ve Hata Yönetimi
SseEmitter Lifecycle
@GetMapping("/stream")
public SseEmitter createStream() {
// Timeout: 5 dakika (milisaniye)
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);
// Timeout: 0 = sınırsız (dikkatli kullan!)
// SseEmitter emitter = new SseEmitter(0L);
// Stream normal şekilde tamamlandığında
emitter.onCompletion(() -> {
System.out.println("Stream tamamlandı");
cleanup(emitter);
});
// Timeout olduğunda
emitter.onTimeout(() -> {
System.out.println("Stream timeout oldu!");
cleanup(emitter);
// Timeout olduğunda emitter otomatik olarak complete() edilir
});
// Herhangi bir hata olduğunda
emitter.onError(throwable -> {
System.err.println("Stream hatası: " + throwable.getMessage());
cleanup(emitter);
});
return emitter;
}Timeout Stratejileri
// Strateji 1: Kısa timeout, client reconnect
// SSE'nin built-in reconnect'i ile güvenli
SseEmitter emitter = new SseEmitter(60_000L); // 1 dakika
// Strateji 2: Uzun timeout, periyodik heartbeat
SseEmitter emitter = new SseEmitter(30 * 60_000L); // 30 dakika
// + Heartbeat gönder (aşağıda)
// Strateji 3: Sınırsız timeout (önerilmez)
SseEmitter emitter = new SseEmitter(0L);Heartbeat ile Bağlantıyı Canlı Tutma
Proxy'ler ve load balancer'lar idle bağlantıları kapatabilir. Periyodik heartbeat göndererek bunu önle:
@Component
public class SseHeartbeat {
private final ScheduledExecutorService scheduler =
Executors.newSingleThreadScheduledExecutor();
private final List<SseEmitter> emitters = new CopyOnWriteArrayList<>();
@PostConstruct
public void startHeartbeat() {
scheduler.scheduleAtFixedRate(() -> {
List<SseEmitter> dead = new ArrayList<>();
emitters.forEach(emitter -> {
try {
emitter.send(SseEmitter.event()
.name("heartbeat")
.data("ping"));
} catch (IOException e) {
dead.add(emitter);
}
});
emitters.removeAll(dead);
}, 0, 15, TimeUnit.SECONDS); // Her 15 saniyede
}
@PreDestroy
public void shutdown() {
scheduler.shutdown();
}
public void addEmitter(SseEmitter emitter) {
emitters.add(emitter);
emitter.onCompletion(() -> emitters.remove(emitter));
emitter.onTimeout(() -> emitters.remove(emitter));
emitter.onError(e -> emitters.remove(emitter));
}
}💡 İpucu: Heartbeat event'inin
namealanını "heartbeat" yap. Client tarafında bu event'i özel olarak dinle ve işleme — kullanıcıya gösterme. Sadece bağlantının canlı olduğunu doğrulamak için kullan.
6. Multiple Client Handling
Thread-Safe Emitter Yönetimi
@Service
public class SseEmitterService {
// CopyOnWriteArrayList — okuma ağırlıklı senaryolarda performanslı
private final List<SseEmitter> emitters = new CopyOnWriteArrayList<>();
// Kullanıcı bazlı emitter'lar
private final Map<String, List<SseEmitter>> userEmitters = new ConcurrentHashMap<>();
public SseEmitter subscribe() {
SseEmitter emitter = new SseEmitter(30 * 60_000L);
emitter.onCompletion(() -> emitters.remove(emitter));
emitter.onTimeout(() -> emitters.remove(emitter));
emitter.onError(e -> emitters.remove(emitter));
emitters.add(emitter);
return emitter;
}
public SseEmitter subscribeUser(String userId) {
SseEmitter emitter = new SseEmitter(30 * 60_000L);
Runnable cleanup = () -> {
List<SseEmitter> userList = userEmitters.get(userId);
if (userList != null) {
userList.remove(emitter);
if (userList.isEmpty()) {
userEmitters.remove(userId);
}
}
};
emitter.onCompletion(cleanup);
emitter.onTimeout(cleanup);
emitter.onError(e -> cleanup.run());
userEmitters.computeIfAbsent(userId, k -> new CopyOnWriteArrayList<>())
.add(emitter);
return emitter;
}
// Herkese gönder
public void broadcast(String eventType, Object data) {
send(emitters, eventType, data);
}
// Belirli kullanıcıya gönder
public void sendToUser(String userId, String eventType, Object data) {
List<SseEmitter> userList = userEmitters.get(userId);
if (userList != null) {
send(userList, eventType, data);
}
}
private void send(List<SseEmitter> targetEmitters, String eventType, Object data) {
List<SseEmitter> dead = new ArrayList<>();
targetEmitters.forEach(emitter -> {
try {
emitter.send(SseEmitter.event()
.name(eventType)
.data(data, MediaType.APPLICATION_JSON));
} catch (IOException e) {
dead.add(emitter);
}
});
targetEmitters.removeAll(dead);
}
public int getActiveConnectionCount() {
return emitters.size();
}
}Controller ile Kullanım
@RestController
@RequestMapping("/api/sse")
public class SseStreamController {
private final SseEmitterService sseService;
public SseStreamController(SseEmitterService sseService) {
this.sseService = sseService;
}
// Genel stream
@GetMapping("/stream")
public SseEmitter subscribe() {
return sseService.subscribe();
}
// Kullanıcıya özel stream
@GetMapping("/stream/{userId}")
public SseEmitter subscribeUser(@PathVariable String userId) {
return sseService.subscribeUser(userId);
}
// REST endpoint'ten event tetikleme
@PostMapping("/notify")
public ResponseEntity<Void> notifyAll(@RequestBody NotificationPayload payload) {
sseService.broadcast("notification", payload);
return ResponseEntity.ok().build();
}
@PostMapping("/notify/{userId}")
public ResponseEntity<Void> notifyUser(@PathVariable String userId,
@RequestBody NotificationPayload payload) {
sseService.sendToUser(userId, "notification", payload);
return ResponseEntity.ok().build();
}
@GetMapping("/stats")
public ResponseEntity<Map<String, Integer>> getStats() {
return ResponseEntity.ok(
Map.of("activeConnections", sseService.getActiveConnectionCount())
);
}
}⚠️ Dikkat:
CopyOnWriteArrayListher yazma işleminde tüm listeyi kopyalar. Az yazma, çok okuma senaryolarında idealdir (emitter listesi: nadiren ekleme/çıkarma, sık sık iterate). Ama çok sık subscribe/unsubscribe olan senaryolardaCollections.synchronizedList()veyaConcurrentLinkedDequedaha performanslı olabilir.
7. SSE + Spring WebFlux: Reactive SSE
Spring WebFlux ile SSE çok daha elegant. SseEmitter yerine Flux<ServerSentEvent> kullanırsın — reactive stream'ler back-pressure destekler ve thread management otomatik.
Bağımlılık
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>Reactive SSE Endpoint
import org.springframework.http.codec.ServerSentEvent;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;
@RestController
@RequestMapping("/api/reactive-sse")
public class ReactiveSseController {
// Hot source — birden fazla subscriber'a yayın yapabilir
private final Sinks.Many<ServerSentEvent<String>> sink =
Sinks.many().multicast().onBackpressureBuffer();
@GetMapping("/stream")
public Flux<ServerSentEvent<String>> stream() {
return sink.asFlux();
}
// Event publish etmek için
@PostMapping("/publish")
public ResponseEntity<Void> publish(@RequestBody String message) {
ServerSentEvent<String> event = ServerSentEvent.<String>builder()
.id(UUID.randomUUID().toString())
.event("message")
.data(message)
.build();
sink.tryEmitNext(event);
return ResponseEntity.ok().build();
}
}Periyodik Event Stream
@GetMapping("/metrics")
public Flux<ServerSentEvent<MetricData>> streamMetrics() {
return Flux.interval(Duration.ofSeconds(5)) // Her 5 saniyede
.map(seq -> {
MetricData data = new MetricData(
Runtime.getRuntime().freeMemory(),
Runtime.getRuntime().totalMemory(),
Thread.activeCount()
);
return ServerSentEvent.<MetricData>builder()
.id(String.valueOf(seq))
.event("metric")
.data(data)
.build();
});
}Heartbeat ile Birleşik Stream
@GetMapping("/combined-stream")
public Flux<ServerSentEvent<String>> combinedStream() {
// Ana veri stream'i
Flux<ServerSentEvent<String>> dataStream = sink.asFlux();
// Heartbeat stream'i — her 15 saniyede
Flux<ServerSentEvent<String>> heartbeat = Flux.interval(Duration.ofSeconds(15))
.map(seq -> ServerSentEvent.<String>builder()
.event("heartbeat")
.data("ping")
.build());
// İkisini birleştir
return Flux.merge(dataStream, heartbeat);
}Database Change Stream
Veritabanı değişikliklerini real-time olarak stream etme (R2DBC ile):
@GetMapping("/orders/stream")
public Flux<ServerSentEvent<Order>> streamNewOrders() {
return Flux.interval(Duration.ofSeconds(2))
.flatMap(tick -> orderRepository.findRecentOrders(
LocalDateTime.now().minusSeconds(2)))
.distinct(Order::getId) // Aynı siparişi tekrar gönderme
.map(order -> ServerSentEvent.<Order>builder()
.id(order.getId())
.event("new-order")
.data(order)
.build());
}💡 İpucu: Reactive SSE'de
Sinks.many().multicast()kullandığında, hiç subscriber yokken emit edilen event'ler kaybolur. Eğer replay istiyorsanSinks.many().replay().limit(10)kullan — son 10 event'i yeni subscriber'lara tekrar gönderir.
SseEmitter vs Flux<ServerSentEvent> Karşılaştırma
| Özellik | SseEmitter (Servlet) | Flux<SSE> (WebFlux) |
|---|---|---|
| Stack | Spring MVC (Servlet) | Spring WebFlux (Reactive) |
| Thread model | Thread-per-request | Non-blocking, event loop |
| Back-pressure | Yok | ✅ Built-in |
| Emitter yönetimi | Manuel (List + cleanup) | Otomatik (Flux lifecycle) |
| Complexity | Düşük | Orta |
| Performance | İyi | Çok iyi (yüksek concurrency) |
| Learning curve | Düşük | Orta-Yüksek (reactive programming) |
8. Browser EventSource API
Temel Kullanım
// EventSource oluştur — otomatik bağlanır
const eventSource = new EventSource('/api/sse/stream');
// Genel message event'i (event tipi belirtilmemiş mesajlar)
eventSource.onmessage = function(event) {
console.log('Mesaj:', event.data);
console.log('ID:', event.lastEventId);
};
// Named event'ler (server'dan event tipi belirtilmiş)
eventSource.addEventListener('new-order', function(event) {
const order = JSON.parse(event.data);
console.log('Yeni sipariş:', order);
addOrderToTable(order);
});
eventSource.addEventListener('stock-alert', function(event) {
const alert = JSON.parse(event.data);
console.log('Stok uyarısı:', alert);
showStockWarning(alert);
});
eventSource.addEventListener('heartbeat', function(event) {
// Heartbeat — sadece bağlantıyı canlı tutar, UI güncelleme yok
console.debug('Heartbeat received');
});
// Bağlantı durumu
eventSource.onopen = function() {
console.log('SSE bağlantısı kuruldu');
updateConnectionStatus('connected');
};
eventSource.onerror = function(event) {
if (eventSource.readyState === EventSource.CONNECTING) {
console.log('Yeniden bağlanıyor...');
updateConnectionStatus('reconnecting');
} else if (eventSource.readyState === EventSource.CLOSED) {
console.error('Bağlantı kalıcı olarak kapandı');
updateConnectionStatus('disconnected');
}
};
// Bağlantıyı kapat
function stopListening() {
eventSource.close();
console.log('SSE bağlantısı kapatıldı');
}Authentication ile EventSource
EventSource API'si custom header gönderemez! Bu büyük bir limitasyon. Çözümler:
// Çözüm 1: Query parameter ile token (güvenlik riski — URL log'lara düşer)
const eventSource = new EventSource('/api/sse/stream?token=' + jwtToken);
// Çözüm 2: Cookie-based authentication (önerilen)
// HttpOnly cookie otomatik gönderilir
const eventSource = new EventSource('/api/sse/stream', {
withCredentials: true // Cookie'leri gönder
});
// Çözüm 3: EventSource polyfill kullan (custom header desteği)
// npm install eventsource-polyfill
const EventSourcePolyfill = require('eventsource-polyfill');
const eventSource = new EventSourcePolyfill('/api/sse/stream', {
headers: {
'Authorization': 'Bearer ' + jwtToken
}
});
// Çözüm 4: fetch API ile manual SSE parsing
async function connectSSE() {
const response = await fetch('/api/sse/stream', {
headers: {
'Authorization': 'Bearer ' + jwtToken,
'Accept': 'text/event-stream'
}
});
const reader = response.body.getReader();
const decoder = new TextDecoder();
while (true) {
const { done, value } = await reader.read();
if (done) break;
const text = decoder.decode(value);
// SSE format parse...
parseSSEText(text);
}
}⚠️ Dikkat:
EventSourcecustom header desteği yoktur. JWT authentication için query parameter kullanmak güvenlik riski taşır (URL log'lara, browser history'ye düşer). Production'da cookie-based auth veya fetch API ile manual SSE parsing tercih edin.
ReadyState Değerleri
EventSource.CONNECTING = 0; // Bağlanıyor veya yeniden bağlanıyor
EventSource.OPEN = 1; // Bağlı, event alıyor
EventSource.CLOSED = 2; // Kalıcı olarak kapandı9. Reconnect Stratejisi: Last-Event-ID
SSE'nin en güçlü özelliklerinden biri: bağlantı koptuğunda tarayıcı otomatik olarak yeniden bağlanır ve son alınan event'in ID'sini Last-Event-ID header'ında gönderir. Sunucu bu ID'den itibaren event'leri tekrar gönderebilir.
Server Tarafı
@GetMapping("/stream")
public SseEmitter streamWithRecovery(
@RequestHeader(value = "Last-Event-ID", required = false) String lastEventId) {
SseEmitter emitter = new SseEmitter(0L);
CompletableFuture.runAsync(() -> {
try {
// Eğer lastEventId varsa, kaçırılan event'leri gönder
if (lastEventId != null) {
List<StoredEvent> missedEvents =
eventStore.getEventsAfter(lastEventId);
for (StoredEvent event : missedEvents) {
emitter.send(SseEmitter.event()
.id(event.getId())
.name(event.getType())
.data(event.getData()));
}
}
// Sonra normal akışa devam et...
} catch (IOException e) {
emitter.completeWithError(e);
}
});
return emitter;
}Event Store
@Service
public class EventStore {
// Basit in-memory store — production'da Redis veya DB kullan
private final Deque<StoredEvent> events = new ConcurrentLinkedDeque<>();
private static final int MAX_EVENTS = 1000;
public void store(StoredEvent event) {
events.addLast(event);
// Eski event'leri temizle
while (events.size() > MAX_EVENTS) {
events.pollFirst();
}
}
public List<StoredEvent> getEventsAfter(String eventId) {
boolean found = false;
List<StoredEvent> result = new ArrayList<>();
for (StoredEvent event : events) {
if (found) {
result.add(event);
}
if (event.getId().equals(eventId)) {
found = true;
}
}
return result;
}
}Client Tarafı — Reconnect Kontrolü
// EventSource otomatik reconnect yapar.
// Default retry süresi genellikle 3 saniye.
// Server "retry:" field'ı ile bu süreyi değiştirebilir.
// Server'dan:
// retry:5000 → 5 saniye sonra tekrar dene
// id:evt-42 → Son alınan event ID'si
// Tarayıcı otomatik olarak şunu yapar:
// 1. Bağlantı kopar
// 2. 5 saniye bekle (retry değeri)
// 3. GET /api/sse/stream + Header: Last-Event-ID: evt-42
// 4. Server evt-42'den sonraki event'leri gönderirReconnect'i Devre Dışı Bırakma
Bazen otomatik reconnect istemezsin (stream tamamlandı, kullanıcı çıkış yaptı):
// Server tarafı: boş data gönderip stream'i kapat
emitter.send(SseEmitter.event()
.name("close")
.data("Stream tamamlandı"));
emitter.complete();// Client tarafı
eventSource.addEventListener('close', function(event) {
eventSource.close(); // Otomatik reconnect'i durdur
console.log('Stream tamamlandı:', event.data);
});💡 İpucu:
Last-Event-IDmekanizması sadece server event'lereidfield'ı verdiğinde çalışır. ID vermezsen tarayıcı reconnect yapar ama server kaçırılan event'leri bilemez. Production'da her event'e mutlaka ID ver!
10. Tam Örnek: Real-Time Dashboard
Backend (Spring MVC)
@RestController
@RequestMapping("/api/dashboard")
public class DashboardSseController {
private final SseEmitterService sseService;
private final AtomicLong eventCounter = new AtomicLong(0);
public DashboardSseController(SseEmitterService sseService) {
this.sseService = sseService;
}
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter subscribe(
@RequestHeader(value = "Last-Event-ID", required = false) String lastEventId) {
SseEmitter emitter = sseService.subscribe();
// İlk bağlantıda dashboard state gönder
try {
emitter.send(SseEmitter.event()
.id(String.valueOf(eventCounter.incrementAndGet()))
.name("dashboard-init")
.data(getDashboardData(), MediaType.APPLICATION_JSON));
} catch (IOException e) {
emitter.completeWithError(e);
}
return emitter;
}
// Simülasyon: Her 5 saniyede rastgele metrik
@Scheduled(fixedRate = 5000)
public void pushMetrics() {
Map<String, Object> metrics = Map.of(
"cpu", Math.random() * 100,
"memory", Math.random() * 100,
"requests", (int)(Math.random() * 1000),
"errors", (int)(Math.random() * 10),
"timestamp", System.currentTimeMillis()
);
sseService.broadcast("metrics", metrics);
}
private Map<String, Object> getDashboardData() {
return Map.of(
"totalOrders", 1256,
"pendingOrders", 42,
"revenue", 125600,
"activeUsers", sseService.getActiveConnectionCount()
);
}
}Frontend Dashboard
<!DOCTYPE html>
<html lang="tr">
<head>
<meta charset="UTF-8">
<title>Real-Time Dashboard</title>
<style>
body {
font-family: 'Segoe UI', sans-serif;
background: #0f0f23; color: #ccc; margin: 0; padding: 20px;
}
h1 { color: #00d4aa; }
.status {
display: inline-block; padding: 4px 12px; border-radius: 20px;
font-size: 0.85em; margin-left: 10px;
}
.status.connected { background: #27ae60; color: white; }
.status.disconnected { background: #e74c3c; color: white; }
.status.reconnecting { background: #f39c12; color: white; }
.cards { display: grid; grid-template-columns: repeat(4, 1fr); gap: 20px; margin: 20px 0; }
.card {
background: #1a1a3e; padding: 25px; border-radius: 12px;
text-align: center; border: 1px solid #2a2a5e;
}
.card .value { font-size: 2.5em; font-weight: bold; color: #00d4aa; }
.card .label { font-size: 0.9em; color: #888; margin-top: 5px; }
.events {
background: #1a1a3e; border-radius: 12px; padding: 20px;
max-height: 300px; overflow-y: auto; border: 1px solid #2a2a5e;
}
.event-item {
padding: 8px 12px; margin: 4px 0; border-radius: 6px;
font-size: 0.9em; background: #0f0f23;
}
.event-item.order { border-left: 3px solid #3498db; }
.event-item.alert { border-left: 3px solid #e74c3c; }
.event-item.metric { border-left: 3px solid #27ae60; }
</style>
</head>
<body>
<h1>📊 Dashboard
<span class="status disconnected" id="connStatus">Bağlanıyor...</span>
</h1>
<div class="cards">
<div class="card">
<div class="value" id="totalOrders">-</div>
<div class="label">Toplam Sipariş</div>
</div>
<div class="card">
<div class="value" id="pendingOrders">-</div>
<div class="label">Bekleyen</div>
</div>
<div class="card">
<div class="value" id="revenue">-</div>
<div class="label">Ciro (₺)</div>
</div>
<div class="card">
<div class="value" id="cpuUsage">-</div>
<div class="label">CPU %</div>
</div>
</div>
<h2>📡 Canlı Event'ler</h2>
<div class="events" id="eventsDiv"></div>
<script>
let eventSource;
function connect() {
eventSource = new EventSource('/api/dashboard/stream');
eventSource.onopen = () => {
setStatus('connected', 'Bağlı ✓');
};
eventSource.onerror = () => {
if (eventSource.readyState === EventSource.CONNECTING) {
setStatus('reconnecting', 'Yeniden bağlanıyor...');
} else {
setStatus('disconnected', 'Bağlı değil ✗');
}
};
// Dashboard ilk yüklemesi
eventSource.addEventListener('dashboard-init', (e) => {
const data = JSON.parse(e.data);
document.getElementById('totalOrders').textContent = data.totalOrders;
document.getElementById('pendingOrders').textContent = data.pendingOrders;
document.getElementById('revenue').textContent =
data.revenue.toLocaleString('tr-TR');
addEvent('Başlangıç verileri yüklendi', 'metric');
});
// Metrik güncellemeleri
eventSource.addEventListener('metrics', (e) => {
const data = JSON.parse(e.data);
document.getElementById('cpuUsage').textContent =
data.cpu.toFixed(1) + '%';
if (data.cpu > 80) {
addEvent('⚠️ CPU yüksek: ' + data.cpu.toFixed(1) + '%', 'alert');
}
});
// Yeni sipariş
eventSource.addEventListener('new-order', (e) => {
const order = JSON.parse(e.data);
addEvent('🛒 Yeni sipariş: ' + order.orderId +
' - ' + order.product, 'order');
// Sipariş sayısını artır
const el = document.getElementById('totalOrders');
el.textContent = parseInt(el.textContent) + 1;
});
// Stok uyarısı
eventSource.addEventListener('stock-alert', (e) => {
const alerts = JSON.parse(e.data);
alerts.forEach(alert => {
addEvent('🔴 Düşük stok: ' + alert.productName +
' (Kalan: ' + alert.currentStock + ')', 'alert');
});
});
// Heartbeat — sessizce yoksay
eventSource.addEventListener('heartbeat', () => {});
}
function setStatus(className, text) {
const el = document.getElementById('connStatus');
el.className = 'status ' + className;
el.textContent = text;
}
function addEvent(text, type) {
const div = document.getElementById('eventsDiv');
const time = new Date().toLocaleTimeString('tr-TR');
div.innerHTML = '<div class="event-item ' + type + '">' +
'<strong>' + time + '</strong> — ' + text + '</div>' + div.innerHTML;
// Max 50 event tut
while (div.children.length > 50) {
div.removeChild(div.lastChild);
}
}
// Başlat
connect();
</script>
</body>
</html>11. SSE ile Dikkat Edilmesi Gerekenler
HTTP/1.1 Connection Limiti
HTTP/1.1'de tarayıcı bir domain'e max 6 TCP bağlantısı açar.
Her SSE stream 1 bağlantı kullanır.
6 SSE bağlantısı = 0 kalan bağlantı = REST API çağrıları bloklanır!
Çözüm 1: HTTP/2 kullan (multiplexing — tek TCP'de yüzlerce stream)
Çözüm 2: Tek SSE endpoint, farklı event type'lar
Çözüm 3: Farklı subdomain'ler (sse.myapp.com)Proxy ve Buffering
Bazı reverse proxy'ler (Nginx) response'ları buffer'lar. SSE için buffering kapatılmalı:
# Nginx konfigürasyonu
location /api/sse/ {
proxy_pass http://backend;
proxy_set_header Connection '';
proxy_http_version 1.1;
chunked_transfer_encoding off;
proxy_buffering off; # Buffer'ı kapat!
proxy_cache off; # Cache'i kapat!
proxy_read_timeout 3600s; # Uzun timeout
}Spring Boot'ta response buffering:
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter stream(HttpServletResponse response) {
response.setHeader("X-Accel-Buffering", "no"); // Nginx buffering off
response.setHeader("Cache-Control", "no-cache"); // Cache off
return sseService.subscribe();
}Async Konfigürasyon
@Configuration
public class AsyncConfig implements WebMvcConfigurer {
@Override
public void configureAsyncSupport(AsyncSupportConfigurer configurer) {
configurer.setDefaultTimeout(30 * 60 * 1000L); // 30 dakika
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("sse-");
executor.initialize();
configurer.setTaskExecutor(executor);
}
}⚠️ Dikkat: Spring MVC'de her SSE bağlantısı bir servlet thread kullanır (async olsa bile request thread'i serbest kalır ama I/O thread hâlâ meşgul). 10.000 eşzamanlı SSE bağlantısı planlıyorsan, Spring WebFlux (reactive) daha uygun. Servlet stack'te thread pool tükenmesi riski var.
Özet
SSE, sunucudan client'a tek yönlü push için WebSocket'ten daha basit ve HTTP-native bir çözüm — dashboard, bildirim, canlı feed senaryolarında ideal
`SseEmitter` ile Spring MVC'de SSE endpoint oluşturmak çok kolay —
onCompletion(),onTimeout(),onError()callback'leriyle yaşam döngüsü kontrol edilir`CopyOnWriteArrayList` ile birden fazla client'a broadcast yapılır — dead emitter'ları temizlemeyi unutma, yoksa memory leak kaçınılmaz
Spring WebFlux ile
Flux<ServerSentEvent>kullanmak daha elegant — back-pressure, otomatik lifecycle yönetimi, yüksek concurrency'de üstün performansBrowser EventSource API otomatik reconnect sağlar —
Last-Event-IDheader'ı ile kaçırılan event'ler kurtarılabilir, her event'e mutlaka ID verSSE + REST combo, birçok senaryoda WebSocket'ten daha pratik — güncellemeleri SSE ile al, aksiyonları REST ile gönder, test ve debug çok daha kolay
AI Asistan
Sorularını yanıtlamaya hazır