0Pricing
Apache Kafka & Stream Processing Fundamentals · 课时

KStream 和 KTable 概念

区分 Kafka Streams 中的 KStream(逐条记录的流)与 KTable(表示物化视图的变更日志流)

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

「KStream 和 KTable 概念」这节课中我会学到什么?

区分 Kafka Streams 中的 KStream(逐条记录的流)与 KTable(表示物化视图的变更日志流) 你通过在浏览器中直接运行的动手代码来练习 Apache Kafka & Stream Processing Fundamentals,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

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

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

「KStream 和 KTable 概念」课时需要多长时间?

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

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

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

此课程中的所有课时

  1. 构建简单的 Kafka Streams 应用
  2. KStream 和 KTable 概念
  3. 无状态操作与有状态操作
  4. Kafka Streams 中的 Serdes 与数据序列化
← 返回 Apache Kafka & Stream Processing Fundamentals