0Pricing
Apache Kafka & Stream Processing Fundamentals · درس

مفاهيم KStream وKTable

ميّز بين KStream (تدفق سجل تلو الآخر) وKTable (تدفق سجل تغييرات يمثّل عرضًا ماديًا) في Kafka Streams

مفاهيم KStream وKTable درس مجاني في Apache Kafka & Stream Processing Fundamentals على CoddyKit. هذا هو الدرس 2 من أصل 4. يمكنك قراءة الدرس كاملاً أدناه مجاناً — ثم تمرن عليه مباشرة في المتصفح باستخدام محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7. هذا الدرس جزء من مسار التعلم في Apache Kafka & Stream Processing Fundamentals، وتقدمك يتزامن عبر الويب وتطبيق CoddyKit. تتضمن دورة Apache Kafka & Stream Processing Fundamentals 4 دروس في المجموع.

بعض أجزاء هذا الدرس لم تُترجم بعد وتظهر باللغة الإنجليزية.

KStream & KTable Unveiled

In Kafka Streams, KStream and KTable are fundamental data abstractions. They represent different ways to view and process your data.

Understanding their differences is key to building powerful, real-time stream processing applications.

KStream: A Stream of Events

A KStream represents an unbounded, immutable sequence of data records. Think of it like a traditional log or event stream.

  • Each record is treated as a distinct, independent event.
  • Records are processed one by one, in the order they arrive.
  • It never "updates" a previous record; new records are always additions.

It's perfect for handling events like clicks, sensor readings, or financial transactions.

KStream in Action: Filtering

Here's a simple Kafka Streams application that uses a KStream to filter messages. It processes each record individually.

This example will filter a stream of text messages, keeping only those that contain the word "event".

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.KStream;

import java.util.Properties;

public class KStreamFilter {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "kstream-filter-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();
        KStream<String, String> sourceStream = builder.stream("input-topic");

        KStream<String, String> filteredStream = sourceStream
            .filter((key, value) -> value.contains("event"));

        filteredStream.to("output-topic");

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        System.out.println("KStream filter topology created.");
        System.out.println("It filters messages containing 'event'.");
        // In a real application, you would call streams.start();
        // and manage its lifecycle, e.g., using a shutdown hook.
    }
}

KTable: A Dynamic View

A KTable represents a changelog stream, where each record signifies an update or deletion to a row in a table. It's like a database table that's constantly being updated.

  • Each record's value is considered the "latest" value for its key.
  • When a new record arrives with an existing key, it overwrites the previous value.
  • Perfect for maintaining the current state of data.

Think of user profiles, stock prices, or current inventory levels.

KTable in Action: Latest State

This example demonstrates a KTable maintaining the latest value for each key. Imagine tracking the most recent status for various sensors.

When a new message arrives for a sensor, its status in the KTable is updated to the new value.

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.KTable;
import org.apache.kafka.streams.kstream.Materialized;

import java.util.Properties;

public class KTableLatestState {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "ktable-latest-state-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();
        KTable<String, String> latestStatusTable = builder
            .table("sensor-updates", Materialized.as("latest-sensor-status"));

        // The KTable is defined. In a real app, you might print
        // its contents to another topic or join it with a KStream.
        // latestStatusTable.toStream().to("latest-status-output-topic");

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        System.out.println("KTable latest state topology created.");
        System.out.println("It tracks the most recent status for each sensor.");
    }
}

KStream vs. KTable: Key Differences

While both process data, their fundamental nature differs significantly:

  • KStream: Each record is an event. It's a sequence of facts. "Something happened."
  • KTable: Each record is an update. It represents the current state. "This is the current value."

Think of KStream as a transaction log and KTable as the current balance sheet.

When to Use KStream

KStreams are ideal when you need to react to individual events or process data without maintaining a long-term state based on keys.

  • Event logging: Storing every user action.
  • Real-time alerts: Notifying immediately when a specific event occurs.
  • Data enrichment (stateless): Adding information to each event based on its content.
  • Filtering and mapping: Transforming events one by one.

When to Use KTable

