← Back to list

Kafka Stream Eğitim Serisi — Bölüm 13 — Grace Period, Geç Gelen Olaylar, Suppression.

Modül 3 — Zaman ve Pencereleme

Sahin Yelkenci · 2026-06-22 07:04 · 0 claps · 8.1 min read
#kafka-streams #windowing #grace-period #late-events #event-stream-processing
Open on Medium ↗
Wiki topics: SOC · Sociology & Politics 📰 · Journalism & News

Kafka Stream Eğitim Serisi — Bölüm 13 — Grace Period, Geç Gelen Olaylar, Suppression.

Modül 3 — Zaman ve Pencereleme

İçindekiler » Geç gelen olay nedir? Kapanmış pencereye gelen kayıt » Grace period: pencereyi ne kadar açık tutalım? » Pencere ne zaman kapanır? Stream-time hatırlatması » Ara sonuç problemi: pencere her güncellemede yayılır » suppress: nihai sonucu tek seferde yaymak » Buffer config: suppress neyle çalışır? » suppress vs caching: hangisi ne zaman? » Lab: pencere kapanınca tek nihai sayım » *https://gitlab.com/sahin.yelkenci2/kafka-streams-canli-oran*

Geç Gelen Olay Nedir? Kapanmış Pencereye Gelen Kayıt

**↗ Bölüm 12'de pencereleri kurarken bilinçli bir şeyi erteledik: gerçek dünyada olaylar her zaman zamanında gelmez. Bir geç gelen olay (late event)**, event-time’ı şu anki stream-time’ın gerisinde kalan bir penceredeki kayıttır — yani ait olduğu pencere “geçilmiş” görünürken çıkagelen bir olay. Bu, bizim sistemimizde teorik bir ihtimal değil, somut bir gerçektir: relay at-least-once çalışır ve hattımızda gecikme vardır, dolayısıyla occurred_at'i 14:30:00 olan bir olay, başka olaylar daha geç occurred_at'lerle önce geldiği için stream-time 14:30:25'e ilerledikten sonra ulaşabilir.

**↗ Bölüm 11**'den hatırlayalım: stream-time, görülen en yüksek event-time’dır ve yalnızca ileri gider. Daha küçük event-time’lı bir kayıt sonradan geldiğinde stream-time geri gitmez; o kayıt “geç” sayılır. Bu geç kaydın iki olası kaderi vardır: ait olduğu pencere hâlâ açıksa, sorunsuz oraya eklenir; ama pencere kapanmışsa, ya düşürülür ya da özel bir tolerans (grace period) varsa hâlâ kabul edilir. İşte bu bölüm, o “kapanmış pencereye gelen kayıt” durumunu doğru yönetmekle ilgili. Aşağıdaki diyagram bir geç olayın varışını gösterir.

Grace Period: Pencereyi Ne Kadar Açık Tutalım?

Geç gelen olayları doğru yönetmenin aracı grace period’dur: bir pencerenin bitiş zamanından sonra, geç kayıtları hâlâ kabul etmeye devam ettiği tolerans süresi. TimeWindows.ofSizeAndGrace(boyut, grace) ile tanımlanır. Diyelim ki 30 saniyelik bir pencerenin grace'i 5 saniye; o zaman [00:00–00:30) penceresi, stream-time 00:35'i geçene kadar geç gelen kayıtları hâlâ kendine ekler. Bu noktadan sonra pencere kapanır ve ona ait yeni varışlar düşürülür (ve bir "late-record-drop" metriği olarak gözlemlenebilir — **↗ Bölüm 33**).

Grace değerini seçmek bir takastır. Uzun bir grace, geç gelen stragglerlar için daha çok doğruluk sağlar — daha az olay kaçar — ama pencereleri daha uzun süre açık ve state’i daha uzun süre dolu tutar, ayrıca nihai sonucun yayınlanmasını geciktirir. Kısa bir grace ise hızlı sonuç verir ama gerçek geç olayları kaçırır. Doğru değer, hattınızın gerçekçi gecikme profiline göre belirlenir: relay ve tüketim gecikmeniz tipik olarak ne kadarsa, grace onu kapsayacak kadar olmalı, ama daha fazlası değil. Aşağıdaki kod 30 saniye boyutlu, 5 saniye grace’li bir pencere tanımlar.

// 30 sn pencere + 5 sn grace: pencere bitişinden 5 sn sonrasına kadar geç kayıt kabul
TimeWindows pencere =
    TimeWindows.ofSizeAndGrace(Duration.ofSeconds(30), Duration.ofSeconds(5));

Çalışan örnek kodlara erişmek için: » https://gitlab.com/sahin.yelkenci2/kafka-streams-canli-oran

Pencere Ne Zaman Kapanır? Stream-Time Hatırlatması

