Ukratko

Najvažnije iz članka

  • Apache Kafka je distribuirana platforma za *streaming* događaja, ključna za real-time obradu velikih količina podataka.
  • Ključni koncepti uključuju proizvođače, potrošače, brokere, topice i particije, a svi zajedno omogućuju skalabilnost i trajnost.
  • Lokalni Kafka klaster se može jednostavno postaviti pomoću Docker Compose-a, uključujući Zookeeper i Kafka UI za upravljanje.
  • Praktični primjeri pokazuju kako kreirati topice, slati i primati poruke pomoću CLI alata te implementirati proizvođače i potrošače u C# .NET Core aplikacijama.
Sadržaj članka
  1. Što je Apache Kafka?
  2. Ključni koncepti Kafka arhitekture
  3. Zašto koristiti Kafka?
  4. Postavljanje Kafka klastera (Lokalno)
  5. Preduvjeti
  6. Konfiguracija docker-compose.yml
  7. Pokretanje klastera
  8. Upravljanje topicima i podacima
  9. Kreiranje Topica
  10. Sluhanje i slanje poruka (Proizvođač i Potrošač)
  11. Primjer C# .NET Core aplikacije: Proizvođač i Potrošač
  12. 1. Kreiranje .NET projekata
  13. 2. Implementacija Proizvođača (KafkaProducerApp/Program.cs)
  14. 3. Implementacija Potrošača (KafkaConsumerApp/Program.cs)
  15. 4. Pokretanje .NET aplikacija
  16. Napredniji koncepti i scenariji
  17. Zaključak

U današnjem digitalnom dobu, količina generiranih podataka raste eksponencijalnom brzinom. Poduzeća se suočavaju s izazovom obrade, analize i djelovanja na tim podacima u realnom vremenu kako bi ostala konkurentna. Tradicionalni sustavi za obradu podataka često su se pokazali nedovoljnima za ove zahtjeve. Tu na scenu stupa Apache Kafka, distribuirana platforma za streaming događaja, koja je postala de facto standard za rukovanje podacima u pokretu.

Što je Apache Kafka?

Apache Kafka je distribuirana streaming platforma otvorenog koda koja je originalno razvijena u LinkedInu, a kasnije je donirana Apache Software Foundationu. Njena primarna funkcija je omogućiti objavljivanje (publish), pretplatu (subscribe), pohranu i obradu event streamova u realnom vremenu. U svojoj srži, Kafka funkcionira kao distribuirani commit log.

Ključni koncepti Kafka arhitekture

Prije nego što zaronimo u praktične primjere, ključno je razumjeti osnovne komponente i koncepte Kafka arhitekture:

  • Proizvođači (Producers): Aplikacije koje objavljuju poruke (zapise/events) u Kafka topice. Proizvođači mogu slati podatke u različite particije topica.
  • Potrošači (Consumers): Aplikacije koje se pretplaćuju na jedan ili više topica i obrađuju pristigle poruke. Potrošači unutar iste potrošačke grupe (Consumer Group) dijele posao čitanja iz particija topica.
  • Brokeri (Brokers): Kafka serveri koji primaju poruke od proizvođača, pohranjuju ih na disk i serviraju ih potrošačima. Klaster Kafka sastoji se od jednog ili više brokera.
  • Topici (Topics): Logičke kategorije ili nazivi feedova za objavljivanje zapisa. Svi zapisi objavljeni od strane proizvođača pripadaju određenom topicu. Topic je podijeljen na particije.
  • Particije (Partitions): Svaki topic je podijeljen na jednu ili više particija. Particije omogućuju paralelizaciju obrade podataka i skalabilnost. Svaka particija je sekvencijski, nepromjenjiv log zapisa, a svaki zapis unutar particije dobiva redni broj (offset). Redoslijed poruka je zagarantiran unutar jedne particije, ali ne i između particija.
  • Offset: Jedinstveni, sekvencijalni identifikator zapisa unutar particije. Potrošači prate svoj offset kako bi znali gdje su stali s čitanjem.
  • Zookeeper: Kafka se tradicionalno oslanja na Apache Zookeeper za upravljanje stanjem klastera, održavanje konfiguracija, koordinaciju brokera i ostalih metadata operacija. NAPOMENA: S verzijom Kafka 2.8+, uvodi se mogućnost KRaft mode (Kafka Raft) koji eliminira potrebu za Zookeeperom, što pojednostavljuje arhitekturu.

Zašto koristiti Kafka?

  • Visoka propusnost (High Throughput): Sposobnost obrade milijuna poruka u sekundi.
  • Skalabilnost (Scalability): Lako se skalira dodavanjem više brokera u klaster.
  • Trajnost (Durability): Podaci se pohranjuju na disk s mogućnošću replikacije, osiguravajući da se poruke ne izgube.
  • Fleksibilnost (Flexibility): Podržava različite use-caseove, od log aggregation do event sourcing i stream processing.
  • Fault-Tolerantnost: Dizajnirana da preživi kvarove brokera bez gubitka podataka.

