Операции без состояния и с состоянием
Изучите, как Kafka Streams обрабатывает данные с сохранением внутреннего состояния записей и без него.
«Операции без состояния и с состоянием» — бесплатный урок Apache Kafka & Stream Processing Fundamentals на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения 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.
Часто задаваемые вопросы
Урок «Операции без состояния и с состоянием» бесплатный?
Да — полный текст урока «Операции без состояния и с состоянием» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Apache Kafka & Stream Processing Fundamentals, подпишись на CoddyKit PRO. Курс Apache Kafka & Stream Processing Fundamentals содержит 4 уроков всего.
Чему я научусь в уроке «Операции без состояния и с состоянием»?
Изучите, как Kafka Streams обрабатывает данные с сохранением внутреннего состояния записей и без него. Ты практикуешь Apache Kafka & Stream Processing Fundamentals с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать Apache Kafka & Stream Processing Fundamentals?
Предыдущий опыт не требуется. Apache Kafka & Stream Processing Fundamentals на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.
Сколько времени занимает урок «Операции без состояния и с состоянием»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке Apache Kafka & Stream Processing Fundamentals?
Да. Каждый урок Apache Kafka & Stream Processing Fundamentals включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Создание простого приложения Kafka Streams
- Понятия KStream и KTable
- Операции без состояния и с состоянием
- Serdes и сериализация данных в Kafka Streams