Dopasowywanie encji metodą probabilistyczną

Zautomatyzowany pipeline dopasowywania encji do łączenia rekordów

Klient posiadał tysiące rekordów pochodzących z różnych źródeł, a każdy z nich nieco się różnił. W jednym widniało „J. Smith”, w innym „John Smith”, a w kolejnym w ogóle brakowało imienia. Gdy pomnożyć to przez dziesiątki tysięcy osób, powstaje chaos. Nikt nie wiedział, który rekord jest poprawny ani czy ta sama osoba nie występuje w danych więcej niż raz. Prowadziło to do błędnych raportów, zduplikowanych wiadomości e-mail, utraconych informacji i niepotrzebnych kosztów.

Klient próbował rozwiązać problem za pomocą własnych algorytmów entity resolution. Opierały się one na prostych regułach — dokładnym dopasowaniu nazw, porównywaniu numerów telefonów i podobnych kryteriach. Działało to częściowo, ale tylko wtedy, gdy dane były czyste. System nie radził sobie z literówkami, brakującymi informacjami ani różnymi formatami. Co gorsza, przestawał działać wydajnie przy dużych zbiorach danych i nie nadawał się do ponownego wykorzystania. Każda kolejna modyfikacja utrudniała jego utrzymanie. Wtedy klient zwrócił się do nas.

Kluczowe wyzwania projektu dopasowywania encji

Niespójne formaty

W części rekordów widniało „Smith, John”, w innych „John S.” albo po prostu „J. Smith”. Czasami imienia nie było wcale. Do tego dochodziło kilka sposobów zapisu daty urodzenia, literówki, brakujące pola i niespójne formatowanie adresów. Normalizacja takich danych była złożonym zadaniem.

Złożoność dopasowywania rekordów

Nie mogliśmy po prostu założyć, że „jeśli imię = imię, to rekordy pasują”. Dane ze świata rzeczywistego tak nie działają. Ludzie zapisują informacje na różne sposoby, przeprowadzają się i zmieniają adresy e-mail. Zbyt niski próg dopasowania prowadziłby do łączenia rekordów dotyczących różnych osób, a zbyt wysoki powodowałby pomijanie prawdziwych dopasowań. Musieliśmy więc precyzyjnie dostroić próg.

Orkiestracja technologii

Część pipeline’u działała w Scali, a część w Pythonie. Całe rozwiązanie musiało działać w AWS, natomiast Airflow odpowiadał za wykonywanie wszystkich etapów zgodnie z harmonogramem, we właściwej kolejności i z wiarygodnymi logami. Awaria jednego elementu mogła zatrzymać cały proces lub doprowadzić do wygenerowania nieprawidłowych danych, dlatego orkiestracja musiała być niezawodna.

Wąskie gardła przy ładowaniu danych

Po oczyszczeniu i dopasowaniu danych trzeba było je załadować — mówimy tu o milionach rekordów każdego miesiąca. Wyzwaniem było przesłanie ich do MariaDB bez timeoutów, awarii ani blokowania bazy danych. Standardowy proces ładowania nie był wystarczająco wydajny, dlatego zbudowaliśmy inteligentniejszy, podzielony na partie proces ingestion, który można skalować bez destabilizowania systemu.

Masz podobne problemy w swoim projekcie?

Skontaktuj się z nami, aby omówić, jak możemy pomóc Ci usunąć wąskie gardła i skuteczniej dopasowywać encje.

Nasze podejście: techniki dopasowywania encji, które sprawdziły się w praktyce

Metodologia rozproszonego dopasowywania encji | NannosTech

Jako główną metodę dopasowywania wykorzystaliśmy probabilistyczne łączenie rekordów z użyciem Splink. W przeciwieństwie do systemów opartych na sztywnych regułach rozwiązanie oblicza prawdopodobieństwo, że dwa rekordy dotyczą tej samej encji — nawet jeśli poszczególne pola nie są identyczne. Zdefiniowaliśmy reguły porównywania imion i nazwisk, dat, numerów telefonów, adresów e-mail, historii zatrudnienia, edukacji i innych danych. Pozwoliło nam to wykrywać prawdziwe dopasowania, które inne rozwiązania pomijały, szczególnie przy niekompletnych i niespójnych danych. Dopasowanie opierało się na punktacji, a nie na decyzji binarnej. Podczas dostrajania ręcznie analizowaliśmy przypadki graniczne i wykorzystywaliśmy je do ulepszania kolejnych uruchomień modelu.

