Zanurzenie w Delta Lake: egzekwowanie schematu i ewolucja

Cześć, Habr! Przedstawiam wam tłumaczenie artykułu „Zanurzenie w Delta Lake: egzekwowanie schematu i ewolucja” autorzy Burak Yavuz, Brenner Heintz i Denny Lee, którzy przygotowali to w przededniu rozpoczęcia kursu Data Engineer od OTUS.

Zanurzenie w Delta Lake: egzekwowanie schematu i ewolucja

Dane, podobnie jak nasze doświadczenie, nieustannie się gromadzą i rozwijają. Aby nie pozostawać w tyle, nasze mentalne modele świata muszą dostosowywać się do nowych danych, z których niektóre zawierają nowe wymiary — nowe sposoby postrzegania rzeczy, o których wcześniej nie mieliśmy pojęcia. Te mentalne modele niewiele różnią się od schematów tabel, które określają, jak klasyfikujemy i przetwarzamy nowe informacje.

Prowadzi nas to do pytania o zarządzanie schematami. W miarę jak zmieniają się zadania biznesowe i wymagania, zmienia się również struktura twoich danych. Delta Lake ułatwia wdrażanie nowych wymiarów w miarę zmiany danych. Użytkownicy mają dostęp do prostych semantyk do zarządzania schematami swoich tabel. Narzędzia te obejmują egzekwowanie schematu (Schema Enforcement), które chroni użytkowników przed przypadkowym zanieczyszczeniem swoich tabel błędami lub zbędnymi danymi, a także ewolucję schematu (Schema Evolution), która pozwala automatycznie dodawać nowe kolumny z cennymi danymi w odpowiednich miejscach. W tym artykule zagłębimy się w wykorzystanie tych narzędzi.

Zrozumienie schematów tabel

Każdy DataFrame w Apache Spark zawiera schemat, który określa formę danych, taką jak typy danych, kolumny i metadane. Przy użyciu Delta Lake schemat tabeli jest przechowywany w formacie JSON wewnątrz dziennika transakcji.

Czym jest egzekwowanie schematu?

Egzekwowanie schematu (Schema Enforcement), znane również jako walidacja schematu (Schema Validation), jest mechanizmem ochronnym w Delta Lake, który zapewnia jakość danych, odrzucając rekordy, które nie pasują do schematu tabeli. Podobnie jak hostess w popularnej restauracji, która przyjmuje tylko rezerwacje, sprawdza, czy każda kolumna danych wprowadzanych do tabeli znajduje się na liście oczekiwanych kolumn (innymi słowy, czy dla każdej z nich istnieje „rezerwacja”), i odrzuca wszelkie rekordy z kolumnami, których nie ma na liście.

Jak działa egzekwowanie schematu?

Delta Lake stosuje weryfikację schematu podczas zapisu, co oznacza, że wszystkie nowe zapisy do tabeli są sprawdzane pod kątem zgodności z schematem docelowej tabeli w momencie zapisu. Jeśli schemat jest niezgodny, Delta Lake całkowicie anuluje transakcję (dane nie są zapisywane) i generuje wyjątek, aby powiadomić użytkownika o niezgodności.
Aby określić zgodność zapisu, Delta Lake stosuje następujące zasady. Zapisywany DataFrame:

  • nie może zawierać dodatkowych kolumn, których nie ma w schemacie docelowej tabeli. Z drugiej strony, wszystko jest w porządku, jeśli dane wejściowe nie zawierają wszystkich kolumn z tabeli — tym kolumnom po prostu zostaną przypisane wartości null.
  • nie może mieć typów danych kolumn, które różnią się od typów danych kolumn w docelowej tabeli. Jeśli kolumna w docelowej tabeli zawiera dane typu StringType, ale odpowiadająca jej kolumna w DataFrame zawiera dane typu IntegerType, wymuszenie schematu spowoduje wyjątek i uniemożliwi wykonanie operacji zapisu.
  • nie może zawierać nazw kolumn, które różnią się tylko wielkością liter. Oznacza to, że nie możesz mieć kolumn o nazwach 'Foo' i 'foo' zdefiniowanych w tej samej tabeli. Chociaż Spark może być używany w trybie czułym lub niewrażliwym na wielkość liter (domyślnie), Delta Lake zachowuje wielkość liter, ale jest niewrażliwy w kontekście przechowywania schematu. Parquet jest czuły na wielkość liter podczas przechowywania i zwracania informacji o kolumnach. Aby uniknąć możliwych błędów, uszkodzenia danych lub ich utraty (z czym osobiście się spotkaliśmy w Databricks), postanowiliśmy dodać to ograniczenie.

