0Pricing
Apache Kafka & Stream Processing Fundamentals · บทเรียน

การดำเนินการแบ่งช่วงเวลาใน Kafka Streams

เรียนรู้การกำหนดช่วงเวลาตามเวลาเพื่อรวมเหตุการณ์ใน Kafka Streams รวมถึงช่วงเวลาแบบรวมกลุ่ม แบบกระโดด และแบบเลื่อน

การดำเนินการแบ่งช่วงเวลาใน Kafka Streams เป็นบทเรียน Apache Kafka & Stream Processing Fundamentals ฟรีบน CoddyKit นี่คือบทเรียนที่ 1 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 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” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส Apache Kafka & Stream Processing Fundamentals ให้อัปเกรดเป็น CoddyKit PRO คอร์ส Apache Kafka & Stream Processing Fundamentals มีบทเรียนทั้งหมด 4 บทเรียน

คุณจะเรียนรู้อะไรในบทเรียน “การดำเนินการแบ่งช่วงเวลาใน Kafka Streams”

เรียนรู้การกำหนดช่วงเวลาตามเวลาเพื่อรวมเหตุการณ์ใน Kafka Streams รวมถึงช่วงเวลาแบบรวมกลุ่ม แบบกระโดด และแบบเลื่อน คุณปฏิบัติ Apache Kafka & Stream Processing Fundamentals ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 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