Grunderna i Apache Kafka och strömbearbetning · Lektion

KStream- och KTable-koncept

Särskilj KStream (en post-för-post-ström) från KTable (en ändringsloggström som representerar en materialiserad vy) i Kafka Streams.

Lektion 2 av 411 steg

KStream- och KTable-koncept är en gratis lektion i Grunderna i Apache Kafka och strömbearbetning på CoddyKit. Detta är lektion 2 av 4. Du kan läsa vilka 3 lektioner som helst i den här lärvägen kostnadsfritt i sin helhet – därefter låser CoddyKit PRO upp alla lektioner, plus praktisk övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Den ingår i lärvägen för Grunderna i Apache Kafka och strömbearbetning, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i Grunderna i Apache Kafka och strömbearbetning innehåller totalt 4 lektioner.

KStream och KTable förklarade

I Kafka Streams är KStream och KTable grundläggande dataabstraktioner. De representerar olika sätt att visa och bearbeta dina data.

Det är viktigt att förstå skillnaderna mellan dem för att kunna bygga kraftfulla strömbehandlingsapplikationer i realtid.

KStream: En ström av händelser

En KStream representerar en obegränsad, oföränderlig sekvens av dataposter. Se den som en traditionell logg eller händelseström.

  • Varje post behandlas som en separat och oberoende händelse.
  • Poster bearbetas en i taget i den ordning de anländer.
  • En tidigare post ”uppdateras” aldrig; nya poster läggs alltid till.

Den passar perfekt för händelser som klick, sensoravläsningar eller finansiella transaktioner.

KStream i praktiken: Filtrering

Här är en enkel Kafka Streams-applikation som använder en KStream för att filtrera meddelanden. Den bearbetar varje post separat.

Exemplet filtrerar en ström av textmeddelanden och behåller endast de meddelanden som innehåller ordet ”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: En dynamisk vy

En KTable representerar en ändringsloggström, där varje post anger en uppdatering eller borttagning av en rad i en tabell. Den liknar en databastabell som uppdateras kontinuerligt.

  • Värdet i varje post betraktas som det ”senaste” värdet för dess nyckel.
  • När en ny post anländer med en befintlig nyckel skriver den över det tidigare värdet.
  • Perfekt för att underhålla det aktuella tillståndet för data.

Tänk på användarprofiler, aktiekurser eller aktuella lagernivåer.

KTable i praktiken: Senaste tillståndet

Det här exemplet visar hur en KTable underhåller det senaste värdet för varje nyckel. Föreställ dig att du följer den senaste statusen för olika sensorer.

När ett nytt meddelande anländer för en sensor uppdateras dess status i KTable till det nya värdet.

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 jämfört med KTable: Viktiga skillnader

Även om båda bearbetar data skiljer sig deras grundläggande natur avsevärt:

  • KStream: Varje post är en händelse. Det är en sekvens av fakta: ”Något hände.”
  • KTable: Varje post är en uppdatering. Den representerar det aktuella tillståndet: ”Det här är det aktuella värdet.”

Tänk på KStream som en transaktionslogg och KTable som den aktuella balansräkningen.

När ska KStream användas

KStreams passar perfekt när du behöver reagera på enskilda händelser eller bearbeta data utan att upprätthålla ett långvarigt tillstånd baserat på nycklar.

  • Händelseloggar: Lagra varje användaråtgärd.
  • Realtidsaviseringar: Skicka en avisering omedelbart när en viss händelse inträffar.
  • Databerikning (tillståndslös): Lägg till information i varje händelse baserat på dess innehåll.
  • Filtrering och mappning: Transformera händelser en i taget.

När ska KTable användas

KTables passar perfekt för applikationer som behöver underhålla och fråga efter det senaste tillståndet för data, ofta för aggregeringar eller join-operationer.

  • Aktuellt lager: Spåra lagernivåer för produkter.
  • Användarprofiler: Lagra de senaste profildetaljerna.
  • Aggregeringar: Räkna unika användare och summera försäljning över tid (när de aggregerade resultaten lagras som en KTable).
  • Joina strömmar med tabeller: Berika en KStream med aktuella KTable-data.

Från ström till tabell: aggregering

Ni kan omvandla en KStream till en KTable med hjälp av tillståndsbaserade operationer som aggregering. Detta omvandlar en serie händelser till ett tillstånd som uppdateras kontinuerligt.

Om ni till exempel räknar förekomsten av ord i en ström av meningar blir resultatet en KTable där nyckeln är ordet och värdet är dess aktuella antal.

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.");
    }
}

Frågesport: KStream kontra KTable

Vilka av följande påståenden beskriver korrekt en KTable?

Sammanfattning av KStream och KTable

Ni har lärt er de viktigaste skillnaderna mellan KStream och KTable i Kafka Streams!

  • KStream: En händelseström som bearbetar enskilda, oföränderliga poster.
  • KTable: En ändringsloggström som representerar en materialiserad, uppdateringsbar vy av data.

Att välja rätt abstraktion är avgörande för en effektiv och meningsfull realtidsbearbetning av data. Nästa steg är att utforska tillståndslösa och tillståndsbaserade operationer i Kafka Streams.

Gratis att börja

Lär dig Grunderna i Apache Kafka och strömbearbetning med en AI-lärare – gratis

Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.

Kurser
12
Lektioner
48

Vanliga frågor

Är lektionen ”KStream- och KTable-koncept” gratis?

Ja – du kan läsa vilka 3 lektioner som helst i lärvägen Grunderna i Apache Kafka och strömbearbetning, inklusive ”KStream- och KTable-koncept”, kostnadsfritt i sin helhet här på webben. Därefter låser CoddyKit PRO upp alla lektioner, plus interaktiv övning med en inbyggd kodredigerare och en AI-lärare dygnet runt. Kursen i Grunderna i Apache Kafka och strömbearbetning innehåller totalt 4 lektioner.

Vad lär jag mig i ”KStream- och KTable-koncept”?

Särskilj KStream (en post-för-post-ström) från KTable (en ändringsloggström som representerar en materialiserad vy) i Kafka Streams. Ni övar på Grunderna i Apache Kafka och strömbearbetning med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.

Behöver jag någon erfarenhet för att börja lära mig Grunderna i Apache Kafka och strömbearbetning?

Du behöver inga förkunskaper. Utbildningen i Grunderna i Apache Kafka och strömbearbetning på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 2 av 4.

Hur lång tid tar lektionen ”KStream- och KTable-koncept”?

De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.

Kan jag skriva och köra kod i den här Grunderna i Apache Kafka och strömbearbetning-lektionen?

Ja. Varje Grunderna i Apache Kafka och strömbearbetning-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.

Alla lektioner i den här kursen

  1. Bygga en enkel Kafka Streams-app
  2. KStream- och KTable-koncept
  3. Tillståndslösa och tillståndsbaserade operationer
  4. Serdes och dataserialisering i Kafka Streams
← Tillbaka till Grunderna i Apache Kafka och strömbearbetning