Pemprosesan Strim dengan KStream & KTable
Pelajari cara menggunakan KStream untuk strim peristiwa tidak boleh ubah dan KTable untuk paparan data berkeadaan yang boleh dikemas kini, termasuk operasi seperti penapisan dan pemetaan.
Pemprosesan Strim dengan KStream & KTable ialah pelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) percuma di CoddyKit. Ini ialah pelajaran 2 daripada 4. Sebanyak 3 pelajaran dalam laluan pembelajaran ini boleh dibaca sepenuhnya secara percuma — selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan praktikal dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Pelajaran ini merupakan sebahagian daripada laluan pembelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka), dan kemajuan anda disegerakkan merentas web serta aplikasi CoddyKit. Kursus Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) merangkumi sejumlah 4 pelajaran.
KStream & KTable Diperkenalkan
Selamat datang! Dalam Kafka Streams, KStream dan KTable ialah alat utama anda untuk memproses data. Kedua-duanya mewakili pandangan yang berbeza terhadap data yang sedang bergerak.
Anggaplah kedua-duanya sebagai dua sisi syiling yang sama, dengan setiap satu sesuai untuk tugas pemprosesan strim yang berbeza. Memahami perbezaannya adalah penting untuk membina aplikasi strim yang berkuasa.
KStream: Peristiwa Tidak Berubah
KStream mewakili jujukan peristiwa yang tidak terhingga dan tidak berubah. Setiap rekod dalam KStream ialah fakta yang lengkap, iaitu peristiwa tersendiri yang berlaku pada titik masa tertentu.
- Ia seperti log transaksi: setelah sesuatu peristiwa ditambah, ia tidak akan diubah.
- Operasi pada KStream menghasilkan KStream baharu dan membiarkan KStream asal tanpa perubahan.
- Ia sesuai untuk memproses peristiwa individu seperti klik, bacaan sensor atau entri log.
Menapis Peristiwa KStream
Salah satu operasi KStream yang biasa ialah penapisan. Anda boleh menyimpan secara terpilih rekod yang sepadan dengan kriteria tertentu, lalu menghasilkan KStream baharu yang hanya mengandungi peristiwa berkaitan.
Berikut ialah contoh mudah untuk menapis mesej yang mengandungi '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!");
}
}Mengubah Nilai KStream
Operasi mapValues mengubah nilai setiap rekod dalam KStream dan menghasilkan KStream baharu dengan nilai yang telah diubah. Kekuncinya kekal tidak berubah.
Ini berguna untuk membersihkan data, menukar format atau memperkaya maklumat tanpa mengubah kunci mesej.
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: Pandangan Terwujud
KTable mewakili strim log perubahan, yang setiap rekodnya ialah kemas kini kepada kunci tertentu. Pada asasnya, ia ialah pandangan terwujud bagi sebuah jadual yang mencerminkan keadaan terkini untuk setiap kunci.
- Ia seperti jadual pangkalan data: kunci mempunyai nilai yang berkaitan dan rekod baharu bagi sesuatu kunci akan menggantikan rekod sebelumnya.
- KTable mengekalkan keadaan, dengan menyimpan nilai terkini bagi setiap kunci dari semasa ke semasa.
- Ia sangat sesuai untuk mengagregatkan data, mengekalkan kiraan atau menyimpan profil pengguna.
Sifat KTable yang Mengekalkan Keadaan
Idea teras di sebalik KTable ialah ia menjejaki nilai terkini bagi setiap kunci unik. Apabila rekod baharu dengan kunci sedia ada tiba, KTable mengemas kini keadaan dalamannya.
Hal ini menjadikan KTable sesuai untuk situasi apabila anda mengambil berat tentang keadaan semasa sesuatu entiti, bukannya setiap peristiwa yang membawa kepada keadaan tersebut.
KStream kepada KTable: Mengira
Anda boleh mengubah KStream menjadi KTable, biasanya untuk melakukan pengagregatan. Contoh yang biasa ialah mengira kejadian kunci menggunakan groupByKey().count().
Setiap kali mesej tiba, kiraan bagi kuncinya dikemas kini dan KTable mengeluarkan jumlah baharu.
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 untuk Pengagregatan
KTable sangat baik untuk pengagregatan berterusan. Selain kiraan mudah, anda boleh menggunakan operasi seperti aggregate untuk mengekalkan jumlah, purata atau agregat tersuai dari semasa ke semasa.
Ini membolehkan aplikasi anda sentiasa mempunyai ringkasan data terkini bagi kunci tertentu.
Bila Perlu Menggunakan Yang Mana?
Pilihan antara KStream dengan KTable bergantung pada keperluan pemprosesan anda:
- Gunakan KStream apabila anda perlu memproses peristiwa individu, bertindak balas terhadap setiap kejadian atau membina saluran transformasi yang tidak bergantung pada keadaan sejarah. Fikirkan tentang amaran masa nyata atau pencatatan peristiwa.
- Gunakan KTable apabila anda perlu mengekalkan keadaan semasa, mengagregatkan data dari semasa ke semasa atau menggabungkan sumber data lain berdasarkan nilai terkini. Fikirkan tentang profil pengguna, harga saham atau metrik terkumpul.
Semakan KStream berbanding KTable
Anda telah mempelajari KStream dan KTable. Mari uji pemahaman anda.
Ulang Kaji KStream & KTable
Tahniah! Anda telah meneroka perbezaan teras dan kegunaan KStream serta KTable.
- KStream mengendalikan peristiwa individu yang tidak berubah, dan sesuai untuk pemprosesan peristiwa demi peristiwa.
- KTable mengekalkan pandangan terwujud dengan menjejaki keadaan terkini bagi setiap kunci, dan sesuai untuk pengagregatan serta pemprosesan yang mengekalkan keadaan.
Kedua-dua asas ini menjadi landasan untuk membina aplikasi pemprosesan strim yang berkuasa dan fleksibel dengan Kafka Streams.
Pelajari Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) dengan tutor kecerdasan buatan — percuma
Tulis dan jalankan kod sebenar dalam pelayar anda, dapatkan bantuan segera daripada tutor kecerdasan buatan yang tersedia 24/7, dan sambung semula dari tempat anda berhenti di web atau dalam aplikasi.
- Kursus
- 12
- Pelajaran
- 48
Soalan Lazim
Adakah pelajaran “Pemprosesan Strim dengan KStream & KTable” percuma?
Ya — sebanyak 3 pelajaran dalam laluan pembelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka), termasuk “Pemprosesan Strim dengan KStream & KTable”, boleh dibaca sepenuhnya secara percuma di web ini. Selepas itu, CoddyKit PRO membuka akses kepada semua pelajaran, serta latihan interaktif dengan penyunting kod terbina dalam dan tutor kecerdasan buatan yang tersedia 24/7. Kursus Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) merangkumi sejumlah 4 pelajaran.
Apakah yang akan saya pelajari dalam “Pemprosesan Strim dengan KStream & KTable”?
Pelajari cara menggunakan KStream untuk strim peristiwa tidak boleh ubah dan KTable untuk paparan data berkeadaan yang boleh dikemas kini, termasuk operasi seperti penapisan dan pemetaan. Anda berlatih Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) menggunakan kod praktikal yang dijalankan terus dalam pelayar, manakala tutor kecerdasan buatan 24/7 menjawab soalan anda semasa anda mengikuti pelajaran.
Adakah saya memerlukan pengalaman untuk memulakan Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka)?
Tiada pengalaman terdahulu diperlukan. Pembelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) di CoddyKit disusun untuk pelajar daripada peringkat pemula hingga lanjutan, jadi anda boleh bermula di sini atau dari awal dan belajar mengikut kadar anda sendiri. Ini ialah pelajaran 2 daripada 4.
Berapa lamakah pelajaran “Pemprosesan Strim dengan KStream & KTable” diambil?
Kebanyakan pelajaran CoddyKit mengambil masa kira-kira 5–10 minit. Setiap pelajaran ringkas dan interaktif, jadi anda boleh membuat kemajuan secara berterusan dan menyambung tepat dari tempat anda berhenti di web atau aplikasi.
Bolehkah saya menulis dan menjalankan kod dalam pelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) ini?
Ya. Setiap pelajaran Spring Boot 4 Lanjutan: Seni Bina Dipacu Peristiwa (Kafka) menyertakan penyunting kod terbina dalam, jadi anda boleh menulis dan menjalankan kod sebenar terus dalam pelayar serta menerima maklum balas kecerdasan buatan serta-merta — tanpa memerlukan persediaan setempat.
Semua pelajaran dalam kursus ini
- Pengenalan kepada Kafka Streams
- Pemprosesan Strim dengan KStream & KTable
- Membina Aplikasi Strim Ringkas
- Tetingkap dan Pengagregatan Berkeadaan dalam Kafka Streams