Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Przyp. tłum.: W tym artykule firma Banzai Cloud dzieli się przykładem wykorzystania swoich specjalnych narzędzi do uproszczenia zarządzania Kafka w ramach Kubernetes. Podane instrukcje ilustrują, jak określić optymalny rozmiar infrastruktury i skonfigurować Kafka w celu osiągnięcia wymaganej przepustowości.

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Apache Kafka to rozproszona platforma streamingu do budowy niezawodnych, skalowalnych i wydajnych systemów strumieniowych czasu rzeczywistego. Jej imponujące możliwości można rozszerzyć za pomocą Kubernetes. Dlatego opracowaliśmy operatora Kafka Open Source oraz narzędzie o nazwie Supertubes. Umożliwia uruchamianie Kafka w Kubernetes oraz korzystanie z jej różnych funkcji, takich jak precyzyjna konfiguracja brokera, skalowanie na podstawie metryk z rebalanceingiem, świadomość rozkładu zasobów (rack awareness), "miękkie" (graceful) wdrażanie aktualizacji itd.

Wypróbuj Supertubes w swoim klastrze:

curl https://getsupertubes.sh | sh oraz supertubes install -a --no-democluster --kubeconfig

Lub skontaktuj się z dokumentacji. Możesz także poczytać o niektórych funkcjonalnościach Kafka, których obsługa została zautomatyzowana za pomocą Supertubes i operatora Kafka. Już o nich pisaliśmy na blogu:

Decydując się na uruchomienie klastra Kafka w Kubernetes, na pewno napotkasz problem określenia optymalnego rozmiaru podstawowej infrastruktury oraz konieczność precyzyjnej konfiguracji Kafka, aby spełnić wymagania dotyczące przepustowości. Maksymalna wydajność każdego brokera zależy od wydajności komponentów infrastruktury, na której działa, takich jak pamięć, procesor, prędkość dysku, przepustowość sieci itd.

Idealnie, konfiguracja brokera powinna być taka, aby wszystkie elementy infrastruktury były wykorzystywane na maksymalnych możliwościach. Jednak w rzeczywistości taka konfiguracja jest dość trudna. Bardziej prawdopodobne jest, że użytkownicy będą konfigurować brokerów w taki sposób, aby maksymalizować wykorzystanie jednego lub dwóch komponentów (dysku, pamięci lub procesora). Ogólnie rzecz biorąc, broker pokazuje maksymalną wydajność, gdy jego konfiguracja pozwala na pełne wykorzystanie najwolniejszego komponentu. W ten sposób możemy uzyskać przybliżone wyobrażenie o obciążeniu, z którym jeden broker jest w stanie sobie poradzić.

Teoretycznie, możemy również oszacować liczbę brokerów potrzebnych do obsługi określonego obciążenia. Jednak w praktyce istnieje tak wiele wariantów konfiguracji na różnych poziomach, że oszacowanie potencjalnej wydajności określonej konfiguracji jest dość trudne (jeśli nie niemożliwe). Innymi słowy, bardzo trudno jest zaplanować konfigurację, opierając się na jakiejś określonej wydajności.

Dla użytkowników Supertubes zwykle stosujemy następujące podejście: zaczynamy od pewnej konfiguracji (infrastruktura + ustawienia), następnie mierzymy jej wydajność, dostosowujemy ustawienia brokera i powtarzamy proces jeszcze raz. Dzieje się to, dopóki potencjał najwolniejszego komponentu infrastruktury nie zostanie całkowicie wykorzystany.

W ten sposób uzyskujemy jaśniejszy obraz tego, ile brokerów jest potrzebnych klasterowi, aby poradzić sobie z określonym obciążeniem (liczba brokerów zależy również od innych czynników, takich jak minimalna liczba replik wiadomości dla zapewnienia niezawodności, liczba liderów partycji itp.). Ponadto uzyskujemy wyobrażenie o tym, dla którego komponentu infrastruktury pożądane byłoby pionowe skalowanie.

W tym artykule omówimy kroki, które podejmujemy, aby „wycisnąć maksimum” z najwolniejszych komponentów w początkowych konfiguracjach i zmierzyć przepustowość klastra Kafka. Wysoko niezawodna konfiguracja wymaga co najmniej trzech działających brokerów (min.insync.replicas=3), rozmieszczonych w trzech różnych strefach dostępności. Do konfiguracji, skalowania i monitorowania infrastruktury Kubernetes wykorzystujemy naszą własną platformę zarządzania kontenerami dla hybrydowych chmur — Pipeline. Wspiera ona on-premise (bare metal, VMware) oraz pięć typów chmur (Alibaba, AWS, Azure, Google, Oracle) i ich wszelkie kombinacje.

