Jak Kafka stała się rzeczywistością

Jak Kafka stała się rzeczywistością

Cześć, Habr!

Pracuję w zespole Tinkoff, który zajmuje się tworzeniem własnego centrum powiadomień. Głównie programuję w Javie z użyciem Spring Boot i rozwiązuję różne problemy techniczne, które pojawiają się w projekcie.

Większość naszych mikrousług komunikuje się ze sobą asynchronicznie za pośrednictwem brokera wiadomości. Wcześniej jako broker używaliśmy IBM MQ, który przestał radzić sobie z obciążeniem, ale miał wysokie gwarancje dostawy.

W zamian zaproponowano nam Apache Kafka, która ma wysoką skalowalność, ale niestety wymaga prawie indywidualnego podejścia do konfiguracji w różnych scenariuszach. Ponadto mechanizm dostawy 'co najmniej raz', który działa w Kafka domyślnie, nie pozwalał utrzymać wymaganego poziomu spójności od razu. Następnie podzielę się naszym doświadczeniem z konfiguracją Kafki, w szczególności opowiem, jak skonfigurować i żyć z dostawą 'exactly once'.

Gwarantowana dostawa i nie tylko

Parametry, o których mowa dalej, pomogą zapobiec szeregowi problemów z domyślnymi ustawieniami połączenia. Ale najpierw chciałbym zwrócić uwagę na jeden parametr, który ułatwi ewentualne debugowanie.

W tym pomoże client.id dla Producenta i Konsumenta. Na pierwszy rzut oka jako wartość można użyć nazwy aplikacji, a w większości przypadków to zadziała. Chociaż sytuacja, gdy aplikacja używa kilku Konsumentów i nadaje im tę samą wartość client.id, prowadzi do następującego ostrzeżenia:

org.apache.kafka.common.utils.AppInfoParser — Błąd rejestracji AppInfo mbean javax.management.InstanceAlreadyExistsException: kafka.consumer:type=app-info,id=kafka.test-0

Jeśli chcesz użyć JMX w aplikacji z Kafką, może to być problem. W takim przypadku najlepiej użyć jako wartości client.id kombinacji nazwy aplikacji i, na przykład, nazwy tematu. Wynik naszej konfiguracji można zobaczyć w wynikach polecenia kafka-consumer-groups z narzędzi Confluent:

Jak Kafka stała się rzeczywistością

Teraz omówmy scenariusz gwarantowanej dostawy wiadomości. Producent Kafki ma parametr acks, który pozwala ustawić, po ilu potwierdzeniach lider klastra powinien uznać wiadomość za pomyślnie zapisaną. Ten parametr może przyjmować następujące wartości:

  • 0 — potwierdzenia nie będą wymagane.
  • 1 — domyślna wartość, wymaga potwierdzenia tylko od 1 repliki.
  • −1 — wymagane akceptacje od wszystkich synchronizowanych replik (konfiguracja klastra min.insync.replicas).

Z wymienionych wartości widać, że acks równe −1 zapewnia najsilniejsze gwarancje, że wiadomość nie zostanie utracona.

Jak wszyscy wiemy, systemy rozproszone są zawodne. Aby zabezpieczyć się przed tymczasowymi awariami, producenci Kafka oferują parametr retries, który pozwala określić liczbę prób ponownego wysłania w ciągu delivery.timeout.ms. Ponieważ parametr retries ma domyślną wartość Integer.MAX_VALUE (2147483647), liczbę powtórzeń wiadomości można kontrolować, zmieniając tylko delivery.timeout.ms.

Przechodzimy do exactly once delivery

Wymienione ustawienia pozwalają naszemu producenciowi dostarczać wiadomości z wysoką gwarancją. Teraz porozmawiajmy o tym, jak zapewnić zapis tylko jednej kopii wiadomości w topiku Kafka? W najprostszej wersji w producencie należy ustawić parametr enable.idempotence na wartość true. Idempotencja zapewnia zapis tylko jednej wiadomości w określonej partycji jednego topika. Wstępnym warunkiem włączenia idempotencji są wartości acks = all, retry > 0, max.in.flight.requests.per.connection ≤ 5. Jeśli te parametry nie są podane przez programistę, to automatycznie zostaną ustawione powyższe wartości.

Kiedy idempotencja jest skonfigurowana, należy zapewnić, aby te same wiadomości trafiały za każdym razem do tych samych partycji. Można to osiągnąć, konfigurując klucz i parametr partitioner.class w producencie. Zacznijmy od klucza. Dla każdej wysyłki musi on być taki sam. Można to łatwo osiągnąć, używając jakiegoś identyfikatora biznesowego z oryginalnej wiadomości. Parametr partitioner.class ma wartość domyślną — DefaultPartitioner. Przy tej strategii partycjonowania domyślnego działamy następująco:

  • Jeśli partycja jest jawnie określona podczas wysyłania wiadomości, to ją wykorzystujemy.
  • Jeśli partycja nie jest określona, ale klucz jest podany — wybieramy partycję według hasha od klucza.
  • Jeśli ani partycja, ani klucz nie są podane — wybieramy partycje na przemian (round-robin).