KTables are perfect for applications that need to maintain and query the latest state of data, often for aggregations or joins.

  • Current inventory: Tracking stock levels for products.
  • User profiles: Storing the most recent profile details.
  • Aggregations: Counting unique users, summing sales over time (when aggregated results are stored as a KTable).
  • Joining streams with tables: Enriching a KStream with current KTable data.

From Stream to Table: Aggregation

You can transform a KStream into a KTable using stateful operations like aggregation. This converts a series of events into a continuously updated state.

For example, counting occurrences of words from a stream of sentences will result in a KTable where the key is the word and the value is its current 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.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;

import java.util.Arrays;
import java.util.Properties;

public class KStreamToKTable {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();
        KStream<String, String> textLines = builder.stream("text-input");

        KTable<String, Long> wordCounts = textLines
            .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
            .groupBy((key, word) -> word)
            .count(Materialized.as("counts-store"));

        wordCounts.toStream().to("word-counts-output", Produced.with(Serdes.String(), Serdes.Long()));

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        System.out.println("KStream to KTable topology created (Word Count).");
        System.out.println("Counts words from 'text-input' and stores counts in a KTable.");
    }
}

KStream vs. KTable Quiz

Which of the following statements accurately describe a KTable?

KStream & KTable Recap

You've learned the core differences between KStream and KTable in Kafka Streams!

  • KStream: An event stream, processing individual, immutable records.
  • KTable: A changelog stream representing a materialized, updatable view of data.

Choosing the right abstraction is crucial for efficient and meaningful real-time data processing. Next, we'll explore stateless vs. stateful operations in Kafka Streams.

الأسئلة الشائعة

هل درس «مفاهيم KStream وKTable» مجاني؟

نعم — نص درس «مفاهيم KStream وKTable» كامل متاح مجاناً هنا على الويب. لتمرينه بشكل تفاعلي (محرر أكواد مدمج ومدرس ذكاء اصطناعي متاح 24/7) وفتح باقي دورة Apache Kafka & Stream Processing Fundamentals، انتقل إلى CoddyKit PRO. تتضمن دورة Apache Kafka & Stream Processing Fundamentals 4 دروس في المجموع.

ماذا ستتعلم في «مفاهيم KStream وKTable»؟

ميّز بين KStream (تدفق سجل تلو الآخر) وKTable (تدفق سجل تغييرات يمثّل عرضًا ماديًا) في Kafka Streams تتمرن على Apache Kafka & Stream Processing Fundamentals مع أكواد عملية تشغلها مباشرة في المتصفح، ومدرس ذكاء اصطناعي متاح 24/7 يجيب على أسئلتك أثناء عملك.

هل أحتاج إلى خبرة سابقة لأبدأ Apache Kafka & Stream Processing Fundamentals؟

لا تُشترط خبرة سابقة. Apache Kafka & Stream Processing Fundamentals على CoddyKit منظم للمبتدئين حتى المتقدمين، لذا يمكنك البدء من هنا أو من البداية والتقدم بسرعتك الخاصة. هذا هو الدرس 2 من أصل 4.

كم من الوقت يستغرق درس «مفاهيم KStream وKTable»؟

معظم دروس CoddyKit تستغرق حوالي 5–10 دقائق. كل منها موجز وتفاعلي، لذا تحرز تقدماً مستمراً وتستأنف من حيث توقفت عبر الويب والتطبيق.

هل يمكنني كتابة وتشغيل أكواد في درس Apache Kafka & Stream Processing Fundamentals هذا؟

نعم. كل درس في Apache Kafka & Stream Processing Fundamentals يتضمن محرر أكواد مدمج، لذا تكتب وتشغل أكواداً حقيقية مباشرة في متصفحك وتحصل على تعليقات فورية من الذكاء الاصطناعي — بدون إعداد محلي.

جميع الدروس في هذه الدورة

  1. بناء تطبيق Kafka Streams بسيط
  2. مفاهيم KStream وKTable
  3. العمليات عديمة الحالة في مقابل العمليات ذات الحالة
  4. Serdes وتسلسل البيانات في Kafka Streams
← العودة إلى Apache Kafka & Stream Processing Fundamentals