Przemyślenia dotyczące infrastruktury i konfiguracji klastra Kafka

W przedstawionych poniżej przykładach wybraliśmy AWS jako dostawcę usług chmurowych oraz EKS jako dystrybucję Kubernetes. Podobną konfigurację można zrealizować z wykorzystaniem PKE — dystrybucji Kubernetes od Banzai Cloud, certyfikowanej przez CNCF.

Dysk

Amazon oferuje różne typy wolumenów EBS. Na podstawie gp2 i io1 opierają się dyski SSD, jednak aby zapewnić wysoką przepustowość gp2 wykorzystuje zgromadzone kredyty (I/O credits), dlatego preferowaliśmy typ io1, który oferuje stabilną wysoką przepustowość.

Typy instancji

Wydajność Kafka mocno zależy od pamięci podręcznej strony systemu operacyjnego, dlatego potrzebujemy instancji z wystarczającą ilością pamięci dla brokerów (JVM) i pamięci podręcznej strony. Instancja c5.2xlarge — to dobry początek, ponieważ ma 16 GB pamięci i jest zoptymalizowana do pracy z EBS. Jej wadą jest to, że może zapewniać maksymalną wydajność przez nie więcej niż 30 minut w ciągu 24 godzin. Jeśli obciążenie wymaga maksymalnej wydajności przez dłuższy czas, warto zastanowić się nad innymi typami instancji. Tak właśnie zrobiliśmy, decydując się na c5.4xlarge. Zapewnia maksymalną przepustowość na poziomie 593,75 MB/s. Maksymalna przepustowość wolumenu EBS io1 jest wyższa niż dla instancji c5.4xlarge, dlatego najwolniejszym elementem infrastruktury zdaje się być przepustowość I/O tego typu instancji (co powinny również potwierdzić wyniki naszych testów obciążeniowych).

Sieć

Przepustowość sieci powinna być wystarczająco duża w porównaniu z wydajnością instancji VM i dysku, w przeciwnym razie sieć staje się wąskim gardłem. W naszym przypadku interfejs sieciowy c5.4xlarge obsługuje prędkość do 10 Gb/s, co znacznie przekracza przepustowość I/O instancji VM.

Rozmieszczanie brokerów

Brokerzy powinni być rozwijani (planowani w Kubernetes) na dedykowanych węzłach, aby uniknąć konkurencji z innymi procesami o zasoby CPU, pamięci, sieci i dysku.

Wersja Java

Logicznym wyborem jest Java 11, ponieważ jest zgodna z Dockerem w tym sensie, że JVM prawidłowo określa procesory i pamięć dostępne dla kontenera, w którym działa broker. Wiedząc, że limity procesorów są ważne, JVM wewnętrznie i przezroczysto ustawia liczbę wątków GC i wątków kompilatora JIT. Użyliśmy obrazu Kafka banzaicloud/kafka:2.13-2.4.0, zawierającego wersję Kafka 2.4.0 (Scala 2.13) na Java 11.

Jeśli chcesz dowiedzieć się więcej o Java/JVM w Kubernetes, zwróć uwagę na nasze publikacje:

Ustawienia pamięci brokera

Istnieją dwa kluczowe aspekty w konfiguracji pamięci brokera: ustawienia dla JVM i dla poda Kubernetes. Limit pamięci ustawiony dla poda powinien być większy niż maksymalny rozmiar heap, aby JVM miało miejsce dla metaprzestrzeni Java, która znajduje się w własnej pamięci, oraz dla pamięci podręcznej strony systemu operacyjnego, której Kafka aktywnie używa. W naszych testach uruchamialiśmy brokerów Kafka z parametrami -Xmx4G -Xms2G, a limit pamięci dla poda wynosił 10 Gi. Zauważ, że ustawienia pamięci dla JVM można uzyskiwać automatycznie za pomocą -XX:MaxRAMPercentage i -X:MinRAMPercentage, na podstawie limitu pamięci dla poda.

Ustawienia procesora brokera

