Apache Kafka to jedno z najpotężniejszych narzędzi w przetwarzaniu danych strumieniowych, często wykorzystywane w nowoczesnych aplikacjach opartych na mikroserwisach. W tym poście pokażemy, jak Kafka integruje się z Symfony – jednym z najpopularniejszych frameworków PHP. Zajmiemy się tematami związanymi z produkcją i konsumpcją wiadomości oraz konfiguracją klastrów, które są kluczowe dla zapewnienia wydajności i niezawodności systemu. Jeśli jesteś gotowy na praktyczne ćwiczenia i chcesz zrozumieć, jak działa Kafka w kontekście aplikacji webowych, to ten artykuł jest dla Ciebie!
1. Produkcja i Konsumpcja Wiadomości w Apache Kafka
Apache Kafka działa na zasadzie komunikacji „publish-subscribe” (publikacja i subskrypcja), co oznacza, że jeden element systemu (producent) wysyła wiadomości, a inne elementy (konsumenci) subskrybują te wiadomości.
Produkcja i Konsumpcja przy użyciu narzędzi CLI
Kafka dostarcza narzędzia wiersza poleceń (CLI), które są niezwykle przydatne na etapie nauki oraz do szybkiego testowania systemu. Pozwalają one na tworzenie tematów, produkcję wiadomości i konsumpcję wiadomości bez konieczności pisania kodu.
Tworzenie Tematu: Aby Kafka wiedziała, gdzie przechowywać wiadomości, musimy utworzyć temat:
bin/kafka-topics.sh --create --topic moj-temat --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1
--create: Tworzy temat.--topic moj-temat: Nazwa tematu.--partitions 3: Liczba partycji (więcej na ten temat w dalszej części).--replication-factor 1: Liczba replik wiadomości.
- Produkcja Wiadomości: Producent jest odpowiedzialny za wysyłanie wiadomości do Kafki:
bin/kafka-console-producer.sh --topic moj-temat --bootstrap-server localhost:9092
- Teraz możesz wpisywać wiadomości, które będą wysyłane do tematu „moj-temat”. Każda wiadomość zostanie umieszczona w jednej z trzech partycji.
- Konsumpcja Wiadomości: Konsument subskrybuje i odbiera wiadomości z danego tematu:
bin/kafka-console-consumer.sh --topic moj-temat --from-beginning --bootstrap-server localhost:9092
--from-beginning: Konsument odbiera wszystkie wiadomości od początku, niezależnie od tego, kiedy został uruchomiony.
Role Producenta i Konsumenta w Kafce
Producent i konsument odgrywają kluczowe role w architekturze Apache Kafka:
- Producent (Producer): To aplikacja, która wysyła wiadomości do tematu w Kafce. Producent decyduje, do której partycji temat zostanie przypisany.
- Konsument (Consumer): To aplikacja, która odbiera wiadomości z danego tematu. Konsumenci mogą należeć do grup konsumentów, co ułatwia równomierne rozdzielanie obciążenia między wielu konsumentów.
Ćwiczenia Praktyczne – Symfony, Producent i Konsument
Teraz przejdźmy do tworzenia aplikacji przy użyciu frameworka Symfony.
Konfiguracja Symfony z Apache Kafka
Aby zintegrować Symfony z Apache Kafka, najpierw musimy zainstalować bibliotekę, która pomoże nam połączyć się z Kafką. Najpopularniejszą opcją jest php-enqueue/enqueue-bundle:
composer require enqueue/enqueue-bundle
Tworzenie Producenta w Symfony
W Symfony producent wiadomości będzie odpowiedzialny za wysyłanie danych do tematu Kafki.
Przykładowy kod Producenta:
// src/Service/KafkaProducer.php
namespace App\Service;
use Enqueue\RdKafka\RdKafkaContext;
class KafkaProducer
{
private $context;
public function __construct(RdKafkaContext $context)
{
$this->context = $context;
}
public function sendMessage(string $topic, string $message): void
{
$producer = $this->context->createProducer();
$topic = $this->context->createTopic($topic);
$message = $this->context->createMessage($message);
$producer->send($topic, $message);
}
}
Krok po kroku:
RdKafkaContext– Jest to kontekst Kafki, który umożliwia tworzenie producenta oraz tematów.createProducer()– Tworzy producenta, który wysyła wiadomości.createTopic($topic)– Tworzy temat, do którego będziemy wysyłać wiadomości.createMessage($message)– Tworzy wiadomość, która następnie jest wysyłana do tematu.
Tworzenie Konsumenta w Symfony
Konsument będzie odpowiedzialny za odbieranie wiadomości z tematu.
Przykładowy kod Konsumenta:
// src/Command/KafkaConsumerCommand.php
namespace App\Command;
use Enqueue\RdKafka\RdKafkaContext;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Output\OutputInterface;
class KafkaConsumerCommand extends Command
{
protected static $defaultName = 'app:kafka:consume';
private $context;
public function __construct(RdKafkaContext $context)
{
parent::__construct();
$this->context = $context;
}
protected function execute(InputInterface $input, OutputInterface $output): int
{
$consumer = $this->context->createConsumer($this->context->createQueue('moj-temat'));
while (true) {
$message = $consumer->receive();
if ($message) {
$output->writeln('Received: ' . $message->getBody());
$consumer->acknowledge($message);
}
}
return Command::SUCCESS;
}
}
Krok po kroku:
KafkaConsumerCommand– Definiuje konsumenta jako komendę, którą możemy uruchomić w konsoli.createConsumer()– Tworzy konsumenta dla wybranego tematu.receive()– Odbiera wiadomości z tematu.acknowledge($message)– Potwierdza, że wiadomość została przetworzona.
2. Konfiguracja Klastrów i Partycjonowanie
Konfiguracja Klastra Kafki
Kafka wspiera konfigurację klastrów, co zapewnia zarówno wysoką dostępność, jak i skalowalność. Klaster to grupa brokerów Kafki, które wspólnie pracują nad przetwarzaniem i przechowywaniem wiadomości.
- Replikacja: Każda partycja w temacie może być replikowana na kilka brokerów. Gwarantuje to, że w razie awarii jednego brokera, dane będą dostępne na innym brokerze.
- Wydajność: Wiadomości są przechowywane w partycjach. Więcej partycji oznacza, że więcej konsumentów może przetwarzać dane równolegle, co zwiększa przepustowość systemu.
Partycje – Klucz do Skalowalności
Partycjonowanie to klucz do skalowalności Apache Kafka. Temat może być podzielony na wiele partycji, co oznacza, że może być przetwarzany przez wielu konsumentów równolegle.
- Ustawienia partycji: Im więcej partycji, tym większa możliwość równoległego przetwarzania wiadomości, ale więcej partycji oznacza również więcej metadanych do zarządzania, co może wpływać na wydajność.
Zarządzanie Danymi i Replikacja
Kafka wspiera konfigurację, w której każda partycja może mieć więcej niż jedną replikę. Replikacja pozwala na:
- Odporność na awarie: Jeśli jeden broker przestanie działać, repliki zapewnią dostępność danych.
- Balansowanie obciążenia: Dzięki replikacji różne brokerzy mogą równomiernie dzielić obciążenie.
Aby skonfigurować replikację, warto ustawić odpowiedni replication-factor przy tworzeniu tematu:
bin/kafka-topics.sh --create --topic moj-temat --bootstrap-server localhost:9092 --partitions 3 --replication-factor 2
To oznacza, że każda partycja będzie miała swoją kopię przechowywaną na innym brokerze.
Podsumowanie
Apache Kafka w połączeniu z Symfony umożliwia tworzenie skalowalnych aplikacji, które mogą przetwarzać duże ilości danych w czasie rzeczywistym. W tym artykule pokazaliśmy, jak tworzyć producentów i konsumentów oraz jak skonfigurować klaster Kafki dla większej niezawodności i wydajności.
Jeśli jesteś nowy w świecie Apache Kafka, polecam zacząć od praktycznych ćwiczeń – uruchom konsumenta i producenta w Symfony, przetestuj przesyłanie wiadomości, a następnie spróbuj skonfigurować klaster z wieloma brokerami. Kafka daje potężne narzędzia, które naprawdę warto wykorzystać w nowoczesnych aplikacjach.