Aby to zilustrować, spójrzmy na to, co dzieje się w poniższym kodzie podczas próby dodania kilku nowo wygenerowanych kolumn do tabeli Delta Lake, która nie jest jeszcze skonfigurowana do ich przyjęcia.

# Сгенерируем DataFrame ссуд, который мы добавим в нашу таблицу Delta Lake
loans = sql("""
            SELECT addr_state, CAST(rand(10)*count as bigint) AS count,
            CAST(rand(10) * 10000 * count AS double) AS amount
            FROM loan_by_state_delta
            """)

# Вывести исходную схему DataFrame
original_loans.printSchema()

root
  |-- addr_state: string (nullable = true)
  |-- count: integer (nullable = true)
 
# Вывести новую схему DataFrame
loans.printSchema()
 
root
  |-- addr_state: string (nullable = true)
  |-- count: integer (nullable = true)
  |-- amount: double (nullable = true) # new column
 
# Попытка добавить новый DataFrame (с новым столбцом) в существующую таблицу
loans.write.format("delta") 
           .mode("append") 
           .save(DELTALAKE_PATH)

Returns:

A schema mismatch detected when writing to the Delta table.
 
To enable schema migration, please set:
'.option("mergeSchema", "true")'
 
Table schema:
root
-- addr_state: string (nullable = true)
-- count: long (nullable = true)
 
Data schema:
root
-- addr_state: string (nullable = true)
-- count: long (nullable = true)
-- amount: double (nullable = true)
 
If Table ACLs are enabled, these options will be ignored. Please use the ALTER TABLE command for changing the schema.

Zamiast automatycznie dodawać nowe kolumny, Delta Lake narzuca schemat i zatrzymuje zapis. Aby pomóc określić, która kolumna (lub wiele z nich) jest przyczyną niezgodności, Spark wypisuje oba schematy z stosu wywołań do porównania.

Jakie są korzyści z wymuszania schematu?

Ponieważ wymuszona aplikacja schematu stanowi dość ścisłą kontrolę, jest doskonałym narzędziem do wykorzystania jako bramka czystego, w pełni przekształconego zbioru danych, gotowego do produkcji lub konsumpcji. Zazwyczaj stosuje się ją do tabel, które bezpośrednio dostarczają dane:

  • Algorytmy uczenia maszynowego
  • Dashbordy BI
  • Analiza danych i narzędzia wizualizacji
  • Jakiejkolwiek systemie produkcyjnym, który wymaga ściśle ustrukturyzowanych, ściśle typowanych schematów semantycznych.

Aby przygotować swoje dane do tego ostatecznego bariery, wielu użytkowników korzysta z prostej architektury „multi-hop”, która stopniowo wprowadza strukturę do ich tabel. Aby dowiedzieć się więcej na ten temat, można zapoznać się z artykułem Uczenie maszynowe na poziomie produkcyjnym z Delta Lake.

Oczywiście wymuszona aplikacja schematu może być używana w dowolnym miejscu Twojego pipeline'u, ale pamiętaj, że strumieniowe zapisywanie do tabeli w takim przypadku może być frustrujące, ponieważ na przykład zapomniałeś, że dodałeś dodatkową kolumnę do danych wejściowych.

