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
- Što je Apache Kafka?
- Ključni koncepti Kafka arhitekture
- Zašto koristiti Kafka?
- Postavljanje Kafka klastera (Lokalno)
- Preduvjeti
- Konfiguracija docker-compose.yml
- Pokretanje klastera
- Upravljanje topicima i podacima
- Kreiranje Topica
- Sluhanje i slanje poruka (Proizvođač i Potrošač)
- Primjer C# .NET Core aplikacije: Proizvođač i Potrošač
- 1. Kreiranje .NET projekata
- 2. Implementacija Proizvođača (KafkaProducerApp/Program.cs)
- 3. Implementacija Potrošača (KafkaConsumerApp/Program.cs)
- 4. Pokretanje .NET aplikacija
- Napredniji koncepti i scenariji
- 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:
zookeeper: Zookeeper instanca koju koristi Kafka.kafka: Kafka broker.KAFKA_ADVERTISED_LISTENERSje važan jer definira kako drugi servisi unutar Docker mreže (PLAINTEXT://kafka:29092) i vanjski klijenti (PLAINTEXT_HOST://localhost:9092) mogu pristupiti brokeru.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.
Komentari