Pipeline danych i optymalizacja wyszukiwarki

Przetwarzanie danych i optymalizacja wyszukiwania w Scali

Nasz klient zmagał się z powolnymi i nieefektywnymi pipeline’ami, które nie nadążały za ogromnym napływem danych. Jego wyszukiwarka była zbyt wolna, zbyt mało elastyczna i niewystarczająco wydajna w stosunku do rosnących potrzeb. Pomimo wykorzystania Hive, HDFS i niestandardowych GlueJobs wydajność się nie poprawiała.

Nierównomierny rozkład danych był poważnym problemem — niektóre zadania Spark trwały znacznie dłużej niż inne, spowalniając przetwarzanie danych i zmniejszając wydajność pipeline’u. Po stronie wyszukiwania zapytania nie były zoptymalizowane. Mimo ustrukturyzowanych danych nieefektywne indeksowanie i partycjonowanie sprawiały, że wyszukiwanie trwało dłużej, niż oczekiwano.

Klient potrzebował w pełni zoptymalizowanego pipeline’u opartego na Scali, wykorzystującego strategie partycjonowania, indeksowania i buforowania, aby skrócić czas przetwarzania i przyspieszyć wyszukiwanie.

Kluczowe wyzwania w projekcie

Duże, nierównomiernie rozłożone zbiory danych

Pipeline’y klienta musiały codziennie pobierać, transformować i analizować terabajty danych, jednak nierównomierny rozkład danych sprawiał, że czas przetwarzania był nieprzewidywalny. Niektóre zadania kończyły się w kilka sekund, podczas gdy inne trwały znacznie dłużej. Zadania Spark nie były odpowiednio zbalansowane, co tworzyło wąskie gardła spowalniające cały workflow.

Szybkość wyszukiwania

Istniejąca wyszukiwarka nie była przystosowana do natychmiastowego zwracania wyników. Przeszukiwanie milionów rekordów trwało zbyt długo. Zapytania nie były zoptymalizowane, indeksowanie było ograniczone, a duże zbiory danych dodatkowo wydłużały czas odpowiedzi. Celem NannosTech było przebudowanie logiki wyszukiwania z wykorzystaniem inteligentnego partycjonowania, indeksowania i optymalizacji bazy danych.

Złożony stos technologiczny

Projekt musiał łączyć Scala, Spark, Python, Java, Hive, Airflow i AWS. Poszczególne komponenty przetwarzały dane z różną szybkością, a zarządzanie zależnościami pomiędzy wieloma frameworkami i rozwiązaniami do przechowywania danych dodatkowo zwiększało złożoność. Naszym zadaniem było więc połączenie wszystkich elementów w dobrze zorkiestrowany system obsługujący zarówno przetwarzanie w czasie rzeczywistym, jak i zadania wsadowe.

Masz podobne problemy w swoim projekcie?

Rozwiążmy je. Skontaktuj się z nami, a pomożemy Ci zbudować wydajne i skalowalne rozwiązanie danych, które sprawdzi się przy każdej skali.

Etapy wdrażania optymalizacji pipeline’u danych i wyszukiwarki

Pipeline przetwarzania danych

Pipeline’y klienta działały wolno, zadania były rozdzielane nierównomiernie, a operacje shuffle generowały wysokie koszty. Aby rozwiązać te problemy, przebudowaliśmy workflow przetwarzania danych, wykorzystując optymalizacje Scala Spark. Zastosowaliśmy salting, aby zrównoważyć rozkład danych. Wspólne partycjonowanie i broadcast joins ograniczyły zbędne przenoszenie danych, a buforowanie i odpowiednia alokacja zasobów zwiększyły ogólną szybkość. Usprawnienia te skróciły czas wykonywania zadań i poprawiły responsywność całego systemu.

Wyszukiwarka

Zwiększyliśmy szybkość wyszukiwania dzięki wielopoziomowemu indeksowaniu i zoptymalizowanym strategiom partycjonowania. System indeksuje dane przed wykonaniem zapytań, co skraca czas skanowania. Zintegrowaliśmy Manticore Search, aby umożliwić indeksowanie w czasie rzeczywistym i wyszukiwanie pełnotekstowe. Dodatkowo dostroiliśmy procedury składowane, dzięki czemu zapytania są wykonywane natychmiast nawet przy milionach rekordów w systemie.

Orkiestracja komponentów

Wykorzystaliśmy Apache Airflow do orkiestracji wszystkich komponentów pipeline’u. Każdy etap — pozyskiwanie danych, transformacja, indeksowanie i przechowywanie — był zarządzany za pomocą skierowanych grafów acyklicznych (DAG), dzięki czemu workflow stały się modułowe i skalowalne. Wbudowany monitoring Airflow pozwolił nam śledzić czas wykonania, wykrywać awarie i uruchamiać alerty umożliwiające szybkie reagowanie. W ten sposób nasz zespół przekształcił wcześniej nieprzewidywalny system w w pełni kontrolowany, wysokowydajny pipeline danych.

