Przetwarzanie strumieni za pomocą KStream i KTable
Dowiedz się, jak używać KStream do obsługi niezmiennych strumieni zdarzeń oraz KTable do tworzenia stanowych, aktualizowalnych widoków danych, wykonując operacje takie jak filtrowanie i mapowanie.
Przetwarzanie strumieni za pomocą KStream i KTable to bezpłatna lekcja Advanced Spring Boot 4: Event-Driven Architecture (Kafka) na CoddyKit. To lekcja 2 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Advanced Spring Boot 4: Event-Driven Architecture (Kafka), a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Advanced Spring Boot 4: Event-Driven Architecture (Kafka) zawiera 4 lekcji w sumie.
KStream i KTable — omówienie
Witamy! W Kafka Streams KStream i KTable to podstawowe narzędzia do przetwarzania danych. Reprezentują różne sposoby postrzegania danych w ruchu.
Można traktować je jak dwie strony tej samej monety — każda z nich sprawdza się w innych zadaniach związanych z przetwarzaniem strumieniowym. Zrozumienie różnic między nimi jest kluczowe przy tworzeniu zaawansowanych aplikacji strumieniowych.
KStream: niezmienne zdarzenia
KStream reprezentuje nieskończoną, niezmienną sekwencję zdarzeń. Każdy rekord w KStream jest samodzielnym faktem, niezależnym zdarzeniem, które wystąpiło w określonym momencie.
- Przypomina dziennik transakcji: po dodaniu zdarzenie nigdy nie jest zmieniane.
- Operacje na KStream tworzą nowe obiekty KStream, pozostawiając oryginał bez zmian.
- Doskonale nadaje się do przetwarzania pojedynczych zdarzeń, takich jak kliknięcia, odczyty z czujników czy wpisy w dziennikach.
Filtrowanie zdarzeń KStream
Jedną z typowych operacji KStream jest filtrowanie. Można wybierać tylko rekordy spełniające określone kryteria, tworząc nowy KStream zawierający wyłącznie istotne zdarzenia.
Oto prosty przykład filtrowania wiadomości zawierających tekst „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!");
}
}Transformowanie wartości KStream
Operacja mapValues przekształca wartość każdego rekordu w KStream, tworząc nowy KStream ze zmodyfikowanymi wartościami. Klucz pozostaje bez zmian.
Jest to przydatne podczas czyszczenia danych, zmiany ich formatu lub wzbogacania informacji bez modyfikowania klucza wiadomości.
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: zmaterializowany widok
KTable reprezentuje strumień zmian, w którym każdy rekord jest aktualizacją dotyczącą określonego klucza. W praktyce jest to zmaterializowany widok tabeli, odzwierciedlający najnowszy stan dla każdego klucza.
- Przypomina tabelę bazy danych: klucze mają przypisane wartości, a nowe rekordy dla danego klucza zastępują poprzednie.
- KTable przechowuje stan, utrzymując w czasie najnowszą wartość dla każdego klucza.
- Doskonale nadaje się do agregowania danych, utrzymywania liczników lub przechowywania profili użytkowników.
Stanowy charakter KTable
Podstawowa idea KTable polega na śledzeniu najnowszej wartości dla każdego unikalnego klucza. Gdy pojawi się nowy rekord z istniejącym kluczem, KTable aktualizuje swój stan wewnętrzny.
Dzięki temu KTable doskonale sprawdza się w sytuacjach, w których interesuje Pana/Panią bieżący stan obiektu, a nie każde pojedyncze zdarzenie, które doprowadziło do jego uzyskania.
KStream do KTable: zliczanie
KStream można przekształcić w KTable, zazwyczaj w celu wykonania agregacji. Typowym przykładem jest zliczanie wystąpień kluczy za pomocą groupByKey().count().
Za każdym razem, gdy nadejdzie wiadomość, licznik dla jej klucza jest aktualizowany, a KTable emituje nową sumę.
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 do agregacji
KTable doskonale nadaje się do ciągłego agregowania danych. Oprócz prostych zliczeń można używać operacji takich jak aggregate do utrzymywania w czasie sum, średnich lub niestandardowych agregatów.
Dzięki temu aplikacja zawsze dysponuje aktualnym podsumowaniem danych dla określonych kluczy.
Kiedy używać którego rozwiązania?
Wybór między KStream a KTable zależy od potrzeb związanych z przetwarzaniem:
- Proszę użyć KStream, gdy trzeba przetwarzać pojedyncze zdarzenia, reagować na każde wystąpienie lub tworzyć potok transformacji niezależny od historycznego stanu. Przykładami są alerty czasu rzeczywistego lub rejestrowanie zdarzeń.
- Proszę użyć KTable, gdy trzeba utrzymywać bieżący stan, agregować dane w czasie lub łączyć je z innymi źródłami na podstawie najnowszej wartości. Przykładami są profile użytkowników, ceny akcji lub zagregowane metryki.
Sprawdzenie: KStream a KTable
Dowiedział(a) się Pan/Pani o KStream i KTable. Sprawdźmy, czy wszystko jest jasne.
Podsumowanie KStream i KTable
Świetna praca! Poznał(a) Pan/Pani najważniejsze różnice między KStream i KTable oraz ich zastosowania.
- KStream obsługuje pojedyncze, niezmienne zdarzenia i doskonale nadaje się do przetwarzania zdarzenie po zdarzeniu.
- KTable utrzymuje zmaterializowany widok, śledząc najnowszy stan każdego klucza, dzięki czemu idealnie nadaje się do agregacji i przetwarzania stanowego.
Te dwa podstawowe typy są fundamentem tworzenia zaawansowanych i elastycznych aplikacji do przetwarzania strumieniowego za pomocą Kafka Streams.
Ucz się Advanced Spring Boot 4: Event-Driven Architecture (Kafka) dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 12
- Lekcje
- 48
Często zadawane pytania
Czy lekcja „Przetwarzanie strumieni za pomocą KStream i KTable” jest bezpłatna?
Tak — pełny tekst „Przetwarzanie strumieni za pomocą KStream i KTable” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Advanced Spring Boot 4: Event-Driven Architecture (Kafka), przejdź na CoddyKit PRO. Kurs Advanced Spring Boot 4: Event-Driven Architecture (Kafka) zawiera 4 lekcji w sumie.
Co nauczysz się w „Przetwarzanie strumieni za pomocą KStream i KTable”?
Dowiedz się, jak używać KStream do obsługi niezmiennych strumieni zdarzeń oraz KTable do tworzenia stanowych, aktualizowalnych widoków danych, wykonując operacje takie jak filtrowanie i mapowanie. Ćwiczysz Advanced Spring Boot 4: Event-Driven Architecture (Kafka) z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Advanced Spring Boot 4: Event-Driven Architecture (Kafka)?
Nie wymagamy żadnego doświadczenia. Advanced Spring Boot 4: Event-Driven Architecture (Kafka) w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 2 z 4.
Ile czasu zajmuje lekcja „Przetwarzanie strumieni za pomocą KStream i KTable”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Advanced Spring Boot 4: Event-Driven Architecture (Kafka)?
Tak. Każda lekcja Advanced Spring Boot 4: Event-Driven Architecture (Kafka) zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Wprowadzenie do Kafka Streams
- Przetwarzanie strumieni za pomocą KStream i KTable
- Tworzenie prostej aplikacji strumieniowej
- Okna czasowe i agregacje stanowe w Kafka Streams