Operacje okien czasowych w Kafka Streams
Dowiedz się, jak definiować okna czasowe do agregowania zdarzeń w Kafka Streams, w tym okna skokowe, przesuwne i kroczące.
Operacje okien czasowych w Kafka Streams to bezpłatna lekcja Apache Kafka & Stream Processing Fundamentals na CoddyKit. To lekcja 1 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Apache Kafka & Stream Processing Fundamentals, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Apache Kafka & Stream Processing Fundamentals zawiera 4 lekcji w sumie.
Części tej lekcji nie zostały jeszcze przetłumaczone i są wyświetlane po angielsku.
Intro to Stream Windows
In stream processing, data arrives continuously. To perform calculations like "total sales per hour" or "average temperature every 5 minutes," we need to group these continuous events into finite, manageable segments.
This grouping of events based on time is called windowing. It allows us to apply aggregations and transformations over specific periods.
Why Windowing Matters
Imagine analyzing website clicks. You might want to know:
- How many clicks happened in the last minute?
- What's the average user activity over the past 5 minutes, updated every minute?
- When did a user become inactive?
Windowing provides the framework to answer these questions by defining boundaries around your data streams.
Event Time vs. Processing Time
Kafka Streams primarily uses event time. This is the timestamp embedded in the event itself, indicating when the event actually happened at the source.
- Event Time: When the event occurred (preferred).
- Processing Time: When the event is processed by the stream application (less reliable for ordering).
Using event time ensures that results are consistent, even if events arrive out of order or with delays.
Tumbling Windows: Fixed & Non-Overlapping
Tumbling windows are like non-overlapping, fixed-size buckets of time. Each event belongs to exactly one window.
- They are fixed in duration (e.g., 5 minutes).
- They do not overlap.
- They are contiguous, covering all time.
Think of them as a series of distinct time slots, where each event falls into one specific slot.
Tumbling Window Code
Here's how to define a 5-second tumbling window in Kafka Streams to count messages. We use TimeWindows.of() and count().
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.TimeWindows;
import java.time.Duration;
import java.util.Properties;
public class TumblingWindowApp {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "tumbling-window-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
builder.stream("input-topic", Consumed.with(Serdes.String(), Serdes.String()))
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofSeconds(5))) // 5-second tumbling window
.count(Materialized.as("tumble-counts"))
.toStream()
.print(org.apache.kafka.streams.kstream.Printed.toSysOut());
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
}
}Hopping Windows: Overlapping & Moving
Hopping windows are also fixed-size, but they can overlap. They "hop" forward by a specified interval, which is typically smaller than the window size.
- Fixed size (e.g., 10 minutes).
- Overlap with previous/next windows.
- Hop interval (e.g., moves every 5 minutes).
This allows for smoother, more frequently updated aggregations, as a new window result is emitted more often.
Hopping Window Code
Here's how to create a 10-second hopping window that advances every 5 seconds. Events will be counted in multiple overlapping windows.
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.TimeWindows;
import java.time.Duration;
import java.util.Properties;
public class HoppingWindowApp {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "hopping-window-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
builder.stream("input-topic", Consumed.with(Serdes.String(), Serdes.String()))
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofSeconds(10)) // 10-second window
.advanceBy(Duration.ofSeconds(5))) // hops every 5 seconds
.count(Materialized.as("hop-counts"))
.toStream()
.print(org.apache.kafka.streams.kstream.Printed.toSysOut());
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
}
}Session Windows: Data-Driven Activity
Session windows are different. They are data-driven, not fixed-time. They group events that occur within a specified "inactivity gap."
- Window closes if no new event for a key arrives within the gap duration.
- Variable size, non-overlapping for a given key.
- Useful for user activity tracking or network sessions.
If a new event arrives for the same key after the inactivity gap, a new session window starts.
Session Window Code
This example shows how to define a session window with a 5-second inactivity gap. If a key doesn't receive an event for 5 seconds, its current session window closes.
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.SessionWindows;
import java.time.Duration;
import java.util.Properties;
public class SessionWindowApp {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "session-window-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
builder.stream("input-topic", Consumed.with(Serdes.String(), Serdes.String()))
.groupByKey()
.windowedBy(SessionWindows.with(Duration.ofSeconds(5))) // 5-second inactivity gap
.count(Materialized.as("session-counts"))
.toStream()
.print(org.apache.kafka.streams.kstream.Printed.toSysOut());
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
}
}Grace Period for Late Events
Events don't always arrive in order or on time. Kafka Streams handles this with a grace period.
- It's extra time a window stays open to accept late-arriving records.
- Defined as part of the window definition (e.g.,
.grace(Duration.ofSeconds(10))). - Events arriving after the grace period are typically dropped or forwarded to a late record topic.
This helps ensure accuracy for aggregations over event time.
Check Your Understanding
Which of the following statements about Kafka Streams windowing are TRUE?
Recap: Windowing in Kafka Streams
We've explored how windowing allows us to perform time-based aggregations on continuous data streams in Kafka Streams.
- Tumbling windows: Fixed, non-overlapping.
- Hopping windows: Fixed, overlapping, advance by an interval.
- Session windows: Data-driven, based on an inactivity gap.
- Grace period: Important for handling late-arriving events.
Mastering these window types is crucial for building robust real-time analytics and stream processing applications.
Ucz się Apache Kafka & Stream Processing Fundamentals dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 12
- Lekcje
- 48
Często zadawane pytania
Czy lekcja „Operacje okien czasowych w Kafka Streams” jest bezpłatna?
Tak — pełny tekst „Operacje okien czasowych w Kafka Streams” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Apache Kafka & Stream Processing Fundamentals, przejdź na CoddyKit PRO. Kurs Apache Kafka & Stream Processing Fundamentals zawiera 4 lekcji w sumie.
Co nauczysz się w „Operacje okien czasowych w Kafka Streams”?
Dowiedz się, jak definiować okna czasowe do agregowania zdarzeń w Kafka Streams, w tym okna skokowe, przesuwne i kroczące. Ćwiczysz Apache Kafka & Stream Processing Fundamentals z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Apache Kafka & Stream Processing Fundamentals?
Nie wymagamy żadnego doświadczenia. Apache Kafka & Stream Processing Fundamentals w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 1 z 4.
Ile czasu zajmuje lekcja „Operacje okien czasowych w Kafka Streams”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Apache Kafka & Stream Processing Fundamentals?
Tak. Każda lekcja Apache Kafka & Stream Processing Fundamentals zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Operacje okien czasowych w Kafka Streams
- Złączenia i agregacje w strumieniach
- Wprowadzenie do KSQL na potrzeby analizy strumieni
- Zapytania interaktywne i magazyny stanu