Burada **↗ Bölüm 11'deki kritik bir gerçeği tekrar vurgulamak şart, çünkü pencere kapanışının kalbinde o yatıyor: bir pencere, stream-time, pencere_bitişi + grace değerini geçtiğinde** kapanır. Ve stream-time duvar saati değildir — yalnızca gelen kayıtlarla ilerler. Yani bir pencere, takvim saati ilerlediği için değil, yeterince yüksek event-time’lı yeni bir kayıt stream-time’ı kapanış noktasının ötesine ittiği için kapanır.

Bunun çok pratik ve sinsi bir sonucu vardır: akış sessizleşirse, son pencere hiç kapanmayabilir. Bir koşu bittiğinde oran akışı durur; eğer stream-time’ı ilerletecek yeni kayıt gelmezse, o koşunun son penceresi açık kalır ve — birazdan göreceğimiz — suppress edilmiş nihai sonuçları asla yayınlanmaz. Bu, “neden son uyarıyı hiç almadım” sorusunun yaygın cevabıdır. Çözüm, duvar saatine bağlı bir punctuator ile pencereleri zorla kapatmak ya da düzenli bir “heartbeat” beslemektir; bu mekanizmayı **↗ Bölüm 22**'de ele alacağız. Şimdilik kuralı aklımızda tutalım: pencereyi kapatan stream-time’dır, ve stream-time’ı ilerleten yeni kayıtlardır.

Ara Sonuç Problemi: Pencere Her Güncellemede Yayılır

Pencereli bir toplulaştırmanın varsayılan davranışında, çoğu zaman istemediğimiz bir özellik vardır. Bir pencere her yeni kayıtla güncellendiğinde, yeni bir sonuç yayınlar (caching ayarına bağlı olarak — **↗ Bölüm 8). Yani 30 saniyelik bir pencerede sayı 1'den 2'ye, 2'den 3'e çıkarken, aşağı akış aynı pencere için üç ayrı yayın görür: “sayı 1”, “sayı 2”, “sayı 3”. [↗ Bölüm 12](https://medium.com/@sahinyelkenci/kafka-stream-e%C4%9Fitim-serisi-b%C3%B6l%C3%BCm-12-pencere-t%C3%BCrleri-tumbling-hopping-sliding-session-404ad3dc54f2)**'deki manipülasyon lab’ımız tam da böyle davranıyordu — pencere doldukça ara sayımları aşağıya akıtıyordu.

Çoğu durumda istediğimiz bu gürültülü ara akış değil, pencere başına tek bir nihai sonuçtur. “Bu 30 saniyelik pencerede toplam 3 değişim oldu” demek isteriz, “1 oldu, sonra 2 oldu, sonra 3 oldu” değil. Özellikle bir karar (manipülasyon var mı?) ya da bir rapor üretiyorsak, ara sonuçlar hem gereksiz trafik hem de yanlış erken kararlar (henüz pencere dolmadan “eşiği geçti” sanmak) yaratır. İşte bu ara sonuç problemini suppress çözer. Aşağıdaki diyagram istenen ile varsayılan davranışı karşılaştırır.

suppress: Nihai Sonucu Tek Seferde Yaymak

suppress, pencereli bir toplulaştırmanın ara sonuçlarını geri tutup yalnızca nihai sonucu yayınlamasını sağlar. suppress(Suppressed.untilWindowCloses(...)) dediğinizde, Streams o pencereye ait tüm ara güncellemeleri bir tampona alır ve hiçbirini aşağıya geçirmez; ancak pencere kapandığında (yani stream-time, bitiş + grace'i geçtiğinde) o pencerenin tek ve nihai sonucunu yayınlar. Böylece gürültülü ara akış, pencere başına tam bir yayına dönüşür.

Bu, “her pencere için kesin olarak bir sonuç” garantisidir — ne eksik ne fazla. Manipülasyon kararımız için bu tam olarak istediğimiz şeydir: pencere kapanana ve tüm geç kayıtlar (grace dahilinde) toplanana kadar bekle, sonra o pencerenin kesin sayısını bir kez söyle. Aşağıdaki kod windowed bir sayıma suppress ekler.

KTable<Windowed<String>, Long> nihaiSayim =
    oranAkisi
        .groupByKey(Grouped.with(Serdes.String(), oddsSerde))
        .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofSeconds(30), Duration.ofSeconds(5)))
        .count(Materialized.as("oran-degisim-30sn"))
        // ara sonuçları tut, yalnızca pencere kapanınca tek nihai sonucu yay
        .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()));

Çalışan örnek kodlara erişmek için: » https://gitlab.com/sahin.yelkenci2/kafka-streams-canli-oran

Buffer Config: suppress Neyle Çalışır?