Bliższe spojrzenie na kompletny pipeline danych

Pozyskiwanie i przetwarzanie danych w pipeline’ie opartym na Scali

Nie można zbudować solidnego pipeline’u danych na danych niskiej jakości. Dlatego po pozyskaniu surowych danych musieliśmy je oczyścić. Zbudowaliśmy warstwę oczyszczania danych opartą na Scala Spark. Pipeline:

  • Wyodrębniał kluczowe wartości z surowych zbiorów danych.
  • Usuwał wartości null, poprawiał formaty i standaryzował pola.
  • Filtrował zbędne rekordy, aby skrócić czas przetwarzania.

Po zakończeniu tego etapu pipeline zawierał ustrukturyzowane dane wysokiej jakości, gotowe do transformacji.

Transformacja danych w pipeline’ie opartym na Scali

Zmapowaliśmy i przekształciliśmy surowe rekordy do ujednoliconego schematu, aby zapewnić kompatybilność w całym systemie. Obejmowało to:

  • Stosowanie reguł transformacji w celu oczyszczania i standaryzacji wartości.
  • Konwersję pól do spójnego formatu, aby różne zbiory danych mogły ze sobą współpracować.
  • Wdrożenie ustrukturyzowanego modelu danych upraszczającego wykonywanie zapytań i przechowywanie.

Po zakończeniu tego etapu dane były oczyszczone, ustrukturyzowane i zoptymalizowane.

Agregacja i łączenie danych w pipeline’ie opartym na Scali

Po transformacji różne zbiory danych zostały połączone w finalny zbiór za pomocą:

  • Operacji Group by umożliwiających logiczne grupowanie danych.
  • Łączenia wielu zbiorów danych na podstawie wcześniej zdefiniowanych kluczy.
  • Operacji Union służących do łączenia ustrukturyzowanych źródeł danych.

Na tym etapie pipeline zawierał ujednolicony zbiór danych zoptymalizowany pod kątem zapytań i gotowy do efektywnego przechowywania.

Ładowanie zbioru danych do systemów przechowywania

Klient potrzebował dostępu do ogromnych zbiorów danych bez opóźnień w wykonywaniu zapytań. Aby to osiągnąć, zoptymalizowaliśmy przechowywanie i wdrażanie danych poprzez:

  • Partycjonowanie i indeksowanie bazy danych w celu przyspieszenia wyszukiwania.
  • Wdrożenie procedur składowanych w MariaDB, MySQL i Manticore.
  • Integrację wyszukiwarki zapewniającej szybkie zwracanie wyników.

Na tym etapie system był gotowy do obsługi dużego obciążenia zapytaniami, wspierając analitykę w czasie rzeczywistym i szybkie pobieranie danych.

Technologie wykorzystane w projekcie

Scala

C#

Java

Apache Spark

Apache Hive

Airflow

MariaDB

Manticore

MySQL

AWS

Rezultaty: krótszy czas przetwarzania i szybsze wyszukiwanie

Pomogliśmy klientowi wyeliminować wąskie gardła, przyspieszyć wyszukiwanie i skalować infrastrukturę. Silnik Scala Spark obsługuje teraz przetwarzanie rozproszone, równoważąc nierównomierne obciążenia za pomocą saltingu, współpartycjonowania i strategii buforowania. Zapytania wyszukiwania są wykonywane natychmiast dzięki zoptymalizowanemu partycjonowaniu, indeksowaniu i procedurom składowanym w Manticore, MariaDB i MySQL. Airflow orkiestruje workflow, zapewniając wykonywanie zadań, automatyczny monitoring i wykrywanie błędów w czasie rzeczywistym. System jest obecnie skalowalny, odporny na awarie i zoptymalizowany pod kątem złożonych zapytań wykonywanych na dużych zbiorach danych.

  • Ponad 50% szybsze przetwarzanie dzięki optymalizacjom Spark
  • Czas odpowiedzi wyszukiwania poniżej jednej sekundy dzięki indeksowanym polom i partycjonowanym zapytaniom
  • Orkiestracja bez przestojów dzięki automatyzacji Apache Airflow
  • Zoptymalizowane pobieranie danych dzięki ustrukturyzowanym procedurom w MariaDB i MySQL
  • Skalowalność bez utraty wydajności, z obsługą milionów rekordów
O 50% szybsze przetwarzanie danych | wyniki projektu NannosTech

Skontaktuj się z nami

support@nannostech.com
+48889712077