无状态操作与有状态操作
了解 Kafka Streams 如何在保留或不保留记录间内部状态的情况下处理数据
无状态操作与有状态操作 是 CoddyKit 上的免费 Apache Kafka & Stream Processing Fundamentals 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Apache Kafka & Stream Processing Fundamentals 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Apache Kafka & Stream Processing Fundamentals 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
Stateless vs. Stateful Streams
In Kafka Streams, how you process data falls into two main categories: stateless and stateful operations.
Understanding this distinction is crucial for building efficient and correct real-time data pipelines. It impacts how your application remembers (or forgets) past events.
What are Stateless Operations?
Stateless operations process each incoming record independently. They don't remember any past records or maintain an internal state across events.
Think of it like a simple function: you give it an input, and it produces an output, without needing any memory of previous inputs.
Stateless Example: Map & Filter
Common stateless operations include map, filter, flatMap, and peek. They transform or filter records one by one.
Here's a simple Kafka Streams app using mapValues to convert all incoming message values to uppercase:
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.Topology;
import org.apache.kafka.streams.kstream.KStream;
import java.util.Properties;
public class StatelessMapApp {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG,
"stateless-map-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> source = builder.stream("input-topic");
// Stateless operation: mapValues
KStream<String, String> upperCaseStream =
source.mapValues(value -> value.toUpperCase());
upperCaseStream.to("output-topic");
Topology topology = builder.build();
KafkaStreams streams = new KafkaStreams(topology, props);
streams.start();
// In a real app, add a shutdown hook.
}
}When to Use Stateless Operations
Stateless operations are ideal for:
- Simple Transformations: Changing data format, type conversion.
- Filtering: Removing unwanted records.
- Data Cleansing: Basic sanitization of individual records.
They are generally simpler to implement and have less overhead because no state needs to be managed.
What are Stateful Operations?
Stateful operations are those that need to remember past records or combine information across multiple records to produce a result.
They maintain an internal state, which is stored locally within the Kafka Streams application instance. This state allows them to perform aggregations, joins, and windowing.
Stateful Example: Counting Events
A classic example of a stateful operation is count(). To count events per key, the application must remember previous counts for each key.
This operation transforms a KStream into a KTable, which represents a changelog of aggregated results.
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.Topology;
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 StatefulCountApp {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG,
"stateful-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> source = builder.stream("input-topic");
// Stateful operation: count by key
KTable<String, Long> countsTable = source
.groupByKey()
.count(Materialized.as("counts-store")); // A named state store
countsTable.toStream().to("output-topic");
Topology topology = builder.build();
KafkaStreams streams = new KafkaStreams(topology, props);
streams.start();
// In a real app, add a shutdown hook.
}
}More Stateful Operations
Besides count(), other common stateful operations include:
- Aggregations:
reduce(),aggregate()(e.g., calculating sums, averages). - Joins: Combining data from two streams or a stream and a table based on a common key.
- Windowing: Grouping records that fall within a defined time frame (e.g., 5-minute window).
These operations all rely on maintaining state to function correctly.
Kafka Streams State Stores
Kafka Streams manages state using internal state stores. These are typically backed by a local key-value store like RocksDB.
For fault tolerance, Kafka Streams also uses internal Kafka topics (called changelog topics) to continuously back up the state store. If an application instance fails, its state can be restored from the changelog topic by a new instance.
Quick Check: Identify Operations
Which of these Kafka Streams operations is considered stateless?
Recap: Stateless vs. Stateful
We've explored the key differences between stateless and stateful operations in Kafka Streams:
- Stateless: Processes records individually, no memory of past events. Ideal for simple transformations and filtering.
- Stateful: Requires internal memory (state stores) to combine or remember information across records. Essential for aggregations, joins, and windowing.
Choosing the right type of operation is fundamental to designing robust and efficient stream processing applications.
常见问题解答
「无状态操作与有状态操作」课时是免费的吗?
是的 — 「无状态操作与有状态操作」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Apache Kafka & Stream Processing Fundamentals 课程的其余内容,请升级到 CoddyKit PRO。 Apache Kafka & Stream Processing Fundamentals 课程共包含 4 节课。
「无状态操作与有状态操作」这节课中我会学到什么?
了解 Kafka Streams 如何在保留或不保留记录间内部状态的情况下处理数据 你通过在浏览器中直接运行的动手代码来练习 Apache Kafka & Stream Processing Fundamentals,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Apache Kafka & Stream Processing Fundamentals 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Apache Kafka & Stream Processing Fundamentals 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「无状态操作与有状态操作」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Apache Kafka & Stream Processing Fundamentals 课中编写并运行代码吗?
能。每节 Apache Kafka & Stream Processing Fundamentals 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。