Postavljanje Kafka klastera (Lokalno)

Za naše primjere, postavit ćemo minimalni Kafka klaster na lokalnom stroju. Koristit ćemo Docker Compose za jednostavniju deployment i upravljanje.

Preduvjeti

  • Docker Desktop instaliran i pokrenut na vašem sustavu.

Konfiguracija docker-compose.yml

Kreirajte datoteku docker-compose.yml sa sljedećim sadržajem:

version: '3.8'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.3
    hostname: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka:
    image: confluentinc/cp-kafka:7.5.3
    hostname: kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
    depends_on:
      - zookeeper

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    hostname: kafka-ui
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local-kafka
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
      KAFKA_CLUSTERS_0_ZOOKEEPERCONNECT: zookeeper:2181
    depends_on:
      - kafka

Ovaj docker-compose.yml definira tri servisa:

  1. zookeeper: Zookeeper instanca koju koristi Kafka.
  2. kafka: Kafka broker. KAFKA_ADVERTISED_LISTENERS je važan jer definira kako drugi servisi unutar Docker mreže (PLAINTEXT://kafka:29092) i vanjski klijenti (PLAINTEXT_HOST://localhost:9092) mogu pristupiti brokeru.
  3. kafka-ui: Korisničko sučelje za pregled i upravljanje Kafka klasterom, vrlo korisno za debugging i nadzor.

Pokretanje klastera

Otvorite terminal u direktoriju gdje ste spremili docker-compose.yml i pokrenite:

docker compose up -d

Ovo će preuzeti potrebne Docker slike i pokrenuti kontejnere u pozadini. Možete provjeriti status s:

docker compose ps

Nakon što se svi servisi pokrenu, možete pristupiti Kafka UI-ju putem preglednika na http://localhost:8080.

Upravljanje topicima i podacima

Kafka CLI alati su esencijalni za interakciju s klasterom. Pristupit ćemo Kafka kontejneru da bismo ih koristili.

Kreiranje Topica

Kreirat ćemo topic naziva moj-prvi-topic s 3 particije i faktorom replikacije 1 (za lokalni setup).

docker exec kafka kafka-topics --create --bootstrap-server kafka:29092 --replication-factor 1 --partitions 3 --topic moj-prvi-topic

Provjera kreiranih topica:

docker exec kafka kafka-topics --list --bootstrap-server kafka:29092

Sluhanje i slanje poruka (Proizvođač i Potrošač)

Otvorite dva nova terminala. U jednom će biti consumer (potrošač), a u drugom producer (proizvođač).

Terminal 1 (Potrošač):

docker exec kafka kafka-console-consumer --bootstrap-server kafka:29092 --topic moj-prvi-topic --from-beginning

Ova naredba će pokrenuti potrošača koji sluša moj-prvi-topic i prikazuje sve poruke od početka topica.

Terminal 2 (Proizvođač):

docker exec -it kafka kafka-console-producer --bootstrap-server kafka:29092 --topic moj-prvi-topic

Sada možete tipkati poruke u Terminalu 2 (proizvođač) i pritiskati Enter. Svaka linija će biti poslana kao poruka u moj-prvi-topic. Vidjet ćete kako se te poruke odmah pojavljuju u Terminalu 1 (potrošač). Ovo demonstrira osnovni real-time streaming podataka.

Primjer C# .NET Core aplikacije: Proizvođač i Potrošač

Koristit ćemo Confluent.Kafka klijentsku biblioteku za .NET Core, koja je u biti wrapper oko librdkafka, performantne C/C++ Kafka klijentske biblioteke.

1. Kreiranje .NET projekata

Kreirajte dva nova konzolna projekta:

dotnet new console -n KafkaProducerApp
dotnet new console -n KafkaConsumerApp
cd KafkaProducerApp
dotnet add package Confluent.Kafka
cd ../KafkaConsumerApp
dotnet add package Confluent.Kafka
cd ..

2. Implementacija Proizvođača (KafkaProducerApp/Program.cs)

using Confluent.Kafka;
using System;
using System.Threading;
using System.Threading.Tasks;

namespace KafkaProducerApp
{
    class Program
    {
        static async Task Main(string[] args)
        {
            var config = new ProducerConfig { BootstrapServers = "localhost:9092" };

            using (var producer = new ProducerBuilder<Null, string>(config).Build())
            {
                Console.WriteLine("Kafka Producer pokrenut. Unosite poruke:");
                while (true)
                {
                    Console.Write("> ");
                    string message = Console.ReadLine();

                    if (string.IsNullOrWhiteSpace(message))
                    {
                        if (message == "exit") break;
                        continue;
                    }

                    try
                    {
                        var dr = await producer.ProduceAsync("moj-prvi-topic", new Message<Null, string> { Value = message });
                        Console.WriteLine($"Dostavljeno '{dr.Value}' u '{dr.TopicPartitionOffset}'");
                    }
                    catch (ProduceException<Null, string> e)
                    {
                        Console.WriteLine($"Greška pri dostavi: {e.Error.Reason}");
                    }
                }
            }
            Console.WriteLine("Producer zaustavljen.");
        }
    }
}

Ovaj proizvođač prima unos od korisnika i šalje ga kao poruku u moj-prvi-topic.

3. Implementacija Potrošača (KafkaConsumerApp/Program.cs)

using Confluent.Kafka;
using System;
using System.Threading;
using System.Threading.Tasks;

namespace KafkaConsumerApp
{
    class Program
    {
        static void Main(string[] args)
        {
            var config = new ConsumerConfig
            {
                BootstrapServers = "localhost:9092",
                GroupId = "moj-potrosacki-klaster", // Jedinstveni GroupId
                AutoOffsetReset = AutoOffsetReset.Earliest // Počinje čitati od početka ako nema zapamćenog offseta
            };

            using (var consumer = new ConsumerBuilder<Ignore, string>(config).Build())
            {
                consumer.Subscribe("moj-prvi-topic");
                var cts = new CancellationTokenSource();
                Console.CancelKeyPress += (_, e) => {
                    e.Cancel = true; // Spriječi odmah prekid programa
                    cts.Cancel();
                };

                Console.WriteLine("Kafka Consumer pokrenut. Čeka poruke...");

                try
                {
                    while (true)
                    {
                        try
                        {
                            var cr = consumer.Consume(cts.Token);
                            Console.WriteLine($"Potrošeno poruku iz '{cr.TopicPartitionOffset}': '{cr.Message.Value}'");
                        }
                        catch (ConsumeException e)
                        {
                            Console.WriteLine($"Greška pri potrošnji: {e.Error.Reason}");
                        }
                    }
                }
                catch (OperationCanceledException)
                {
                    // Korisnik je pritisnuo Ctrl+C
                    Console.WriteLine("Gašenje potrošača...");
                    consumer.Close();
                }
            }
            Console.WriteLine("Consumer zaustavljen.");
        }
    }
}

Ovaj potrošač se pretplaćuje na moj-prvi-topic i ispisuje sve pristigle poruke.

4. Pokretanje .NET aplikacija

Otvorite dva nova terminala. U jednom pokrenite proizvođača, a u drugom potrošača:

Terminal 3 (Proizvođač):

cd KafkaProducerApp
dotnet run

Terminal 4 (Potrošač):

cd KafkaConsumerApp
dotnet run

Opet, kao i s CLI alatima, moći ćete unijeti poruke u proizvođaču i vidjeti ih kako se odmah pojavljuju u potrošaču. Ovo pokazuje kako programski možete integrirati Kafka u svoje aplikacije.

Napredniji koncepti i scenariji

  • Schema Registry: Za osiguravanje kompatibilnosti podataka i evolucije sheme, često se koristi Confluent Schema Registry zajedno s Apache Avro, Protobuf ili JSON Schema. Proizvođači i potrošači koriste sheme za serijalizaciju/deserijalizaciju podataka.
  • Kafka Streams API: Biblioteka za izgradnju stream processing aplikacija direktno na Kafki. Omogućuje pisanje kompleksnih operacija poput agregacije, spajanja streamova i transformacija.
  • Kafka Connect: Alat za skalabilan i pouzdan streaming podataka između Apache Kafka i drugih sustava (baza podataka, messaging sustava, itd.).
  • Replikacija i Visoka Dostupnost: Pravilna konfiguracija faktora replikacije particija (obično 3) osigurava da podaci nisu izgubljeni čak i ako neki brokera padnu.
  • Monitoring: Praćenje Kafka klastera je ključno. Alati poput Prometheus i Grafana, zajedno s JMX metrikama, pružaju uvid u performanse i zdravlje klastera.

Zaključak

Apache Kafka je moćan alat za real-time data streaming i obradu događaja. Njegova skalabilnost, trajnost i visoka propusnost čine ga idealnim rješenjem za širok spektar aplikacija, od prikupljanja logova do izgradnje event-driven mikroservisa. Razumijevanjem ključnih koncepata i stjecanjem praktičnog iskustva s postavljanjem i programiranjem, možete iskoristiti puni potencijal Kafke za transformaciju načina na koji vaše aplikacije rukuju podacima. Ovaj tutorial pružio je osnovu za početak, a daljnje istraživanje naprednih značajki i ekosustava oko Kafke otvara vrata za još složenije i robusnije arhitekture.

Izvori i dodatno čitanje

  1. Apache Kafka Documentation
  2. Confluent Developer - Kafka Tutorials
  3. Confluent.Kafka GitHub Repository
B
Uredništvo portala

BAJT

Službeni autorski profil redakcije portala BAJT. Sadržaj priprema i provjerava uredništvo portala.