Apache Kafka ja suoratoistonkäsittelyn perusteet · Oppitunti

KStream- ja KTable-käsitteet

Erota toisistaan KStream (tietue kerrallaan etenevä stream) ja KTable (materialisoitua näkymää edustava muutosloki-stream) Kafka Streamsissa.

Oppitunti 2/411 vaihetta

KStream- ja KTable-käsitteet on ilmainen Apache Kafka ja suoratoistonkäsittelyn perusteet-oppitunti CoddyKitissä. Tämä on oppitunti 2/4. Voit lukea tästä oppimispolusta kokonaan mitkä tahansa 3 oppituntia ilmaiseksi — sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä käytännön harjoittelun sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. Oppitunti kuuluu Apache Kafka ja suoratoistonkäsittelyn perusteet-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Apache Kafka ja suoratoistonkäsittelyn perusteet-kurssilla on yhteensä 4 oppituntia.

KStream ja KTable – perusteet

Kafka Streamsissa KStream ja KTable ovat keskeisiä datan abstraktioita. Ne kuvaavat erilaisia tapoja tarkastella ja käsitellä dataa.

Niiden erojen ymmärtäminen on tärkeää tehokkaiden reaaliaikaisten suorakäsittelysovellusten rakentamisessa.

KStream: tapahtumien virta

KStream edustaa rajatonta, muuttumatonta datatietueiden sarjaa. Ajatelkaa sitä perinteisen lokin tai tapahtumavirran kaltaisena.

  • Jokaista tietuetta käsitellään erillisenä, itsenäisenä tapahtumana.
  • Tietueet käsitellään yksi kerrallaan saapumisjärjestyksessä.
  • Se ei koskaan "päivitä" aiempaa tietuetta, vaan uudet tietueet ovat aina lisäyksiä.

Se sopii erinomaisesti esimerkiksi klikkausten, anturilukemien tai rahansiirtojen käsittelyyn.

KStream käytännössä: suodatus

Tässä on yksinkertainen Kafka Streams -sovellus, joka käyttää KStream-virtaa viestien suodattamiseen. Se käsittelee jokaisen tietueen erikseen.

Tässä esimerkissä tekstiviestien virrasta säilytetään vain viestit, jotka sisältävät sanan "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: dynaaminen näkymä

KTable edustaa muutoslokia, jossa jokainen tietue ilmaisee taulun rivin päivityksen tai poiston. Se muistuttaa jatkuvasti päivittyvää tietokantataulua.

  • Kunkin tietueen arvoa pidetään sen avaimen "uusimpana" arvona.
  • Kun olemassa olevalla avaimella saapuu uusi tietue, se korvaa aiemman arvon.
  • Se sopii erinomaisesti datan nykyisen tilan ylläpitämiseen.

Ajatelkaa käyttäjäprofiileja, osakekursseja tai varaston nykyisiä määriä.

KTable käytännössä: uusin tila

Tässä esimerkissä KTable ylläpitää kunkin avaimen uusinta arvoa. Kuvitelkaa, että seuraatte eri antureiden viimeisintä tilaa.

Kun anturilta saapuu uusi viesti, sen tila KTable-taulussa päivitetään uudeksi arvoksi.

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 ja KTable: keskeiset erot

Vaikka molemmat käsittelevät dataa, niiden perusluonne eroaa merkittävästi:

  • KStream: Jokainen tietue on tapahtuma. Se on tosiasioiden sarja: "Jotain tapahtui."
  • KTable: Jokainen tietue on päivitys. Se edustaa nykyistä tilaa: "Tämä on nykyinen arvo."

Ajatelkaa KStreamia tapahtumalokina ja KTablea nykyisenä taseena.

Milloin KStreamia kannattaa käyttää

KStreamit sopivat erinomaisesti tilanteisiin, joissa yksittäisiin tapahtumiin on reagoitava tai dataa on käsiteltävä ylläpitämättä avainten perusteella pitkän aikavälin tilaa.

  • Tapahtumien kirjaaminen: Jokaisen käyttäjän toiminnon tallentaminen.
  • Reaaliaikaiset hälytykset: Ilmoituksen lähettäminen heti tietyn tapahtuman tapahtuessa.
  • Datan täydentäminen (tilaton): Tietojen lisääminen kuhunkin tapahtumaan sen sisällön perusteella.
  • Suodatus ja kuvaus: Tapahtumien muuntaminen yksi kerrallaan.

Milloin KTablea kannattaa käyttää

KTablet sopivat erinomaisesti sovelluksiin, joiden on ylläpidettävä ja kyseltävä tietojen uusinta tilaa, usein koostamista tai liitoksia varten.

  • Nykyinen varastotilanne: Tuotteiden varastomäärien seuranta.
  • Käyttäjäprofiilit: Profiilitietojen uusimman version tallentaminen.
  • Koostaminen: Ainutkertaisten käyttäjien laskeminen ja myynnin summaaminen ajan kuluessa (kun koostetut tulokset tallennetaan KTableen).
  • Virtojen liittäminen tauluihin: KStreamin rikastaminen KTablen nykyisillä tiedoilla.

Virta tauluksi: koostaminen

Voit muuntaa KStreamin KTableksi tilallisten toimintojen, kuten koostamisen, avulla. Näin tapahtumasarja muunnetaan jatkuvasti päivittyväksi tilaksi.

Jos esimerkiksi lasketaan sanojen esiintymät lausevirrasta, tuloksena on KTable, jossa avaimena on sana ja arvona sen nykyinen esiintymismäärä.

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

KStreamin ja KTablen vertailu – tietovisa

Mitkä seuraavista väitteistä kuvaavat KTablea oikein?

KStreamin ja KTablen kertaus

Olet oppinut KStreamin ja KTablen keskeiset erot Kafka Streamsissa!

  • KStream: Tapahtumavirta, joka käsittelee yksittäisiä muuttumattomia tietueita.
  • KTable: Muutoslokia sisältävä virta, joka edustaa tietojen materialisoitua ja päivitettävää näkymää.

Oikean abstraktion valitseminen on ratkaisevan tärkeää tehokkaan ja merkityksellisen reaaliaikaisen tietojenkäsittelyn kannalta. Seuraavaksi tutustumme tilattomiin ja tilallisiin toimintoihin Kafka Streamsissa.

Aloita maksutta

Opi Apache Kafka ja suoratoistonkäsittelyn perusteet tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
12
Oppitunnit
48

Usein kysytyt kysymykset

Onko oppitunti ”KStream- ja KTable-käsitteet” ilmainen?

Kyllä — voit lukea täällä verkossa kokonaan ilmaiseksi mitkä tahansa Apache Kafka ja suoratoistonkäsittelyn perusteet-oppimispolun 3 oppituntia, myös oppitunnin “KStream- ja KTable-käsitteet”. Sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä interaktiiviset harjoitukset sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. Apache Kafka ja suoratoistonkäsittelyn perusteet-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”KStream- ja KTable-käsitteet”?

Erota toisistaan KStream (tietue kerrallaan etenevä stream) ja KTable (materialisoitua näkymää edustava muutosloki-stream) Kafka Streamsissa. Harjoittelet Apache Kafka ja suoratoistonkäsittelyn perusteet-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni Apache Kafka ja suoratoistonkäsittelyn perusteet-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin Apache Kafka ja suoratoistonkäsittelyn perusteet-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 2/4.

Kuinka kauan ”KStream- ja KTable-käsitteet”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä Apache Kafka ja suoratoistonkäsittelyn perusteet-oppitunnilla?

Kyllä. Jokainen Apache Kafka ja suoratoistonkäsittelyn perusteet-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. Yksinkertaisen Kafka Streams -sovelluksen rakentaminen
  2. KStream- ja KTable-käsitteet
  3. Tilattomat ja tilalliset operaatiot
  4. Serdes ja datan serialisointi Kafka Streamsissa
← Takaisin: Apache Kafka ja suoratoistonkäsittelyn perusteet