Потапяне в Delta Lake: задължително приложение и еволюция на схемата

Здравей, Хабр! Представям ви превода на статията „Потапяйки в Delta Lake: Задание и Развитие на Картине“ автори Burak Yavuz, Brenner Heintz и Denny Lee, подготвил се като подготвка за старта на курса „Data Engineer“ от OTUS.

Потапяне в Delta Lake: задължително приложение и еволюция на схемата

Данните, както и нашият опит, постоянно се натрупват и развиват. За да не изоставаме, нашите ментални модели на света трябва да се адаптират към новите данни, част от които съдържат нови измерения — нови начини за наблюдение на неща, за които преди не сме имали представа. Тези ментални модели са малко различни от схемите на таблици, които определят как класифицираме и обработваме нова информация.

Това ни подтиква към въпроса за управлението на схемите. С времето, когато бизнес задачите и изискванията се променят, структурите на вашите данни също се променят. Delta Lake улеснява внедряването на нови измерения при промяна на данните. Потребителите имат достъп до проста семантика за управление на схемите на своите таблици. Тези инструменти включват принудително прилагане на схемата (Schema Enforcement), което защитава потребителите от непреднамерено замърсяване на таблиците с грешки или ненужни данни, както и еволюция на схемата (Schema Evolution), която позволява автоматично добавяне на нови колони със значими данни на подходящите места. В тази статия ще се задълбочим в употребата на тези инструменти.

Разбиране на схемите на таблиците

Всеки DataFrame в Apache Spark съдържа схема, която определя формата на данните, като типове данни, колони и метаданни. С Delta Lake схемата на таблицата се съхранява в формат JSON в журнала на транзакциите.

Какво е принудително прилагане на схемата?

Принудителното прилагане на схемата (Schema Enforcement), известно също като проверка на схемата (Schema Validation), е защитен механизъм в Delta Lake, който гарантира качеството на данните, отхвърляйки записи, които не отговарят на схемата на таблицата. Подобно на хостеса на рецепцията в популярния ресторант, който приема само по предварителна резервация, той проверява дали всеки колона от данните, въведени в таблицата, е в съответния списък на очакваните колони (с други думи, дали за всяка от тях има „резервация“) и отхвърля всякакви записи с колони, които не са в списъка.

Как работи принудителното прилагане на схемата?

Delta Lake използва проверка на схемата при запис, което означава, че всички нови записи в таблицата се проверяват за съвместимост със схемата на целевата таблица по време на записа. Ако схемата е несъвместима, Delta Lake изцяло отменя транзакцията (данните не се записват) и генерира изключение, за да информира потребителя за несъответствието.
За определяне на съвместимостта на записа с таблицата, Delta Lake използва следните правила. Записваният DataFrame:

  • не може да съдържа допълнителни колони, които не са в схемата на целевата таблица. И обратното, всичко е наред, ако входящите данни не съдържат абсолютно всички колони от таблицата — на тези колони просто ще бъдат присвоени нулеви стойности.
  • не може да има типове данни на колоните, които са различни от типовете данни на колоните в целевата таблица. Ако колона в целевата таблица съдържа данни от тип StringType, но съответстващата колона в DataFrame съдържа данни от тип IntegerType, принудителното прилагане на схемата ще предизвика изключение и ще попречи на изпълнението на операцията по запис.
  • не може да съдържа имена на колони, които се различават само по регистър. Това означава, че не можете да имате колони с имената ‘Foo’ и ‘foo’, дефинирани в една таблица. Въпреки че Spark може да се използва в режим, чувствителен или нечувствителен към регистъра (по подразбиране), Delta Lake запазва регистра, но е нечувствителен в рамките на съхранението на схемата. Parquet е чувствителен към регистъра при съхранение и извличане на информация за колоната. За да избегнем възможни грешки, повреди на данните или загуба на информация (с които лично сме се сблъсквали в Databricks), решихме да добавим това ограничение.

За да илюстрираме това, нека погледнем какво се случва в долния код при опит за добавяне на някои наскоро генерирани колони в таблица Delta Lake, която все още не е настроена да ги приема.

# Сгенерируем 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.

Вместо автоматично добавяне на нови колони, Delta Lake налага схемата и спира записа. За да помогне да се определи коя колона (или множество от тях) е причина за несъответствието, Spark извежда двете схеми от стека на грешките за сравнение.

Каква е ползата от принудителното прилагане на схемата?

Тъй като принудителното прилагане на схема представлява доста строг контрол, то е отличен инструмент за използване като врата за чист, напълно трансформиран набор от данни, готов за производство или потребление. Обикновено се прилага към таблици, които пряко подават данни:

  • Алгоритми за машинно обучение
  • BI табла за управление
  • Анализ на данни и инструменти за визуализация
  • Всяка производствена система, която изисква строго структурирани, строго типизирани семантични схеми.

За да подготвят данните си за тази финална бариера, много потребители използват проста архитектура на „много скокове“, която постепенно внася структура в таблиците им. За да научите повече за това, можете да се запознаете със статията Машинно обучение на производствено ниво с Delta Lake.

Разбира се, принудителното прилагане на схема може да се използва навсякъде във вашия пайплайн, но имайте предвид, че потоковото записване в таблица в такъв случай може да бъде разочароващо, защото например сте забравили, че сте добавили още една колона в входящите данни.

Предотвратяване на разреждането на данни

В този момент можете да се запитате, каква е такава истерия? В крайна сметка, понякога неочаквана грешка „несъответствие на схемата“ може да ви поднови в работния процес, особено ако сте новак в Delta Lake. Защо просто не позволите на схемата да се променя така, както е необходимо, за да мога да запиша своя DataFrame, независимо от всичко?