Normalizacja danych na potrzeby entity resolution | NannosTech

Przed dopasowaniem oczyściliśmy i ustandaryzowaliśmy każde pole. Imiona i nazwiska występowały w najróżniejszych formatach — „Doe, Jane”, „J. Doe”, „Jane D.” — dlatego usunęliśmy tytuły, ujednoliciliśmy wielkość liter i rozdzieliliśmy poszczególne części nazw. Daty zostały ujednolicone do formatu ISO, a numery telefonów przeformatowane. Historię edukacji i zatrudnienia znormalizowaliśmy przy użyciu list referencyjnych. Adresy przetworzyliśmy za pomocą zewnętrznego narzędzia do oczyszczania danych. Dzięki temu rekordy były porównywane na podstawie rzeczywistej treści.

Entity resolution z wykorzystaniem Spark | NannosTech

Zbudowaliśmy w pełni zautomatyzowany pipeline oparty na nowoczesnym, rozproszonym stosie technologicznym. Apache Spark obsługiwał operacje wymagające intensywnego przetwarzania danych i wykonywał probabilistyczne dopasowywanie za pomocą Splink. Airflow orkiestruje cały workflow — od pobierania surowych danych, przez uruchamianie zadań Spark i walidację wyników, aż po ładowanie finalnych zbiorów do MariaDB. Na każdym etapie zapisywane są logi wspierające śledzenie procesu i debugowanie.

Implementacja łączy Scala i Python, przy czym każdy z tych języków został wykorzystany w tej części pipeline’u, do której najlepiej pasował. Całość działa w AWS. Rezultatem jest powtarzalny i odporny na awarie system, który zgodnie z harmonogramem przetwarza miliony rekordów bez ręcznej ingerencji. Przykładowo, na klastrze składającym się z 8 maszyn (po 38 CPU i 418 GB RAM każda) system dopasował 400 mln rekordów w około 40 minut.

Technologie wykorzystane do stworzenia pipeline’u entity resolution

Scala

Python

Airflow

AWS

Splink

Spark

MariaDB

Cały proces przebiegał według ustrukturyzowanego schematu

1. Normalizacja

Oczyściliśmy i ustandaryzowaliśmy pola w całym zbiorze danych — imiona i nazwiska, daty, numery telefonów, stanowiska, historię edukacji oraz adresy. Ten etap ujednolicił formatowanie danych pochodzących z niespójnych źródeł i umożliwił ich wiarygodne porównywanie podczas dopasowywania.

2. Dopasowywanie

Korzystając ze Splink na Spark, zastosowaliśmy probabilistyczne łączenie rekordów na podstawie wielu pól. Każda para rekordów otrzymywała wynik dopasowania, a pary przekraczające próg 85% były traktowane jako duplikaty.

3. Tworzenie zbioru danych

Dopasowane rekordy zostały połączone w ujednolicone wpisy. Dodaliśmy metadane zapewniające pełną identyfikowalność, dzięki czemu każdy finalny rekord można prześledzić aż do jego źródeł.

4. Ładowanie danych

Finalne zbiory danych były ładowane do MariaDB za pomocą stworzonego przez nas równoległego narzędzia do ingestion, które obsługiwało duże wolumeny danych i pozwalało uniknąć wąskich gardeł wydajnościowych.

Rezultaty: 76% deduplikacji przy progu pewności ≥85%

System przetwarza miliony rekordów podczas każdego uruchomienia i działa w pełni automatycznie zgodnie z miesięcznym harmonogramem. Pipeline rozwiązał 76% duplikatów przy zastosowaniu wysokiego progu pewności — akceptowane były wyłącznie dopasowania z wynikiem co najmniej 85%.

  • Poziom deduplikacji: 76%
  • Dopasowania wymagały prawdopodobieństwa ≥85%
  • 400 mln rekordów dopasowanych w 40 minut na 8-węzłowym klastrze Spark
  • Wysoki potencjał ponownego wykorzystania: komponenty normalizacji i dopasowywania można zastosować w innych projektach po niewielkich modyfikacjach.
Przykłady dopasowywania tożsamości wraz z wynikami | NannosTech

Skontaktuj się z nami

support@nannostech.com
+48889712077