Zapobieganie rozcieńczaniu danych

Na tym etapie możesz zadać sobie pytanie, skąd ten cały rozgłos? W końcu czasami nieoczekiwany błąd „niezgodności schematu” może podstawić Ci nogę w Twoim przepływie pracy, zwłaszcza jeśli jesteś nowicjuszem w Delta Lake. Dlaczego po prostu nie pozwolić schematowi zmieniać się w sposób, w jaki muszę zapisać swój DataFrame, niezależnie od wszystkiego?

Jak mówi stare przysłowie, „uncja zapobiegania jest warta funta leczenia”. W pewnym momencie, jeśli nie zadbasz o zastosowanie swojego schematu, nieprzyjemne problemy związane z kompatybilnością typów danych wyjdą na jaw — na pierwszy rzut oka jednorodne źródła surowych danych mogą zawierać przypadki brzegowe, uszkodzone kolumny, źle sformułowane mapowania lub inne straszne rzeczy, które śnią się w koszmarach. Najlepszym podejściem jest powstrzymanie tych wrogów u bram — za pomocą wymuszonej aplikacji schematu — i zajęcie się nimi w świetle dziennym, a nie później, gdy zaczną węszyć w ciemnych zakamarkach Twojego kodu roboczego.

Wymuszone stosowanie schematu zapewnia pewność, że schemat Twojej tabeli nie zmieni się, chyba że sam potwierdzisz wariant zmiany. Zapobiega to 'rozcieńczeniu' danych, które może wystąpić, gdy nowe kolumny są dodawane tak często, że wcześniej cenne, skompresowane tabele tracą swoją wartość i użyteczność z powodu zalania danymi. Zachęcając Cię do bycia zamierzonym, ustalania wysokich standardów i oczekiwania wysokiej jakości, wymuszone stosowanie schematu spełnia dokładnie to, do czego zostało zaprojektowane — pomaga Ci pozostać rzetelnym, a Twoim tabelom — czystymi.

Jeśli przy dalszym rozważaniu zdecydujesz, że w rzeczywistości trzeba chcesz dodać nową kolumnę — nie ma problemu, poniżej znajduje się jednolinijkowe rozwiązanie. Rozwiązanie — ewolucja schematu!

Czym jest ewolucja schematu?

Ewolucja schematu to funkcja, która pozwala użytkownikom łatwo zmieniać aktualny schemat tabeli w odpowiedzi na zmieniające się dane w czasie. Najczęściej jest używana podczas wykonywania operacji dodawania lub zapisania, aby automatycznie dostosować schemat do uwzględnienia jednej lub kilku nowych kolumn.

Jak działa ewolucja schematu?

Podążając za przykładem z poprzedniej sekcji, programiści mogą łatwo używać ewolucji schematu do dodawania nowych kolumn, które wcześniej zostały odrzucone z powodu niezgodności ze schematem. Ewolucja schematu jest aktywowana przez dodanie .option('mergeSchema', 'true') do Twojej komendy Spark .write lub .writeStream.

# Добавьте параметр mergeSchema
loans.write.format("delta") 
           .option("mergeSchema", "true") 
           .mode("append") 
           .save(DELTALAKE_SILVER_PATH)

Aby zobaczyć wykres, wykonaj następujące zapytanie Spark SQL

# Создайте график с новым столбцом, чтобы подтвердить, что запись прошла успешно
%sql
SELECT addr_state, sum(`amount`) AS amount
FROM loan_by_state_delta
GROUP BY addr_state
ORDER BY sum(`amount`)
DESC LIMIT 10

Zanurzenie w Delta Lake: egzekwowanie schematu i ewolucja
Alternatywnie możesz ustawić tę opcję dla całej sesji Spark, dodając spark.databricks.delta.schema.autoMerge = True do konfiguracji Spark. Używaj jej jednak ostrożnie, ponieważ wymuszone stosowanie schematu nie będzie już ostrzegać Cię o niezamierzonych niezgodnościach ze schematem.