Както казва старата поговорка, „унция превенция струва фунт лечение“. В определен момент, ако не се погрижите за прилагането на схемата си, ще изникнат проблеми със съвместимостта на типовете данни - на пръв поглед хомогенни източници на необработени данни могат да съдържат гранични случаи, повредени колони, неправилно формулирани отображения или други страшни неща, които сънувате в кошмари. Най-добрият подход е да спирате тези врагове при портите - с помощта на принудителното прилагане на схема - и да се заемете с тях на светло, а не по-късно, когато те започнат да шетат в тъмните дълбини на работния ви код.

Принудителното прилагане на схемата осигурява увереност, че схемата на таблицата ви няма да се промени, освен ако сами не одобрите вариант за промяна. Това предотвратява "разреждането" на данните, което може да се случи, когато нови колони се добавят толкова често, че преди ценни, компресирани таблици губят стойността и полезността си поради наводняване с данни. Поощрявайки ви да бъдете целенасочени, да задавате високи стандарти и да очаквате високо качество, принудителното прилагане на схемата прави точно това, за което е предназначено — да ви помогне да останете добросъвестни, а таблиците ви — чисти.

Ако при по-нататъшно разглеждане решите, че всъщност необходимо искате да добавите нова колона — няма проблем, по-долу е посочен едноредовият фикс. Решението е еволюция на схемата!

Какво е еволюция на схемата?

Еволюцията на схемата е функция, която позволява на потребителите лесно да променят текущата схема на таблицата в съответствие с данните, които се променят с времето. Най-често тя се използва при извършване на операция по добавяне или презапис, за да се адаптира автоматично схемата за включване на една или няколко нови колони.

Как работи еволюцията на схемата?

Следвайки примера от предишния раздел, разработчиците могат лесно да използват еволюцията на схемата за добавяне на нови колони, които преди това са били отхвърлени поради несъответствие на схемата. Еволюцията на схемата се активира чрез добавяне на .option('mergeSchema', 'true') к вашия Spark команден ред .write или .writeStream.

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

За да видите графика, изпълнете следната 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

Потапяне в Delta Lake: задължително приложение и еволюция на схемата
Като алтернатива, можете да зададете тази опция за цялата Spark сесия, добавяйки spark.databricks.delta.schema.autoMerge = True в конфигурацията на Spark. Но бъдете внимателни, тъй като принудителното прилагане на схемата вече няма да ви предупреждава за непреднамерени несъответствия на схемата.

Включвайки параметъра в заявката mergeSchema, всички колони, които присъстват в DataFrame, но липсват в целевата таблица, автоматично се добавят в края на схемата в рамките на транзакцията за запис. Могат да бъдат добавени и вложени полета, и те също ще бъдат добавени в края на съответните колони на структурата.

Инженерите и учените могат да използват тази опция, за да добавят нови колони (например, наскоро наблюдавана метрика или колонка с показатели за продажби през този месец) в своите съществуващи производствени таблици за машинно обучение, без да нарушават съществуващите модели, основани на старите колони.

Следните типове промени в схемата са допустими в контекста на еволюцията на схемите при добавяне или презаписване на таблица:

  • Добавяне на нови колони (това е най-често срещаният сценарий)
  • Промяна на типовете данни от NullType -> всеки друг тип или повишаване от ByteType -> ShortType -> IntegerType

Други промени, които не са допустими в рамките на еволюцията на схемата, изискват схемата и данните да бъдат презаписани чрез добавяне .option("overwriteSchema", "true"). Например, ако колоната „Foo“ първоначално е била integer, а новата схема е от тип string, тогава всички Parquet файлове (данни) трябва да бъдат презаписани. Такива промени включват:

  • премахване на колона
  • промяна на типа данни на съществуваща колона (на място)
  • преименуване на колони, които се различават само по регистър (например „Foo“ и „foo“)

Накрая, с следващото издание Spark 3.0 ще бъде напълно поддържан явен DDL (чрез използване на ALTER TABLE), което ще позволи на потребителите да извършват следните действия върху схемите на таблиците:

  • добавяне на колони
  • промяна на коментарите към колоните
  • настояване на свойства на таблицата, определящи поведението на таблицата, например, задаване на продължителността на съхранение на журнал на транзакции.

Каква е ползата от еволюцията на схемата?

Еволюцията на схемата може да се използва всеки път, когато предвиждате да промените схемата на таблицата си (в контекста на случаи, когато случайно сте добавили колони в DataFrame, които не трябва да ги има). Това е най-простият начин за мигриране на схемата ви, тъй като автоматично добавя правилните имена на колони и типове данни, без да е необходимо да ги обявявате явност.

Заключение

Принудителното прилагане на схемата отхвърля всякакви нови колони или други изменения в схемата, които не са съвместими с вашата таблица. Като установяват и поддържат тези високи стандарти, анализаторите и инженерите могат да разчитат на изключително високо ниво на целостта на данните, което им позволява да вземат по-ефективни бизнес решения.

От друга страна, еволюцията на схемата допълва принудителното прилагане, опростявайки предположителните автоматични изменения на схемата. В крайна сметка, добавянето на колона не трябва да е сложно.

Принудителното прилагане на схемата е янь, докато еволюцията на схемата е инь. Когато се използват заедно, тези функции опростяват подавянето на шум и настройката на сигнала.

Също така искаме да благодарим на Мукул Мурти и Пранава Ананд за техния принос към тази статия.

Други статии от тази серия:

Поглед към Delta Lake: анализ на журнала на транзакциите

Пуснете видеото

Статии по темата

Машинно обучение на производствено ниво с Delta Lake

Какво е езеро от данни?

Научете повече за курса

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster