Stream-käsittely KStreamilla ja KTablella
Opi käyttämään KStreamia muuttumattomille tapahtumavirroille ja KTablea tilallisille, päivitettäville datanäkymille sekä suorittamaan toimintoja, kuten suodatusta ja muunnettaessa.
Stream-käsittely KStreamilla ja KTablella on ilmainen Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-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 Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-kurssilla on yhteensä 4 oppituntia.
KStream ja KTable esittelyssä
Tervetuloa! Kafka Streamsissa KStream ja KTable ovat ensisijaiset työkalunne tietojen käsittelyyn. Ne kuvaavat datan liikkeessä olevia eri näkymiä.
Ajatelkaa niitä saman kolikon kahtena puolena, joista kumpikin soveltuu erilaisiin virtakäsittelytehtäviin. Niiden erojen ymmärtäminen on tärkeää tehokkaiden virtasovellusten rakentamisessa.
KStream: muuttumattomat tapahtumat
KStream edustaa ääretöntä ja muuttumatonta tapahtumien sarjaa. Jokainen KStreamin tietue on itsenäinen tosiasia eli erillinen tapahtuma, joka tapahtui tiettynä ajanhetkenä.
- Se muistuttaa tapahtumalokia: kun tapahtuma on lisätty, sitä ei enää muuteta.
- KStreamiin kohdistetut operaatiot tuottavat uusia KStream-virtoja alkuperäisen säilyessä muuttumattomana.
- Se sopii erinomaisesti yksittäisten tapahtumien, kuten napsautusten, anturilukemien tai lokimerkintöjen, käsittelyyn.
KStream-tapahtumien suodatus
Yksi yleinen KStream-operaatio on suodatus. Voitte säilyttää valikoivasti tietueet, jotka vastaavat tiettyjä ehtoja, ja luoda uuden KStreamin, joka sisältää vain olennaiset tapahtumat.
Tässä on yksinkertainen esimerkki viestien suodattamisesta siten, että mukaan otetaan viestit, jotka sisältävät merkkijonon '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!");
}
}KStream-arvojen muuntaminen
mapValues-operaatio muuntaa KStreamin jokaisen tietueen arvon ja tuottaa uuden KStreamin, joka sisältää muunnetut arvot. Avain säilyy muuttumattomana.
Tämä on hyödyllistä tietojen puhdistamiseen, muotojen muuttamiseen tai tietojen rikastamiseen ilman, että viestin avainta muutetaan.
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: materialisoitu näkymä
KTable edustaa muutoslokia, jossa jokainen tietue on tiettyyn avaimeen kohdistuva päivitys. Se on käytännössä taulun materialisoitu näkymä, joka heijastaa kunkin avaimen uusinta tilaa.
- Se muistuttaa tietokantataulua: avaimiin liittyy arvoja, ja avaimen uudet tietueet korvaavat aiemmat arvot.
- KTable on tilallinen ja ylläpitää kunkin avaimen uusinta arvoa ajan kuluessa.
- Se sopii erinomaisesti tietojen koostamiseen, lukumäärien ylläpitämiseen tai käyttäjäprofiilien tallentamiseen.
KTablen tilallinen luonne
KTablen keskeinen ajatus on, että se seuraa kunkin yksilöllisen avaimen uusinta arvoa. Kun olemassa olevan avaimen sisältävä uusi tietue saapuu, KTable päivittää sisäisen tilansa.
Tämän ansiosta KTable sopii tilanteisiin, joissa olette kiinnostuneita entiteetin nykyisestä tilasta kaikkien siihen johtaneiden yksittäisten tapahtumien sijaan.
KStreamistä KTableksi: laskeminen
Voitte muuntaa KStreamin KTableksi tyypillisesti koostamista varten. Yleinen esimerkki on avainten esiintymiskertojen laskeminen käyttämällä groupByKey().count()-kutsua.
Aina kun viesti saapuu, sen avaimen lukumäärä päivitetään, ja KTable lähettää uuden kokonaismäärän.
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 koostamiseen
KTablet sopivat erinomaisesti jatkuvaan koostamiseen. Yksinkertaisten lukumäärien lisäksi voitte käyttää esimerkiksi aggregate-operaatiota summien, keskiarvojen tai mukautettujen koostetietojen ylläpitämiseen ajan kuluessa.
Näin sovelluksellanne on aina ajantasainen yhteenveto tiettyjen avainten tiedoista.
Milloin mitäkin käytetään?
Valinta KStreamin ja KTablen välillä riippuu käsittelytarpeistanne:
- Käyttäkää KStreamia, kun haluatte käsitellä yksittäisiä tapahtumia, reagoida jokaiseen esiintymään tai rakentaa muunnosputken, joka ei riipu aiemmasta tilasta. Esimerkkejä ovat reaaliaikaiset hälytykset ja tapahtumien kirjaaminen.
- Käyttäkää KTablea, kun haluatte ylläpitää nykyistä tilaa, koostaa tietoja ajan kuluessa tai yhdistää tietoja muihin lähteisiin uusimman arvon perusteella. Esimerkkejä ovat käyttäjäprofiilit, osakekurssit ja koostetut mittarit.
KStreamin ja KTablen testi
Olette oppineet KStreamista ja KTablesta. Testataanpa ymmärrystänne.
KStreamin ja KTablen kertaus
Hienoa työtä! Tutustuitte KStreamin ja KTablen keskeisiin eroihin ja käyttötapoihin.
- KStream käsittelee yksittäisiä, muuttumattomia tapahtumia ja sopii erinomaisesti tapahtumien käsittelyyn yksi kerrallaan.
- KTable ylläpitää materialisoitua näkymää ja seuraa kunkin avaimen uusinta tilaa, joten se sopii koostamiseen ja tilalliseen käsittelyyn.
Nämä kaksi perusrakennetta muodostavat perustan tehokkaiden ja joustavien virtakäsittelysovellusten rakentamiselle Kafka Streamsilla.
Opi Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka) 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 ”Stream-käsittely KStreamilla ja KTablella” ilmainen?
Kyllä — voit lukea täällä verkossa kokonaan ilmaiseksi mitkä tahansa Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-oppimispolun 3 oppituntia, myös oppitunnin “Stream-käsittely KStreamilla ja KTablella”. Sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä interaktiiviset harjoitukset sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-kurssilla on yhteensä 4 oppituntia.
Mitä opin oppitunnilla ”Stream-käsittely KStreamilla ja KTablella”?
Opi käyttämään KStreamia muuttumattomille tapahtumavirroille ja KTablea tilallisille, päivitettäville datanäkymille sekä suorittamaan toimintoja, kuten suodatusta ja muunnettaessa. Harjoittelet Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.
Tarvitsenko kokemusta aloittaakseni Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-opiskelun?
Aiempi kokemus ei ole tarpeen. CoddyKitin Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 2/4.
Kuinka kauan ”Stream-käsittely KStreamilla ja KTablella”-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ä Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-oppitunnilla?
Kyllä. Jokainen Spring Boot 4:n edistyneet aiheet: tapahtumavetoinen arkkitehtuuri (Kafka)-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
- Johdatus Kafka Streamsiin
- Stream-käsittely KStreamilla ja KTablella
- Yksinkertaisen stream-sovelluksen rakentaminen
- Ikkunointi ja tilalliset aggregoinnit Kafka Streamsissa