Włączając do zapytania parametr mergeSchema, wszystkie kolumny, które są obecne w DataFrame, ale brakują w docelowej tabeli, zostaną automatycznie dodane na końcu schematu w ramach transakcji zapisu. Mogą być także dodawane zagnieżdżone pola, które również zostaną dodane na końcu odpowiednich kolumn struktury.

Inżynierowie i naukowcy mogą korzystać z tej opcji, aby dodać nowe kolumny (np. ostatnio monitorowaną metrykę lub kolumnę wskaźników sprzedaży w tym miesiącu) do swoich istniejących tabel produkcyjnych dla uczenia maszynowego, nie naruszając istniejących modeli opartych na starych kolumnach.

Następujące typy zmian schematu są dozwolone w ramach ewolucji schematu podczas dodawania lub nadpisywania tabeli:

  • Dodawanie nowych kolumn (to najczęstszy scenariusz)
  • Zmiana typów danych z NullType -> jakikolwiek inny typ lub podnoszenie z ByteType -> ShortType -> IntegerType

Inne zmiany, niedozwolone w ramach ewolucji schematu, wymagają, aby schema i dane były nadpisane poprzez dodanie .option("overwriteSchema", "true"). Na przykład, w przypadku, gdy kolumna „Foo” początkowo była typu integer, a nowa schema miałaby typ danych string, wszystkie pliki Parquet (dane) musiałyby być nadpisane. Do takich zmian należą:

  • usunięcie kolumny
  • zmiana typu danych istniejącej kolumny (na miejscu)
  • zmiana nazw kolumn, które różnią się tylko wielkością liter (np. „Foo” i „foo”)

Na koniec, w następnej wersji Spark 3.0 będzie w pełni wspierane jawne DDL (przy użyciu ALTER TABLE), co umożliwi użytkownikom wykonywanie następujących działań na schematach tabel:

  • dodawanie kolumn
  • zmiana komentarzy do kolumn
  • konfigurowanie właściwości tabeli, definiujących zachowanie tabeli, na przykład ustawienie długości przechowywania dziennika transakcji.

Jakie są korzyści z ewolucji schematu?

Ewolucję schematu można stosować zawsze, gdy zamierzasz zmienić schemat swojej tabeli (w przeciwieństwie do sytuacji, w których przypadkowo dodałeś kolumny do swojego DataFrame, których tam nie powinno być). To najprostszy sposób na migrację swojej schemy, ponieważ automatycznie dodaje odpowiednie nazwy kolumn i typy danych bez konieczności ich jawnego deklarowania.

Podsumowanie

Wymuszone zastosowanie schematu odrzuca wszelkie nowe kolumny lub inne zmiany schematu, które nie są zgodne z Twoją tabelą. Ustanawiając i utrzymując te wysokie standardy, analitycy i inżynierowie mogą polegać na tym, że ich dane mają najwyższy poziom integralności, co pozwala im podejmować bardziej efektywne decyzje biznesowe.

Z drugiej strony, ewolucja schematu uzupełnia wymuszone zastosowanie, upraszczając domniemane automatyczne zmiany schematu. W końcu nie powinno to być skomplikowane — dodać kolumnę.

Wymuszone zastosowanie schematu to yin, podczas gdy ewolucja schematu to yang. Wspólnie te funkcje jak nigdy dotąd uproszczają tłumienie szumów i dostosowanie sygnału.

Chcielibyśmy również podziękować Mukulowi Murti i Pranavowi Anandowi za ich wkład w ten artykuł.

Inne artykuły z tej serii:

Zanurzenie w Delta Lake: rozpakowanie dziennika transakcji

Odtwarzaj wideo

Artykuły tematyczne

Uczenie maszynowe na poziomie produkcyjnym z Delta Lake

Czym jest jezioro danych?

Dowiedz się więcej o kursie

Ź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