Ogólnie rzecz biorąc, można poprawić wydajność, zwiększając równoległość poprzez zwiększenie liczby wątków używanych przez Kafka. Im więcej procesorów dostępnych dla Kafka, tym lepiej. W naszym teście zaczęliśmy od limitu 6 procesorów i stopniowo (iteracyjnie) zwiększyliśmy ich liczbę do 15. Ponadto ustawiliśmy num.network.threads=12 w ustawieniach brokera, aby zwiększyć liczbę wątków przyjmujących dane z sieci i je przesyłających. Natychmiast zauważyliśmy, że brokerzy-zwrotnicy nie mogą otrzymywać replik wystarczająco szybko, więc podnieśliśmy num.replica.fetchers do 4, aby zwiększyć prędkość, z jaką brokerzy-zwrotnicy replikowali wiadomości od liderów.

Narzędzie do generowania obciążenia

Należy upewnić się, że potencjał wybranego generatora obciążenia nie wyczerpie się przed osiągnięciem przez klaster Kafka maksymalnego obciążenia, na którym przeprowadzane są benchmarki. Innymi słowy, konieczne jest wcześniejsze ocenienie możliwości narzędzia generującego obciążenie oraz dobór instancji z odpowiednią liczbą procesorów i pamięci. W tym przypadku nasze narzędzie będzie generować większe obciążenie, niż klaster Kafka jest w stanie znieść. Po wielu eksperymentach zatrzymaliśmy się na trzech instancjach c5.4xlarge, w każdej z których uruchomiono generator.

Benchmarking

Pomiar wydajności to iteracyjny proces, który obejmuje następujące etapy:

  • konfiguracja infrastruktury (klaster EKS, klaster Kafka, narzędzie generujące obciążenie oraz Prometheus i Grafana);
  • generowanie obciążenia przez określony czas w celu filtrowania przypadkowych odchyleń w zbieranych wynikach wydajności;
  • dostosowanie infrastruktury i konfiguracji brokera na podstawie obserwowanych wyników wydajności;
  • powtarzanie procesu, aż osiągnięty zostanie wymagany poziom przepustowości klastra Kafka. Musi on być stabilnie reprodukowalny i wykazywać minimalne wahania przepustowości.

W następnym rozdziale opisano kroki, jakie zostały podjęte w trakcie benchmarku testowego klastra.

Narzędzia

Do szybkiego uruchomienia podstawowej konfiguracji, generowania obciążenia i pomiaru wydajności użyto następujących narzędzi:

  • Banzai Cloud Pipeline do zarządzania klastrem EKS od Amazon i Prometheus (do zbierania metryk Kafka i infrastruktury) oraz Grafana (do wizualizacji tych metryk). Skorzystaliśmy z zintegrowanych do Pipeline usług, które zapewniają federacyjne monitorowanie, centralizowane zbieranie logów, skanowanie podatności, przywracanie po awarii, bezpieczeństwo na poziomie korporacyjnym i wiele innych.
  • Sangrenel to narzędzie do testowania obciążenia klastra Kafka.
  • Pulpity Grafana do wizualizacji metryk Kafka i infrastruktury: Kubernetes Kafka, Node Exporter.
  • Supertubes CLI dla maksymalnie prostego konfigurowania klastra Kafka w Kubernetes. Zookeeper, operator Kafka, Envoy i wiele innych komponentów zostało zainstalowanych i odpowiednio skonfigurowanych do uruchomienia klastra Kafka gotowego do produkcji w Kubernetes.
    • klonujemy odpowiednie repozytorium z github: supertubes CLI skorzystaj z instrukcji podanych tutaj.

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Klaster EKS

Przygotuj klaster EKS z wydzielonymi węzłami roboczymi c5.4xlarge w różnych strefach dostępności dla podów z brokerami Kafka oraz dedykowanymi węzłami dla generatora obciążenia i infrastruktury monitorującej.

banzai cluster create -f https://raw.githubusercontent.com/banzaicloud/kafka-operator/master/docs/benchmarks/infrastructure/cluster_eks_202001.json

Gdy klaster EKS będzie gotowy, włącz jego zintegrowaną usługę monitorowania — ona wdroży Prometheus i Grafana w klastrze.

Komponenty systemowe Kafka

Zainstaluj komponenty systemowe Kafka (Zookeeper, kafka-operator) w EKS za pomocą supertubes CLI:

supertubes install -a --no-democluster --kubeconfig

Klaster Kafka

Domyślnie w EKS używane są wolumeny EBS typu gp2, dlatego należy utworzyć oddzielną klasę przechowywania na podstawie wolumenów io1 dla klastra Kafka:

kubectl create -f - <<EOF
apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
  name: fast-ssd
provisioner: kubernetes.io/aws-ebs
parameters:
  type: io1
  iopsPerGB: "50"
  fsType: ext4
