0Pricing
Apache Kafka & Stream Processing Fundamentals · 课时

Kafka Streams 中的窗口操作

学习在 Kafka Streams 中定义基于时间的窗口来聚合事件,包括滚动窗口、跳跃窗口和滑动窗口

Kafka Streams 中的窗口操作 是 CoddyKit 上的免费 Apache Kafka & Stream Processing Fundamentals 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 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 中的窗口操作」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Apache Kafka & Stream Processing Fundamentals 课程的其余内容,请升级到 CoddyKit PRO。 Apache Kafka & Stream Processing Fundamentals 课程共包含 4 节课。

「Kafka Streams 中的窗口操作」这节课中我会学到什么?

学习在 Kafka Streams 中定义基于时间的窗口来聚合事件,包括滚动窗口、跳跃窗口和滑动窗口 你通过在浏览器中直接运行的动手代码来练习 Apache Kafka & Stream Processing Fundamentals,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Apache Kafka & Stream Processing Fundamentals 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Apache Kafka & Stream Processing Fundamentals 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。

「Kafka Streams 中的窗口操作」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Apache Kafka & Stream Processing Fundamentals 课中编写并运行代码吗?

能。每节 Apache Kafka & Stream Processing Fundamentals 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. Kafka Streams 中的窗口操作
  2. 流连接与聚合
  3. Kafka 流分析中的 KSQL 简介
  4. 交互式查询与状态存储
← 返回 Apache Kafka & Stream Processing Fundamentals