Apache Kafka & Stream Processing Fundamentals · Ders

Akışlarda Birleştirmeler ve Toplamalar

Akışlar ve tablolar arasında karmaşık birleştirmeler gerçekleştirin ve gerçek zamanlı olarak anlamlı içgörüler elde etmek için verileri toplayın.

2. ders / 411 adım

Akışlarda Birleştirmeler ve Toplamalar, CoddyKit'te ücretsiz bir Apache Kafka & Stream Processing Fundamentals dersidir. Bu, 4 dersinin 2. dersidir. Aşağıdan dersin tamamını ücretsiz okuyabilir, sonra tarayıcıda yerleşik kod editörü ve 7/24 yapay zeka koçu ile uygulamalı olarak pratik yapabilirsin. Bu, Apache Kafka & Stream Processing Fundamentals öğrenme yolunun bir parçasıdır ve ilerlemeniz web ve CoddyKit uygulaması arasında senkronize olur. Apache Kafka & Stream Processing Fundamentals kursu toplamda 4 dersten oluşur.

Bu dersin bazı bölümleri henüz çevrilmemiş olup İngilizce olarak gösterilmektedir.

Combining Data Streams

In real-world data processing, you often need to combine information from different sources. Imagine tracking user clicks and matching them with user profiles, or correlating an order with its payment details.

Kafka Streams provides powerful operations to join different data streams (KStreams) and tables (KTables) based on a common key. This allows you to enrich your data and derive more complete insights in real-time.

KStream-KStream Joins

A KStream-KStream join combines records from two KStreams based on their shared key. Since KStreams represent unbounded, continuous event streams, these joins require a time window.

  • Events from both streams must arrive within this defined time window to be considered for a join.
  • If an event from one stream arrives outside the window of its matching event in the other stream, they won't be joined.
  • This is crucial for correlating events that happen close together, like a user clicking an ad and then visiting a product page.

KStream-KStream Join Example

This example demonstrates joining two KStreams, streamA and streamB, using a 10-second time window. Only records with the same key arriving within this window will be combined.

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.JoinWindows;
import org.apache.kafka.streams.kstream.KStream;
import java.time.Duration;
import java.util.Properties;

public class StreamStreamJoin {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, "stream-join-app");
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());
    props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());

    StreamsBuilder builder = new StreamsBuilder();
    KStream<String, String> streamA = builder.stream("topic-A");
    KStream<String, String> streamB = builder.stream("topic-B");

    KStream<String, String> joined = streamA.join(
        streamB,
        (valA, valB) -> "Joined: " + valA + "-" + valB,
        JoinWindows.of(Duration.ofSeconds(10))
    );
    joined.to("joined-topic");

    KafkaStreams streams = new KafkaStreams(builder.build(), props);
    streams.start();
  }
}

KStream-KTable Joins

A KStream-KTable join combines an event stream (KStream) with a materialized view or state table (KTable). This is a very common pattern for data enrichment.

  • When a new record arrives on the KStream, it's joined with the current state of the KTable for the matching key.
  • No time window is explicitly needed for the KTable side, as it always represents the latest known state.
  • Think of it as looking up additional details for an event from a constantly updating database.

KStream-KTable Join Example

Here, a stream of transactions is enriched with data from a user-profiles KTable. Each transaction record gets the latest profile information for its user.

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 java.util.Properties;

public class StreamTableJoin {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, "stream-table-join-app");
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());
    props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());

    StreamsBuilder builder = new StreamsBuilder();
    KStream<String, String> transactions = builder.stream("transactions");
    KTable<String, String> userProfiles = builder.table("user-profiles");

    KStream<String, String> enrichedTransactions = transactions.join(
        userProfiles,
        (transactionVal, profileVal) -> "Tx: " + transactionVal + ", User: " + profileVal
    );
    enrichedTransactions.to("enriched-transactions");

    KafkaStreams streams = new KafkaStreams(builder.build(), props);
    streams.start();
  }
}

KTable-KTable Joins

A KTable-KTable join combines two materialized views (KTables) based on their shared key. This is similar to joining two constantly updating database tables.

  • Whenever a record in either KTable is updated, the join operation is re-evaluated for that key.
  • The output KTable will reflect the combined latest state of the matching records from both input KTables.
  • This is useful for combining different aspects of an entity, like product pricing and inventory levels.

KTable-KTable Join Example

This code joins product-prices and product-stocks KTables. Any update to a product's price or stock will trigger an update to the joined-products KTable.

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 java.util.Properties;

public class TableTableJoin {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, "table-table-join-app");
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());
    props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());

    StreamsBuilder builder = new StreamsBuilder();
    KTable<String, String> productPrices = builder.table("product-prices");
    KTable<String, String> productStocks = builder.table("product-stocks");

    KTable<String, String> joinedProducts = productPrices.join(
        productStocks,
        (price, stock) -> "Price: " + price + ", Stock: " + stock
    );
    joinedProducts.toStream().to("joined-products");

    KafkaStreams streams = new KafkaStreams(builder.build(), props);
    streams.start();
  }
}