volumeBindingMode: WaitForFirstConsumer
EOF

Ustaw parametr dla brokerów min.insync.replicas=3 i wdroż pod'y brokerów na węzłach w trzech różnych strefach dostępności:

supertubes cluster create -n kafka --kubeconfig  -f https://raw.githubusercontent.com/banzaicloud/kafka-operator/master/docs/benchmarks/infrastructure/kafka_202001_3brokers.yaml --wait --timeout 600

Tematy

Równocześnie uruchomiliśmy trzy egzemplarze generatora obciążenia. Każdy z nich pisze do swojego tematu, co oznacza, że potrzebujemy w sumie trzech tematów:

supertubes cluster topic create -n kafka --kubeconfig  -f -<<EOF
apiVersion: kafka.banzaicloud.io/v1alpha1
kind: KafkaTopic
metadata:
  name: perftest1
spec:
  name: perftest1
  partitions: 12
  replicationFactor: 3
  retention.ms: '28800000'
  cleanup.policy: delete
EOF

supertubes cluster topic create -n kafka --kubeconfig  -f -<<EOF
apiVersion: kafka.banzaicloud.io/v1alpha1
kind: KafkaTopic
metadata:
    name: perftest2
spec:
  name: perftest2
  partitions: 12
  replicationFactor: 3
  retention.ms: '28800000'
  cleanup.policy: delete
EOF

supertubes cluster topic create -n kafka --kubeconfig  -f -<<EOF
apiVersion: kafka.banzaicloud.io/v1alpha1
kind: KafkaTopic
metadata:
  name: perftest3
spec:
  name: perftest3
  partitions: 12
  replicationFactor: 3
  retention.ms: '28800000'
  cleanup.policy: delete
EOF

Dla każdego tematu współczynnik replikacji wynosi 3 — minimalnie zalecana wartość dla wysoko dostępnych systemów produkcyjnych.

Narzędzie do generowania obciążenia

Uruchomiliśmy trzy egzemplarze generatora obciążenia (każdy pisał w oddzielnym temacie). Aby pod’y generatora obciążenia zostały zaplanowane tylko na przydzielone węzły, należy określić affinność węzłów.

apiVersion: extensions/v1beta1
kind: Deployment
metadata:
  labels:
    app: loadtest
  name: perf-load1
  namespace: kafka
spec:
  progressDeadlineSeconds: 600
  replicas: 1
  revisionHistoryLimit: 10
  selector:
    matchLabels:
      app: loadtest
  strategy:
    rollingUpdate:
      maxSurge: 25%
      maxUnavailable: 25%
    type: RollingUpdate
  template:
    metadata:
      creationTimestamp: null
      labels:
        app: loadtest
    spec:
      affinity:
        nodeAffinity:
          requiredDuringSchedulingIgnoredDuringExecution:
            nodeSelectorTerms:
            - matchExpressions:
              - key: nodepool.banzaicloud.io/name
                operator: In
                values:
                - loadgen
      containers:
      - args:
        - -brokers=kafka-0:29092,kafka-1:29092,kafka-2:29092,kafka-3:29092
        - -topic=perftest1
        - -required-acks=all
        - -message-size=512
        - -workers=20
        image: banzaicloud/perfload:0.1.0-blog
        imagePullPolicy: Always
        name: sangrenel
        resources:
          limits:
            cpu: 2
            memory: 1Gi
          requests:
            cpu: 2
            memory: 1Gi
        terminationMessagePath: /dev/termination-log
        terminationMessagePolicy: File
      dnsPolicy: ClusterFirst
      restartPolicy: Always
      schedulerName: default-scheduler
      securityContext: {}
      terminationGracePeriodSeconds: 30

Kilka kwestii, na które należy zwrócić uwagę:

  • Generator obciążenia generuje wiadomości o długości 512 bajtów i publikuje je w Kafka w partiach po 500 wiadomości.
  • Za pomocą argumentu -required-acks=all publikacja jest uznawana za udaną, gdy wszystkie zsynchronizowane repliki wiadomości zostały odebrane i potwierdzone przez brokerów Kafka. Oznacza to, że w benchmarku mierzyliśmy nie tylko prędkość działania liderów otrzymujących wiadomości, ale także ich następców replikujących wiadomości. Celem tego testu nie jest ocena prędkości odczytu przez konsumentów (consumers) niedawno przyjętych wiadomości, które wciąż pozostają w pamięci podręcznej systemu operacyjnego, oraz jej porównanie z prędkością odczytu wiadomości przechowywanych na dysku.
  • Generator obciążenia równolegle uruchamia 20 worker'ów (-workers=20). Każdy worker zawiera 5 producerów, którzy wspólnie korzystają z połączenia worker'a z klastrem Kafka. W efekcie każdy generator ma 100 producerów, a wszyscy oni wysyłają wiadomości do klastra Kafka.

