0Pricing
PHP Academy · Lekcja

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 IDE

Publikowanie 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 — lz4 lub zstd znacznie ogranicza koszt sieciowy.
  • acks — all dla trwałości, 1 dla 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_level i 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.size i compression.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

  1. Dlaczego warto używać komunikacji asynchronicznej
  2. Praca z RabbitMQ w PHP
  3. Apache Kafka z PHP
  4. Tworzenie przepływów pracy sterowanych zdarzeniami
← Powrót do PHP Academy