Apache Kafka z PHP
Przesyłanie strumieniowe zdarzeń o dużej przepustowości za pomocą Kafka
Apache Kafka z PHP to bezpłatna lekcja PHP Academy na CoddyKit. To lekcja 3 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej PHP Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs PHP Academy zawiera 4 lekcji w sumie.
Kafka w PHP
Kafka nie jest kolejką zadań — to rozproszony dziennik tylko do dopisywania. Producenci dopisują rekordy do tematów, a konsumenci odczytują je od wybranego przez siebie offsetu i mogą ponownie odtworzyć historię. Dzięki temu Kafka doskonale nadaje się do strumieni zdarzeń o dużej przepustowości, event sourcingu oraz zasilania wielu niezależnych grup konsumentów z jednego strumienia.
W PHP komunikacja z Kafką odbywa się za pośrednictwem ext-rdkafka — bindingu do sprawdzonej biblioteki C librdkafka.
Dziennik, nie kolejka
Zmiana sposobu myślenia przy przejściu z RabbitMQ na Kafkę:
- Wiadomości nie są usuwane po odczytaniu; wygasają zgodnie z polityką retencji (czasową lub rozmiarową).
- Każdy konsument śledzi własny offset — swoją pozycję w dzienniku.
- Temat jest podzielony na partycje; kolejność jest gwarantowana wyłącznie w obrębie jednej partycji.
- Wiele grup konsumentów niezależnie od siebie odczytuje ten sam temat.
Jeśli potrzebne jest ponowne odtwarzanie, fan-out do wielu odbiorców lub ogromna przepustowość, Kafka będzie dobrym wyborem. Jeśli potrzebne jest kierowanie poszczególnych wiadomości i TTL, lepszy będzie RabbitMQ.
Instalowanie ext-rdkafka
Najpierw należy zainstalować bibliotekę natywną, następnie rozszerzenie PECL, a na końcu — opcjonalnie — wrapper wyższego poziomu.
# Debian/Ubuntu
apt-get install -y librdkafka-dev
pecl install rdkafka
echo "extension=rdkafka.so" >> php.ini
# Optional ergonomic wrapper
composer require longlang/phpkafka # or arnaud-lb/php-rdkafka-stubs for IDEPublikowanie rekordów
Należy utworzyć obiekt RdKafka\Producer, uzyskać uchwyt tematu i wywołać produce(). Klucz decyduje o tym, do której partycji trafi rekord — ten sam klucz oznacza tę samą partycję i zachowanie kolejności. Przed zakończeniem procesu zawsze należy wywołać flush(), w przeciwnym razie buforowane rekordy zostaną utracone.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('orders');
// RD_KAFKA_PARTITION_UA = let the partitioner choose by key
$key = 'order-42';
$topic->produce(RD_KAFKA_PARTITION_UA, 0, json_encode(['id' => 42]), $key);
$producer->poll(0);
$result = $producer->flush(10000); // wait up to 10s
if ($result !== RD_KAFKA_RESP_ERR_NO_ERROR) {
throw new RuntimeException('Failed to flush');
}Partycjonowanie według klucza
Partycjonowanie jest podstawą skalowalności i zachowania kolejności w Kafce. Domyślny partycjoner oblicza hash klucza rekordu: partition = hash(key) % numPartitions. Wybór dobrego klucza ma znaczenie:
- Klucz
customerId→ wszystkie zdarzenia danego klienta zachowują kolejność i pozostają na jednej partycji. - Brak klucza → rozdzielanie round-robin między partycje (maksymalna przepustowość, brak gwarancji kolejności).
Liczby partycji tematu nie można zmniejszyć, a dodanie partycji zmienia mapowanie wartości hashujących — dlatego liczbę partycji należy z góry dobrać do maksymalnego poziomu równoległości.
Grupy konsumentów i offsety
Należy użyć wysokopoziomowego KafkaConsumer wraz z group.id. Kafka rozdziela partycje między członków grupy i wykonuje rebalansowanie, gdy członkowie dołączają do grupy lub ją opuszczają. Każdy członek odczytuje tylko przypisane mu partycje, co zapewnia skalowanie poziome bez dodatkowej pracy.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('group.id', 'order-emailers');
$conf->set('auto.offset.reset', 'earliest'); // start of log if no offset
$conf->set('enable.auto.commit', 'false'); // we commit manually
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['orders']);
while (true) {
$msg = $consumer->consume(5000);
if ($msg->err === RD_KAFKA_RESP_ERR_NO_ERROR) {
handle($msg->payload);
$consumer->commit($msg); // commit offset AFTER work
}
}
function handle(string $p): void {}Kiedy zatwierdzać
Zatwierdzenie offsetu określa semantykę dostarczania:
- Zatwierdzenie po przetworzeniu → co najmniej raz (awaria przed zatwierdzeniem spowoduje ponowne odtworzenie rekordu).
- Zatwierdzenie przed przetworzeniem → co najwyżej raz (awaria spowoduje jego utratę).
Automatyczne zatwierdzanie (enable.auto.commit=true) zatwierdza offset według zegara niezależnie od tego, czy praca została zakończona — jest wygodne, ale podczas awarii może po cichu utracić rekordy. Gdy poprawność ma znaczenie, należy je wyłączyć i zatwierdzać offsety ręcznie.
Obsługa błędów odbierania
Nie każda wartość zwrócona przez consume() oznacza wiadomość. Należy rozgałęzić logikę na podstawie kodu błędu — __PARTITION_EOF i __TIMED_OUT to normalne sygnały sterujące, a nie awarie.
<?php
$msg = $consumer->consume(2000);
switch ($msg->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
echo "Got: {$msg->payload} @ offset {$msg->offset}\n";
break;
case RD_KAFKA_RESP_ERR__PARTITION_EOF:
echo "Reached end of partition\n"; // caught up, keep polling
break;
case RD_KAFKA_RESP_ERR__TIMED_OUT:
echo "No message this poll\n";
break;
default:
throw new \Exception($msg->errstr(), $msg->err);
}Dostrajanie przepustowości
Przepustowość Kafki wynika z grupowania rekordów w partie. Najważniejsze ustawienia producenta:
linger.ms— krótkie oczekiwanie pozwalające zgromadzić więcej rekordów w jednym żądaniu (np. 5–20 ms).batch.size/queue.buffering.max.messages— większe bufory i mniej operacji komunikacji z brokerem.compression.type—lz4lubzstdznacznie ogranicza koszt sieciowy.acks—alldla trwałości,1dla mniejszych opóźnień.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('compression.type', 'lz4');
$conf->set('linger.ms', '10');
$conf->set('batch.size', '65536');
$conf->set('acks', 'all');
$producer = new RdKafka\Producer($conf);Callbacki raportu dostarczenia
Ponieważ produce() działa asynchronicznie, nieudane wysłanie nie zgłosi błędu od razu. Należy zarejestrować callback raportu dostarczenia w konfiguracji producenta, aby poznać los każdego rekordu — w PHP jest to jedyny niezawodny sposób wykrywania cichych awarii producenta.
<?php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->setDrMsgCb(function ($producer, $msg) {
if ($msg->err) {
fwrite(STDERR, 'Delivery FAILED: ' . rd_kafka_err2str($msg->err) . "\n");
} else {
echo "Delivered to partition {$msg->partition} @ offset {$msg->offset}\n";
}
});
$producer = new RdKafka\Producer($conf);
// poll() services the callback queue; call it after producing
$producer->poll(0);Pułapki specyficzne dla PHP
Kafka zakłada istnienie klientów działających przez długi czas, a cykl życia żądania w PHP temu nie sprzyja:
- Producenci asynchronicznie buforują rekordy — przed zakończeniem skryptu zawsze należy wywołać
flush(), w przeciwnym razie rekordy zostaną utracone. - Należy uruchamiać konsumentów jako trwałe workery CLI pod kontrolą supervisora, nigdy wewnątrz żądania internetowego.
- Rebalansowanie wstrzymuje odbieranie; praca nad pojedynczą wiadomością powinna być krótka albo należy ustawić odpowiednio dużą wartość
max.poll.interval.ms, aby konsument nie został usunięty z grupy. - Należy ustawić
log_leveli zarejestrować callback raportu dostarczenia producenta, aby wykrywać ciche awarie.
Szybkie sprawdzenie
Gwarancje kolejności w Kafce.
Podsumowanie
Kafka w PHP:
- Kafka to dziennik, który można odtwarzać, a nie kolejka; konsumenci śledzą offsety.
- Partycjonowanie według klucza zapewnia zachowanie kolejności dla danego klucza i skalowanie.
- Grupy konsumentów dzielą między siebie partycje i automatycznie wykonują rebalansowanie.
- Offsety należy zatwierdzać po wykonaniu pracy, aby uzyskać dostarczanie co najmniej raz; automatyczne zatwierdzanie należy wyłączyć, aby zachować kontrolę.
- Należy dostroić
linger.ms,batch.sizeicompression.type, a także zawsze wywoływaćflush().
Następny temat: połączenie tych podstawowych mechanizmów w niezawodne przepływy pracy sterowane zdarzeniami.
Często zadawane pytania
Czy lekcja „Apache Kafka z PHP” jest bezpłatna?
Tak — pełny tekst „Apache Kafka z PHP” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu PHP Academy, przejdź na CoddyKit PRO. Kurs PHP Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Apache Kafka z PHP”?
Przesyłanie strumieniowe zdarzeń o dużej przepustowości za pomocą Kafka Ćwiczysz PHP Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć PHP Academy?
Nie wymagamy żadnego doświadczenia. PHP Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 3 z 4.
Ile czasu zajmuje lekcja „Apache Kafka z PHP”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji PHP Academy?
Tak. Każda lekcja PHP Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Dlaczego warto używać komunikacji asynchronicznej
- Praca z RabbitMQ w PHP
- Apache Kafka z PHP
- Tworzenie przepływów pracy sterowanych zdarzeniami