Understanding Aggregations

Aggregations are operations that summarize data from a stream or table over a specific key or time window. Common aggregations include:

  • Counting: How many events occurred for a key?
  • Summing: What is the total value for a key?
  • Averaging: What is the average value for a key?
  • Reducing: Combining values using a custom logic.

Aggregations are fundamental for building real-time dashboards, metrics, and summary statistics from continuous data streams.

KStream Aggregation Example

This example demonstrates a simple aggregation: counting user events per user ID within 10-second tumbling windows. The groupByKey() and windowedBy() methods are key here.

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.Materialized;
import org.apache.kafka.streams.kstream.TimeWindows;
import java.time.Duration;
import java.util.Properties;

public class StreamAggregation {
  public static void main(String[] args) {
    Properties props = new Properties();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, "stream-aggregation-app");
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());
    props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_BY_KEY_CONFIG, Serdes.String().getClass());

    StreamsBuilder builder = new StreamsBuilder();
    KStream<String, String> userEvents = builder.stream("user-events");

    userEvents
        .groupByKey()
        .windowedBy(TimeWindows.of(Duration.ofSeconds(10)))
        .count(Materialized.as("user-event-counts"))
        .toStream((windowedKey, count) -> windowedKey.key())
        .to("user-event-counts-output");

    KafkaStreams streams = new KafkaStreams(builder.build(), props);
    streams.start();
  }
}

Join & Aggregate Check

Which type of Kafka Streams join is typically used for data enrichment, where an incoming event stream is combined with the latest state of a dataset?

Joins & Aggregations Summary

You've learned how Kafka Streams allows you to combine and summarize data in powerful ways:

  • KStream-KStream Joins: Correlate events from two streams within a time window.
  • KStream-KTable Joins: Enrich stream events with the latest state from a table.
  • KTable-KTable Joins: Combine two continuously updating materialized views.
  • Aggregations: Summarize data (e.g., count, sum) over keys and windows to derive real-time insights.

These operations are key to building sophisticated real-time analytics and data processing pipelines with Kafka Streams.

Başlamak ücretsiz

Yapay zeka eğitmeniyle Apache Kafka & Stream Processing Fundamentals öğren — ücretsiz

Tarayıcında gerçek kod yaz ve çalıştır, 7/24 yapay zeka eğitmeninden anında yardım al; web'de ya da uygulamada kaldığın yerden devam et.

Kurslar
12
Dersler
48

Sıkça Sorulan Sorular

“Akışlarda Birleştirmeler ve Toplamalar” dersi ücretsiz mi?

Evet — “Akışlarda Birleştirmeler ve Toplamalar” dersin tüm metni burada web'de ücretsiz olarak okunabilir. Etkileşimli olarak pratik yapmak (yerleşik kod editörü ve 7/24 yapay zeka koçu) ve Apache Kafka & Stream Processing Fundamentals kursunun geri kalanını açmak için CoddyKit PRO'ya yükselt. Apache Kafka & Stream Processing Fundamentals kursu toplamda 4 dersten oluşur.

“Akışlarda Birleştirmeler ve Toplamalar” dersinde ne öğreneceğim?

Akışlar ve tablolar arasında karmaşık birleştirmeler gerçekleştirin ve gerçek zamanlı olarak anlamlı içgörüler elde etmek için verileri toplayın. Apache Kafka & Stream Processing Fundamentals ile uygulamalı kodu tarayıcıda doğrudan çalıştırarak pratik yaparsın ve 7/24 yapay zeka koçu dersi çalışırken sorularını yanıtlar.

Apache Kafka & Stream Processing Fundamentals öğrenmeye başlamak için deneyim gerekli mi?

Önceden deneyim gerekmez. CoddyKit'te Apache Kafka & Stream Processing Fundamentals, başlangıçtan ileri seviyeye kadar yapılandırıldığı için buradan başlayabilir veya başından başlayıp kendi hızında ilerleme yapabilirsin. Bu, 4 dersinin 2. dersidir.

“Akışlarda Birleştirmeler ve Toplamalar” dersi ne kadar sürer?

Çoğu CoddyKit dersi yaklaşık 5–10 dakika sürer. Her biri kısa ve etkileşimli olduğu için sabit ilerleme yaparsın ve web ile uygulama arasında tam olarak bıraktığın yerden devam edebilirsin.

Bu Apache Kafka & Stream Processing Fundamentals dersinde kod yazıp çalıştırabilir miyim?

Evet. Her Apache Kafka & Stream Processing Fundamentals dersi yerleşik bir kod editörü içerir, bu sayede tarayıcıda gerçek kod yazıp çalıştırabilir ve anlık yapay zeka geri bildirimi alırsın — yerel kurulum gerekli değildir.

Bu kursun tüm dersleri

  1. Kafka Streams'te Pencereleme İşlemleri
  2. Akışlarda Birleştirmeler ve Toplamalar
  3. Akış Analizi için KSQL'e Giriş
  4. Etkileşimli Sorgular ve Durum Depoları
← Apache Kafka & Stream Processing Fundamentals Sayfasına Dön