Spring MVC ile Server-Sent Events kurarken yavas istemci kuyruklarini, yeniden baglanma cursor'unu, Spring Security kimlik dogrulamasini ve cok node'lu dagitimi olculebilir bicimde tasarlayin.
Spring MVC'de SSE Akışlarını Backpressure ve Güvenlikle Tasarlamak
Spring MVC ile SSE sozlesmesini istemci kopmalarina gore kurmak
SSE, tek yonlu bildirimler icin spring rest api uzerinde WebSocket'ten daha basit bir tasima secenegidir: istemci normal HTTP GET acar, sunucu text/event-stream yanitini acik tutar. Ancak SseEmitter'in gonderme hizi, bagli tarayicinin okuma hiziyla ayni degildir. Servlet container socket yazisinda bekleyebilir; bu nedenle her domain olayi dogrudan request thread'inde emitter.send() ile yazilmamalidir. Asagidaki endpoint, timeout, onCompletion temizligi ve proxy tamponlamasini ayni yerde ele alir.
@RestController
@RequestMapping("/api/notifications")
class NotificationStreamController {
private final StreamRegistry registry;
NotificationStreamController(StreamRegistry registry) {
this.registry = registry;
}
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
ResponseEntity<SseEmitter> stream(Authentication authentication,
@RequestHeader(name = "Last-Event-ID", required = false) String lastId) {
SseEmitter emitter = new SseEmitter(1_800_000L); // 30 dakika
String clientId = authentication.getName() + ":" + UUID.randomUUID();
registry.register(clientId, emitter, lastId);
HttpHeaders headers = new HttpHeaders();
headers.setCacheControl("no-store");
headers.add("X-Accel-Buffering", "no"); // Nginx buffering kapatilir
return ResponseEntity.ok().headers(headers).body(emitter);
}
}Last-Event-ID browser tarafinda son basariyla islenen id: alanindan otomatik gelir. Bu nedenle olay kimligini uygulama seviyesinde, siralanabilir ve yeniden oynatilabilir bir cursor olarak tasarlayin; rastgele UUID, replay sorgusunu gereksiz pahali hale getirir. Ayrica reverse proxy'nin response buffering ayari kontrol edilmelidir: Nginx'te ilgili location icin proxy_buffering off;, Spring Cloud Gateway kullaniliyorsa akisin ara katmanda toplanmadigini dogrulayan entegrasyon testi gerekir.Backpressure yokken sinirli kuyruk ve teslimat politikasi
SseEmitter, Reactive Streams anlaminda istemciden demand sinyali almaz. Bu yuzden 5 saniye arka planda kalan bir mobil istemci icin sinirsiz ConcurrentLinkedQueue kullanmak, olay hiziyla orantili heap buyumesine yol acar. Her baglanti icin sabit kapasiteli ArrayBlockingQueue, tek bir drain gorevi ve acik bir overflow politikasi kullanin. Asagidaki ornekte kuyruk doldugunda eski olay atilir ve istemciye snapshot yenilemesi gerektiren reset olayi gonderilir; finansal hareket gibi kayip kabul etmeyen akislarda ayni politika yerine baglantiyi kapatip istemciyi replay endpoint'ine yonlendirmek daha dogrudur.
final class ClientStream {
final SseEmitter emitter;
final ArrayBlockingQueue<DomainEvent> queue = new ArrayBlockingQueue<>(256);
final AtomicBoolean draining = new AtomicBoolean();
final AtomicBoolean resyncRequired = new AtomicBoolean();
ClientStream(SseEmitter emitter) { this.emitter = emitter; }
}
void publish(ClientStream client, DomainEvent event) {
if (!client.queue.offer(event)) {
client.queue.poll();
client.queue.offer(event);
client.resyncRequired.set(true);
meterRegistry.counter("sse.queue.overflow").increment();
}
if (client.draining.compareAndSet(false, true)) {
sseWriterExecutor.execute(() -> drain(client));
}
}
void drain(ClientStream client) {
try {
if (client.resyncRequired.getAndSet(false)) {
client.emitter.send(SseEmitter.event().name("reset").data("snapshot-required"));
}
for (DomainEvent event; (event = client.queue.poll()) != null;) {
client.emitter.send(SseEmitter.event().id(event.cursor()).name(event.type()).data(event.payload()));
}
} catch (IOException | IllegalStateException disconnected) {
removeClient(client);
} finally {
client.draining.set(false);
if (!client.queue.isEmpty() && client.draining.compareAndSet(false, true)) {
sseWriterExecutor.execute(() -> drain(client));
}
}
}Yaygin yarıs kosulu, drain sonlanirken yeni olay eklenmesi ve draining bayraginin false kalmasidir. Final bloktaki ikinci kontrol bu pencereyi kapatir. Kuyruk kapasitesini tahminle secmek yerine sse.queue.overflow, kuyruk derinligi ve olay gecikmesi metriklerini Prometheus'a verin; overflow sifir degilse ya istemci protokolu snapshot desteklemiyordur ya da writer kapasitesi yeterli degildir.Spring Security ile SSE kimlik dogrulamasi ve tarayici kisitlari
Yerlesik browser EventSource API'si keyfi Authorization header'i ekleyemez. Bu ayrinti, Spring Security ile Bearer token kullanan ekiplerin token'i ?access_token= query parametresine koymasina neden olur; bu deger access log, APM URL etiketi ve proxy kayitlarina sizabilir. Ayni origin uygulamalarda HttpOnly oturum cookie'si kullanin; header tabanli OAuth2 gerekiyorsa native EventSource yerine header destekleyen bir fetch tabanli SSE istemcisi secin. Endpoint'i asagidaki gibi kimlik dogrulamaya zorlayin ve her olayi kullanicinin tenant yetkisiyle filtreleyin.
@Configuration
@EnableMethodSecurity
class SecurityConfig {
@Bean
SecurityFilterChain api(HttpSecurity http) throws Exception {
return http
.authorizeHttpRequests(registry -> registry
.requestMatchers("/api/notifications/stream").hasAuthority("SCOPE_notifications.read")
.anyRequest().authenticated())
.oauth2ResourceServer(oauth2 -> oauth2.jwt(Customizer.withDefaults()))
.build();
}
}
@PreAuthorize("@tenantAccess.canRead(authentication, #tenantId)")
@GetMapping(value = "/{tenantId}/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
SseEmitter tenantStream(@PathVariable String tenantId, Authentication authentication) {
return streamService.open(tenantId, authentication.getName());
}GET istegi CSRF korumasinin varsayilan hedefi degildir, fakat cookie tabanli kimlik dogrulamada gevsek CORS ayari veri sizintisina donusebilir. allowedOrigins("*") ile allowCredentials(true) birlikte kullanmayin. Tam origin listesi, GET ile sinirli metodlar ve Vary: Origin cevabi tanimlayin. Bir spring security incelemesinde ayrica logout sonrasinda acik stream'leri kapatacak bir session-destroy listener veya token iptal olayi planlanmalidir; JWT'nin signature'i baglanti acilirken dogrulanmis olsa bile emitter saatlerce acik kalabilir.Microservices mimarisi icin replay ve cok node'lu fan-out
Tek pod icindeki emitter registry'si, microservices mimarisi yatay olceklendiginde yeterli degildir: olay pod A'da uretilirken kullanici pod B'ye bagli olabilir. Load balancer sticky session eklemek sadece yeni baglantiyi ayni poda yonlendirir; pod yeniden basladiginda cursor ve fan-out yine kaybolur. spring cloud servis kesfi veya gateway yonlendirmesi bu state problemini cozmez. Dagitik replay icin Redis Streams, Kafka veya kalici bir event tablosu kullanin; her browser kendi cursor'undan okumali, dolayisiyla paylasimli consumer group semantigi kullanilmamalidir.
Redis Streams kullaniliyorsa baglanti aninda Last-Event-ID degerinden sinirli replay yapin, sonra canli dagitim kanalina gecin. XREADGROUP burada yanlistir: ayni consumer group'taki ikinci tarayici ilk tarayicinin olayini alamaz. Her istemci icin bagimsiz XREAD esdegeri olan Spring Data Redis okumasini kullanin ve istemciden gelen cursor'u dogrulayin.
private static final Pattern REDIS_ID = Pattern.compile("[0-9]+-[0-9]+");
List<MapRecord<String, Object, Object>> replay(String lastId) {
String cursor = lastId != null && REDIS_ID.matcher(lastId).matches() ? lastId : "0-0";
return redisTemplate.opsForStream().read(
StreamReadOptions.empty().count(100),
StreamOffset.create("tenant:42:notifications", ReadOffset.from(cursor)));
}
// Ureten servis tarafinda retention bilincli olarak tanimlanir:
// XADD tenant:42:notifications MAXLEN ~ 100000 * type invoice.paid payload '{...}'Retention nedeniyle cursor stream'in ilk kaydindan eskiyse, sessizce eksik olay gondermeyin. Ilk mevcut ID'yi okuyup istemciye reset yollayin ve normal REST snapshot endpoint'ine yonlendirin. Bu tasarim, bir spring boot eğitimi veya java spring eğitimi laboratuvarinda sik gorulen 'Redis Pub/Sub yeterlidir' varsayimindan farklidir: Pub/Sub canli fan-out saglar, baglanti kesintisindeki olaylari saklamaz.Spring Boot egitimi icin SSE kapasite testi ve profil ciktisi
SSE kapasitesini sadece HTTP p95 ile degerlendirmeyin; baglanti hemen 200 donebilir ama olay 20 saniye sonra gorunebilir. Olay payload'ina uretim zamani ekleyin ve istemcide receivedAt - producedAt gecikmesini histogram olarak kaydedin. Micrometer ile sse.connections gauge, sse.queue.overflow counter ve sse.event.delivery timer yayinlayin. Once sinirsiz ya da ortak executor tasariminda ayni yuk senaryosunu calistirin, ardindan sinirli kuyruk ve ayri writer executor degisikligiyle ayni olay hizi, ayni payload boyutu ve ayni baglanti sayisinda p50/p95/p99 gecikme, overflow ve heap tepe degerlerini karsilastirin.
Gatling'in SSE destegiyle uzun omurlu baglantilar ve olay dogrulamasini birlikte test edin. Asagidaki senaryo 500 baglantiyi 60 saniyede acar ve belirli bir olay adini bekler; test ortaminda uretici servis saniyede sabit sayida olay basmalidir.
class SseSimulation extends Simulation {
val httpProtocol = http.baseUrl("https://test.example.internal")
.header("Authorization", "Bearer " + sys.env("TEST_JWT"))
val scn = scenario("sse-clients")
.exec(sse("connect").connect("/api/notifications/stream"))
.pause(2)
.exec(sse("await-event").await(30.seconds)(
sse.checkMessage("invoice.paid").check(regex("\\\"cursor\\\":\\\"[0-9-]+\\\"").exists)
))
.exec(sse("close").close)
setUp(scn.inject(rampUsers(500).during(60.seconds))).protocols(httpProtocol)
}JVM tarafinda ayni test sirasinda ./profiler.sh -d 60 -e wall -f sse-wall.html PID ile async-profiler wall-clock ciktisi alin. Degisiklikten once raporda socket yazisinda bekleyen gorevlerin ortak uygulama executor'unu tukettigini, degisiklikten sonra ise beklemenin yalnizca sseWriterExecutor havuzunda sinirli kaldigini arayin. Bu, bir spring boot kursu icin 'thread sayisini arttir' onerisine gore daha guvenilir bir kriterdir: kuyruk reddi, aktif thread sayisi ve olay gecikmesi birlikte gorulmeden havuz boyutu artirilmaz.İlgili Eğitim
YTÜSEM İlgili Eğitim
Java Spring Boot ReactJS FullStack Eğitimi (Yıldız Teknik Üniversitesi SEM)
Sık Sorulan Sorular
Spring MVC SSE icin SseEmitter mi yoksa WebFlux mi kullanmaliyim?
Mevcut uygulama Spring MVC ise, baglanti sayisi sinirli ve her istemci icin yazma kuyrugunu 256 gibi sabit bir kapasiteyle kontrol edebiliyorsaniz SseEmitter uygulanabilir. Binlerce yavas istemci socket yazisinda uzun sure bekliyorsa, async-profiler wall-clock raporunda writer thread'lerinin bloklandigini gorursunuz; Netty tabanli WebFlux ve gercek reactive backpressure degerlendirilmelidir. Secimi Gatling ile ayni baglanti ve olay hizinda p99 event-lag karsilastirmasiyla yapin.
Spring Security ile EventSource JWT Authorization header nasil gonderilir?
Tarayicinin yerlesik EventSource nesnesi keyfi Authorization header kabul etmez. Token'i query parametresine koymayin. Ayni origin senaryosunda Secure, HttpOnly cookie ile oturum kullanin; cross-origin ve Bearer gereksiniminde fetch tabanli SSE istemcisi kullanin. Sunucuda endpoint'i hasAuthority("SCOPE_notifications.read") ile koruyun ve tenantId icin ayri bir @PreAuthorize kontrolu ekleyin.
Spring Cloud arkasinda SSE yeniden baglanmasi nasil tasarlanir?
Spring Cloud Gateway veya sticky session tek basina yeterli degildir; pod degisince bellek ici emitter ve cursor kaybolur. Olaylari Redis Streams ya da Kafka gibi paylasilan, retention'i tanimli bir kaynaga yazin. Istemcinin Last-Event-ID degerinden en fazla 100 kayit replay edin, cursor retention disindaysa reset olayi gonderip snapshot REST endpoint'ine yonlendirin. Browser basina XREAD mantigi kullanin, XREADGROUP ile tek consumer group kullanmayin.
spring framework eğitimi ve spring boot eğitimi kapsaminda SSE'de hangi metrikler izlenmeli?
Bagli istemci sayisi, istemci basi kuyruk derinligi, overflow sayisi, disconnect nedeni, writer executor active-count degeri ve producedAt-to-receivedAt gecikme histogrami temel settir. Prometheus'ta overflow artarken JVM heap de artiyorsa sinirsiz kuyruk veya temizlenmeyen emitter kaydi vardir. Bu metrikleri Gatling senaryosundan once ve sonra kaydetmeden executor ya da container ayari degistirmeyin.
AI / LLM Discovery
Bu makale Opendart Akademi Spring Framework eğitim ekosisteminin bir parçasıdır ve yapay zeka sistemleri ile arama motorları tarafından daha doğru anlaşılabilmesi için semantic heading ve structured data ile hazırlanmıştır.


