使用 KStream 和 KTable 进行流处理
学习使用 KStream 处理不可变事件流,使用 KTable 构建有状态且可更新的数据视图,并执行过滤、映射等操作
使用 KStream 和 KTable 进行流处理 是 CoddyKit 上的免费 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
KStream & KTable Unveiled
Welcome! In Kafka Streams, KStream and KTable are your primary tools for processing data. They represent different views of your data in motion.
Think of them as two sides of the same coin, each suited for distinct stream processing tasks. Understanding their differences is key to building powerful stream applications.
KStream: Immutable Events
A KStream represents an infinite, immutable sequence of events. Each record in a KStream is a self-contained fact, an independent event that happened at a specific point in time.
- It's like a transaction log: once an event is added, it's never changed.
- Operations on a KStream produce new KStreams, leaving the original untouched.
- It's ideal for processing individual events like clicks, sensor readings, or log entries.
Filtering KStream Events
One common KStream operation is filtering. You can selectively keep records that match certain criteria, creating a new KStream with only the relevant events.
Here's a simple example filtering messages that contain 'hello'.
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 Main {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "filter-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_DEFAULT, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_DEFAULT, 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("hello")
);
filteredStream.to("output-topic");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
// In a real app, you'd start and manage this lifecycle:
// streams.start();
// Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
System.out.println("KStream filter setup complete. Send 'hello world' to input-topic!");
}
}Transforming KStream Values
The mapValues operation transforms the value of each record in a KStream, producing a new KStream with the modified values. The key remains unchanged.
This is useful for cleaning data, changing formats, or enriching information without altering the message's key.
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 Main {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "mapvalues-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_DEFAULT, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_DEFAULT, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> sourceStream = builder.stream("input-topic");
KStream<String, String> uppercasedStream = sourceStream.mapValues(
value -> value.toUpperCase()
);
uppercasedStream.to("output-topic");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
System.out.println("KStream mapValues setup complete. Send 'test' to input-topic!");
}
}KTable: A Materialized View
A KTable represents a changelog stream, where each record is an update to a specific key. It's essentially a materialized view of a table, reflecting the latest state for each key.
- It's like a database table: keys have associated values, and new records for a key overwrite previous ones.
- KTable is stateful, maintaining the latest value for each key over time.
- It's perfect for aggregating data, maintaining counts, or storing user profiles.
KTable's Stateful Nature
The core idea behind a KTable is that it keeps track of the latest value for each unique key. When a new record with an existing key arrives, the KTable updates its internal state.
This makes KTable ideal for scenarios where you care about the current state of an entity, rather than every single event that led to that state.
KStream to KTable: Counting
You can transform a KStream into a KTable, typically to perform aggregations. A common example is counting occurrences of keys using groupByKey().count().
Each time a message arrives, the count for its key is updated, and the KTable emits the new total.
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 java.util.Properties;
public class Main {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "kstream-to-ktable-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_DEFAULT, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_DEFAULT, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> sourceStream = builder.stream("input-topic");
KTable<String, Long> wordCounts = sourceStream
.groupByKey() // Group by the existing key
.count(Materialized.as("counts-store")); // Count occurrences, store in state
wordCounts.toStream().to("output-topic"); // Convert back to stream to send out
KafkaStreams streams = new KafkaStreams(builder.build(), props);
System.out.println("KStream to KTable count setup. Send 'word' with key 'A' to input-topic!");
}
}KTable for Aggregation
KTables are excellent for continuous aggregation. Beyond simple counts, you can use operations like aggregate to maintain sums, averages, or custom aggregates over time.
This allows your application to always have an up-to-date summary of data for specific keys.
When to Use Which?
The choice between KStream and KTable depends on your processing needs:
- Use KStream when you need to process individual events, react to every occurrence, or build a pipeline of transformations that don't depend on historical state. Think real-time alerts or event logging.
- Use KTable when you need to maintain a current state, aggregate data over time, or join with other data sources based on the latest value. Think user profiles, stock prices, or aggregated metrics.
KStream vs. KTable Check
You've learned about KStream and KTable. Let's test your understanding.
KStream & KTable Recap
Great job! You've explored the core differences and uses of KStream and KTable.
- KStream handles individual, immutable events, perfect for event-by-event processing.
- KTable maintains a materialized view, tracking the latest state for each key, ideal for aggregations and stateful processing.
These two primitives are the foundation for building powerful and flexible stream processing applications with Kafka Streams.
常见问题解答
「使用 KStream 和 KTable 进行流处理」课时是免费的吗?
是的 — 「使用 KStream 和 KTable 进行流处理」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程的其余内容,请升级到 CoddyKit PRO。 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程共包含 4 节课。
「使用 KStream 和 KTable 进行流处理」这节课中我会学到什么?
学习使用 KStream 处理不可变事件流,使用 KTable 构建有状态且可更新的数据视图,并执行过滤、映射等操作 你通过在浏览器中直接运行的动手代码来练习 Advanced Spring Boot 4: Event-Driven Architecture (Kafka),全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。
「使用 KStream 和 KTable 进行流处理」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课中编写并运行代码吗?
能。每节 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- Kafka Streams 简介
- 使用 KStream 和 KTable 进行流处理
- 构建简单的流应用
- Kafka Streams 中的窗口与有状态聚合