suppress, henüz kapanmamış pencerelerin "bekleyen" sonuçlarını bir tamponda tutmak zorundadır — ve bu tamponu Suppressed.BufferConfig ile yapılandırırsınız. İki boyutu vardır: tamponun ne kadar büyüyebileceği (kayıt sayısı maxRecords ya da bayt maxBytes) ve dolduğunda ne yapılacağı. untilWindowCloses modu için iki seçenek geçerlidir: ya unbounded() (sınırsız tampon — bellek riski var ama "yalnızca nihai" garantisini korur), ya da maxBytes(...).shutDownWhenFull() gibi sınırlı ama dolduğunda uygulamayı durduran bir yapı.

Burada kritik bir kısıt vardır: untilWindowCloses ile erken yayın yapamazsınız (emitEarlyWhenFull kullanamazsınız), çünkü erken yaymak "yalnızca nihai sonuç" garantisini bozardı — o yüzden ya sınırsız tampona ya da dolunca kapanmaya razı olursunuz. Kapasite açısından bu önemlidir: tampon, aynı anda açık olan tüm pencerelerin bekleyen sonuçlarını tutar; çok sayıda eşzamanlı açık pencereniz varsa (çok koşu, uzun grace) tampon büyür. Bu boyutlandırmayı **↗ Bölüm 34**'te ele alacağız. Aşağıdaki kod iki buffer config seçeneğini gösterir.

// Seçenek 1: sınırsız — "yalnızca nihai" garantisini korur, bellek riski taşır
Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded());
// Seçenek 2: sınırlı - dolunca uygulamayı durdur (fail-fast), bellek tahmin edilebilir
Suppressed.untilWindowCloses(
    Suppressed.BufferConfig.maxBytes(64 * 1024 * 1024).shutDownWhenFull());

suppress vs Caching: Hangisi Ne Zaman?

suppress ile **↗ Bölüm 8'de gördüğümüz caching kolayca karıştırılır, ama ikisi farklı şeylerdir ve bu ayrım önemlidir. Caching, ara yayınları fırsatçı biçimde azaltır — tampon boyutuna ve commit aralığına bağlı olarak bazı ara güncellemeleri birleştirip atlar. Ama bu en iyi çaba** (best-effort) bir azaltmadır; zamanlamaya bağlıdır ve "yalnızca nihai sonuç" garantisi vermez. Caching açıkken yine de pencere için birden çok ara yayın görebilirsiniz, sadece daha az.

suppress(untilWindowCloses) ise deterministiktir: pencere başına kesinlikle bir nihai yayın, hiç ara yayın yok — caching ayarından bağımsız olarak. Yani garanti istiyorsanız suppress kullanırsınız; sadece gürültüyü azaltmak yetiyorsa ve kesin garanti gerekmiyorsa caching yeterli olabilir. İkisi birlikte de kullanılabilir, ama "tek nihai sonuç" gereksinimini karşılayan suppress'tir. Kural: bir karar ya da rapor üretiyorsanız (manipülasyon var mı, pencere toplamı kaç) suppress; yalnızca downstream trafiğini hafifletmek istiyorsanız caching.

Lab: Pencere Kapanınca Tek Nihai Sayım

**↗ Bölüm 12'deki Senaryo B’yi olgun haline getirelim. Orada manipülasyon sayımını pencere doldukça ara ara yayınlıyorduk — gürültülü ve erken-karar riskli. Şimdi onu, her 30 saniyelik pencere için tek bir nihai karar** üretecek şekilde yeniden kuruyoruz: 5 saniyelik grace ile geç gelen olayları da topluyor, suppress(untilWindowCloses) ile pencere kapanana kadar bekleyip yalnızca nihai sayıyı yayınlıyoruz. Böylece "bu 30 saniyede toplam şu kadar değişim oldu, dolayısıyla manipülasyon var/yok" kararı bir kez, doğru ve gürültüsüz çıkar.

Tek bir operasyonel uyarıyı tekrar hatırlatalım: bu hat, son pencereyi ancak stream-time onu kapanış noktasının ötesine ittiğinde yayınlar; akış sessizleşirse son pencere asılı kalabilir, bunun çözümü **↗ Bölüm 22**'deki wall-clock punctuator’dır. Aşağıdaki topology Senaryo B’nin nihai halini kurar.

StreamsBuilder builder = new StreamsBuilder();
Serde<FixedOddsChanged> oddsSerde = new FixedOddsChangedSerde();
Serde<ManipulasyonUyarisi> uyariSerde = new ManipulasyonUyarisiSerde();
KStream<String, FixedOddsChanged> oranAkisi =
    builder.stream("fixed-odds-changed",
        Consumed.with(Serdes.String(), oddsSerde)
                .withTimestampExtractor(new OccurredAtExtractor()));   // event-time
