0Pricing
Apache Kafka & Stream Processing Fundamentals · Урок

Оконные операции в Kafka Streams

Научитесь определять временные окна для агрегации событий в Kafka Streams, включая фиксированные, скользящие и раздвижные окна.

«Оконные операции в Kafka Streams» — бесплатный урок Apache Kafka & Stream Processing Fundamentals на CoddyKit. Это урок 1 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Apache Kafka & Stream Processing Fundamentals, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Apache Kafka & Stream Processing Fundamentals содержит 4 уроков всего.

Части этого урока еще не переведены и отображаются на английском.

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.

Часто задаваемые вопросы

Урок «Оконные операции в Kafka Streams» бесплатный?

Да — полный текст урока «Оконные операции в Kafka Streams» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Apache Kafka & Stream Processing Fundamentals, подпишись на CoddyKit PRO. Курс Apache Kafka & Stream Processing Fundamentals содержит 4 уроков всего.

Чему я научусь в уроке «Оконные операции в Kafka Streams»?

Научитесь определять временные окна для агрегации событий в Kafka Streams, включая фиксированные, скользящие и раздвижные окна. Ты практикуешь Apache Kafka & Stream Processing Fundamentals с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Apache Kafka & Stream Processing Fundamentals?

Предыдущий опыт не требуется. Apache Kafka & Stream Processing Fundamentals на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 1 из 4.

Сколько времени занимает урок «Оконные операции в Kafka Streams»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке Apache Kafka & Stream Processing Fundamentals?

Да. Каждый урок Apache Kafka & Stream Processing Fundamentals включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Оконные операции в Kafka Streams
  2. Объединения и агрегации в потоках
  3. Введение в KSQL для потоковой аналитики
  4. Интерактивные запросы и хранилища состояния
← Назад к Apache Kafka & Stream Processing Fundamentals