Obserwacja stanu klastra

Podczas testowania obciążeniowego klastra Kafka monitorowaliśmy również jego stan, aby upewnić się, że nie występują ponowne uruchomienia podów, niesynchronizowane repliki oraz maksymalna przepustowość przy minimalnych wahaniach.

  • Generator obciążenia zapisuje standardowe statystyki dotyczące liczby opublikowanych wiadomości i poziomu błędów. Procent błędów powinien pozostawać na poziomie 0,00%.
  • Cruise Control, wdrożony przez kafka-operatora, oferuje pulpit nawigacyjny, na którym możemy również obserwować stan klastra. Aby otworzyć ten pulpit, wykonaj:
    supertubes cluster cruisecontrol show -n kafka --kubeconfig
  • Poziom ISR (liczba replik „in-sync”) shrink i expansion wynosi 0.

Wyniki pomiarów

3 brokerów, rozmiar wiadomości — 512 bajtów

Z partycjami równomiernie rozłożonymi na trzech brokerach udało nam się osiągnąć wydajność ~500 Mb/s (około 990 tys. wiadomości na sekundę):

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Zużycie pamięci przez maszynę wirtualną JVM nie przekroczyło 2 GB:

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Przepustowość dysku osiągnęła maksymalną przepustowość I/O węzła na wszystkich trzech instancjach, na których działały brokery:

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Z danych dotyczących wykorzystania pamięci przez węzły wynika, że buforowanie systemowe i pamięć podręczna zajęły ~10-15 GB:

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

3 brokerów, rozmiar wiadomości — 100 bajtów

Zmniejszenie rozmiaru wiadomości powoduje spadek przepustowości o około 15-20%: wpływa to na czas potrzebny na przetworzenie każdej wiadomości. Ponadto obciążenie procesora wzrosło prawie dwukrotnie.

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Ponieważ na węzłach brokerów wciąż są niewykorzystane rdzenie, wydajność można zwiększyć poprzez zmianę konfiguracji Kafka. To nie jest proste zadanie, dlatego aby zwiększyć przepustowość lepiej pracować z wiadomościami o większym rozmiarze.

4 brokerów, rozmiar wiadomości — 512 bajtów

Można łatwo zwiększyć wydajność klastra Kafka, po prostu dodając nowych brokerów i utrzymując równowagę partycji (co zapewnia równomierne rozłożenie obciążenia między brokerami). W naszym przypadku po dodaniu brokera przepustowość klastra wzrosła do ~580 Mb/s (~1,1 mln wiadomości na sekundę). Wzrost okazał się mniejszy niż oczekiwano: głównie z powodu dysbalansu partycji (nie wszyscy brokerzy pracują na maksymalnej wydajności).

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Zarządzanie pamięcią JVM pozostało poniżej 2 GB:

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Praca brokerów z magazynami została dotknięta nierównowagą partycji:

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Określamy odpowiedni rozmiar klastra Kafka w Kubernetes

Wnioski

Przedstawione powyżej podejście iteracyjne można rozszerzyć, aby objąć bardziej złożone scenariusze, obejmujące setki konsumentów, repartycjonowanie, aktualizacje na gorąco, ponowne uruchamianie podów itp. Wszystko to pozwala nam ocenić granice możliwości klastra Kafka w różnych warunkach, zidentyfikować wąskie gardła w jego działaniu i znaleźć sposoby ich przezwyciężenia.

Opracowaliśmy Supertubes, aby szybko i łatwo wdrażać klaster, konfigurować go, dodawać/usuwać brokerów i tematy, reagować na powiadomienia oraz zapewniać prawidłowe działanie Kafka w Kubernetes w całości. Naszym celem jest pomoc w skupieniu się na głównym zadaniu („generować” i „konsumpować” komunikaty Kafka), podczas gdy całą ciężką pracę wykonają Supertubes i operator Kafka.

Jeśli interesują Cię technologie i projekty Open Source Banzai Cloud, śledź firmę na GitHub, LinkedIn lub Twitter.

P.S. od tłumacza

Przeczytaj także na naszym blogu:

Ź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