builder // tek nihai sayım üreten hat
    .stream("fixed-odds-changed",
        Consumed.with(Serdes.String(), oddsSerde).withTimestampExtractor(new OccurredAtExtractor()))
    .groupByKey(Grouped.with(Serdes.String(), oddsSerde))
    .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofSeconds(30), Duration.ofSeconds(5)))
    .count(Materialized.as("oran-degisim-30sn"))
    // ara sonuçları tut; yalnızca pencere kapanınca (end + 5sn grace) tek nihai yay
    .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
    .toStream()
    .filter((winKey, sayi) -> sayi != null && sayi >= 3)            // nihai sayı ≥ 3
    .map((winKey, sayi) -> KeyValue.pair(
        winKey.key(),
        ManipulasyonUyarisi.newBuilder()
            .setAnahtar(winKey.key())
            .setPencereBaslangic(winKey.window().start())
            .setPencereBitis(winKey.window().end())
            .setDegisimSayisi(sayi)
            .build()))
    .to("manipulasyon-uyarisi", Produced.with(Serdes.String(), uyariSerde));
Topology topology = builder.build();
log.info("Topology:\n{}", topology.describe());  // window store + suppress buffer görünür

Çalışan örnek kodlara erişmek için: » https://gitlab.com/sahin.yelkenci2/kafka-streams-canli-oran

Bölüm Özeti ve Sıradaki Adım

Bu bölümde pencerelemenin gerçek dünyayla yüzleştiği yeri ele aldık: olaylar her zaman zamanında gelmez. Geç gelen olayın, event-time’ı şu anki stream-time’ın gerisinde kalan bir kayıt olduğunu — ve at-least-once relay ile hattımızda bunun gerçek olduğunu — gördük. Grace period’un, bir pencerenin bitişinden sonra geç kayıtları kabul etmeye devam ettiği tolerans süresi olduğunu, ve değerinin hattın gecikme profiline göre bir takas olduğunu öğrendik. Bir pencerenin stream-time, bitiş + grace’i geçtiğinde kapandığını ve stream-time yalnızca yeni kayıtlarla ilerlediği için sessiz bir akışta son pencerenin asılı kalabileceğini (çözümü ↗ Bölüm 22) vurguladık. Ara sonuç problemini — pencerenin her güncellemede yayın yapmasını — tanıdık ve suppress(untilWindowCloses) ile pencere başına tek nihai sonuç üretmeyi, buffer config seçeneklerini ve suppress ile caching arasındaki "garanti vs en iyi çaba" farkını netleştirdik.

Lab’da Senaryo B’yi olgunlaştırdık: 5 saniye grace ile geç olayları toplayan, suppress ile her 30 saniyelik pencere için tek ve gürültüsüz bir manipülasyon kararı üreten bir hat. Artık zamanı doğru okumayı (Bölüm 11), pencerelemeyi (Bölüm 12) ve geç olaylarla nihai yayını (bu bölüm) biliyoruz. Ama hâlâ örtük bir varsayım var: olayların aşağı yukarı sıralı geldiğini varsaydık ve geç olanları "istisna" saydık. Gerçekte, at-least-once bir relay ile sıra bozulması ve duplicate'ler kuraldır, istisna değil. Sıradaki bölüm bu gerçekle doğrudan yüzleşiyor.

» Sıradaki — ↗ Bölüm 14: Sıra Bozulması (Out-of-Order). Relay at-least-once çalıştığında sıranın neden bozulduğunu ve duplicate’lerin neden kural olduğunu, Streams’in bu gerçeğe event-time ve stream-time ile nasıl yanıt verdiğini, ve out-of-order bir oran akışında doğruluğu nasıl koruyacağımızı, hep oran verimiz üzerinden işleyeceğiz.

**İçindekiler… « Önceki [Bölüm 12 — Pencere Türleri. Tumbling, hopping, sliding, session.] » Sonraki [Bölüm 14 — Sıra Bozulması (Out-of-Order).] » Proje Kodları…** » https://gitlab.com/sahin.yelkenci2/kafka-streams-canli-oran » https://gitlab.com/sahin.yelkenci2/kafka-streams-spring-oran


메타데이터
post_id
efd5c00da243
slug
kafka-stream-eğitim-serisi-bölüm-13-grace-period-geç-gelen-olaylar-suppression-efd5c00da243
url
https://medium.com/@sahinyelkenci/kafka-stream-e%C4%9Fitim-serisi-b%C3%B6l%C3%BCm-13-grace-period-ge%C3%A7-gelen-olaylar-suppression-efd5c00da243
canonical_url
https://medium.com/@sahinyelkenci/kafka-stream-e%C4%9Fitim-serisi-b%C3%B6l%C3%BCm-13-grace-period-ge%C3%A7-gelen-olaylar-suppression-efd5c00da243
author_url
https://medium.com/@sahinyelkenci
status
ok
fetched_at
2026-07-08 21:20:17