Ponadto, użycie klucza i idempotentne wysyłanie z parametrem max.in.flight.requests.per.connection = 1 zapewnia uporządkowane przetwarzanie wiadomości w Consumerze. Należy pamiętać, że jeśli w klastrze skonfigurowano zarządzanie dostępem, potrzebne będą uprawnienia do idempotentnego zapisu w temacie.

Jeśli brakuje ci możliwości idempotentnego wysyłania według klucza lub logika po stronie Producenta wymaga zachowania spójności danych między różnymi partycjami, z pomocą przychodzą transakcje. Ponadto, za pomocą transakcji łańcuchowej można warunkowo zsynchronizować zapis w Kafka, na przykład z zapisem w bazie danych. Aby włączyć transakcyjne wysyłanie, Producent musi mieć idempotentność i dodatkowo ustawić transactional.id. Jeśli w klastrze Kafka skonfigurowano zarządzanie dostępem, to dla zapisu transakcyjnego, tak jak dla idempotentnego, wymagane będą uprawnienia do zapisu, które mogą być przyznane za pomocą maski z użyciem wartości przechowywanej w transactional.id.

Formalnie jako identyfikator transakcji można używać dowolnego ciągu, na przykład nazwy aplikacji. Jednak jeśli uruchamiasz kilka instancji tej samej aplikacji z takim samym transactional.id, pierwsza uruchomiona instancja zostanie zatrzymana z błędem, ponieważ Kafka uzna ją za proces zombie.

org.apache.kafka.common.errors.ProducerFencedException: Producent próbował wykonać operację z przestarzałą epoką. Może istnieć nowszy producent o tym samym transactionalId lub transakcja producenta została wygaśnięta przez brokera.

Aby rozwiązać ten problem, dodajemy do nazwy aplikacji sufiks w postaci nazwy hosta, którą uzyskujemy z zmiennych środowiskowych.

Producent jest skonfigurowany, ale transakcje w Kafka zarządzają tylko zakresem widoczności wiadomości. Niezależnie od statusu transakcji, wiadomość trafia od razu do tematu, ale ma dodatkowe atrybuty systemowe.

Aby takie wiadomości nie były odczytywane przez Consumer’a przed czasem, musi on ustawić parametr isolation.level na wartość read_committed. Taki Consumer będzie mógł czytać nietransakcyjne wiadomości jak wcześniej, a transakcyjne — tylko po zatwierdzeniu.
Jeśli ustawiłeś wszystkie wymienione wcześniej ustawienia, to skonfigurowałeś exactly once delivery. Gratulacje!

Ale jest jeszcze jeden aspekt. Transactional.id, który konfigurowaliśmy wcześniej, w rzeczywistości jest prefiksem transakcji. Do menedżera transakcji dodawany jest numer porządkowy. Uzyskany identyfikator jest wydawany na transactional.id.expiration.ms, który jest konfigurowany na klastrze Kafka i domyślnie ma wartość „7 dni”. Jeśli przez ten czas aplikacja nie otrzyma żadnych wiadomości, to przy następnej próbie wysłania transakcji otrzymasz InvalidPidMappingException. Po tym koordynator transakcji wyda nowy numer porządkowy dla następnej transakcji. Wiadomość może być utracona, jeśli InvalidPidMappingException nie zostanie prawidłowo obsłużony.

Zamiast podsumowania

Jak można zauważyć, nie wystarczy po prostu wysyłać wiadomości do Kafka. Należy dobrać odpowiednią kombinację parametrów i być przygotowanym na wprowadzenie szybkich zmian. W tym artykule starałem się szczegółowo pokazać konfigurację exactly once delivery i opisałem kilka problemów z konfiguracjami client.id i transactional.id, z którymi się spotkaliśmy. Poniżej w skrócie przedstawione są ustawienia Producenta i Konsumenta.

Producent:

  1. acks = all
  2. retries > 0
  3. enable.idempotence = true
  4. max.in.flight.requests.per.connection ≤ 5 (1 — dla uporządkowanego wysyłania)
  5. transactional.id = ${application-name}-${hostname}

Konsument:

  1. isolation.level = read_committed

Aby zminimalizować błędy w przyszłych aplikacjach, stworzyliśmy nasz wrapper nad konfiguracją spring, gdzie już ustawione są wartości dla niektórych z wymienionych parametrów.

Oto kilka materiałów do samodzielnego przestudiowania:

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster