Практическое применение Spark schemaEvolution

Уважаеми читатели, добър ден!

В настоящата статия водещият консултант в бизнес направление Big Data Solutions на компания „Неофлекс“ подробно описва вариантите за изграждане на витрини с променлива структура, използвайки Apache Spark.

В рамките на проекта за анализ на данни често възниква задачата за изграждане на витрини на основата на слабо структурирани данни.

Обикновено това са логове или отговори от различни системи, запазвани във формат JSON или XML. Данните се изнасят в Hadoop, след което от тях трябва да се изградят витрини. Можем да организираме достъп до създадената витрина, например, чрез Impala.

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

Например, днес е регистриран такъв отговор:

{source: "app1", error_code: ""}

а утре от тази съща система идва такъв отговор:

{source: "app1", error_code: "error", description: "Network error"}

В резултат на това в витрината трябва да се добави още едно поле — description, и дали ще дойде или не, никой не знае.

Задачата за създаване на витрина с такива данни е доста стандартна и Spark разполага с редица инструменти за това. За парсиране на изходни данни има поддръжка и за JSON, и за XML, а за неизвестна предварително схема е предвидена поддръжка schemaEvolution.

От пръв поглед решението изглежда просто. Трябва да вземем папка с JSON и да я прочетем в dataframe. Spark ще създаде схема, а вложените данни ще бъдат преобразувани в структури. След това всичко трябва да се запази в parquet, който включително се поддържа и в Impala, регистрирайки витрината в Hive metastore.

На пръв поглед изглежда просто.

Въпреки това, от кратките примери в документацията не е ясно какво да се прави с редица проблеми на практика.

В документацията се описва подход, който не е за създаване на витрина, а за четене на JSON или XML в dataframe.

А именно, просто се посочва как да се прочете и разпарсира JSON:

df = spark.read.json(path...)

Това е достатъчно, за да направи данните достъпни за Spark.

На практика сценарият е много по-сложен, отколкото просто да се прочетат JSON файлове от папка и да се създаде dataframe. Ситуацията изглежда така: вече има определена витрина, всеки ден идват нови данни, които трябва да се добавят в витрината, без да забравяме, че схемата може да се различава.

Обикновената схема за изграждане на витрина изглежда така:

Стъпка 1. Данните се зареждат в Hadoop с последващо ежедневно обновяване и се складират в нова партиция. Получава се партиционирана по дни папка с изходни данни.

Стъпка 2. По време на инициализиращото зареждане тази папка се чете и парсира с помощта на Spark. Полученият dataframe се запазва в формат, достъпен за анализ, например в parquet, който след това може да се импортира в Impala. Така се създава целева витрина с всички данни, натрупани до този момент.

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

Да дадем пример. Да предположим, че е реализирана първата стъпка за изграждане на хранилище и е настроена изваждането на JSON файлове в папка.

Създаването на dataframe от тях, за да бъде запазен като витрина, не представлява проблем. Това е именно първата стъпка, която лесно може да бъде намерена в документацията на Spark:

df = spark.read.option("mergeSchema", True).json(".../*")
df.printSchema()

root
|-- a: long (nullable = true)
|-- b: string (nullable = true)
|-- c: struct (nullable = true)
|-- d: long (nullable = true)

Всичко изглежда наред.

Прочетохме и парснахме JSON, след което запазваме dataframe като parquet, регистрирайки в Hive по удобен начин:

df.write.format(“parquet”).option('path','').saveAsTable('')

Получаваме витрина.

Но на следващия ден нови данни са добавени от източника. Имаме папка с JSON и витрина, създадена на базата на тази папка. След зареждането на следващия пакет данни от източника, витрината не разполага с данни за един ден.

Логично решение е да партиционираме витрината по дни, което ще позволи на всеки следващ ден да добавя нова партиция. Този механизъм също е добре познат, Spark позволява записването на партиции отделно.

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

df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable)

На следващия ден зареждаме само новата партиция:

df.coalesce(1).write.mode("overwrite").parquet(dbpath + "\/" + db + "\/" + destTable +"\/date_load=" + date_load + "\/" )

Остава само да перерегистрираме в Hive, за да обновим схемата.
Въпреки това, тук и възникват проблеми.

Проблема номер едно. Рано или късно полученото parquet няма да може да бъде прочетено. Това се дължи на различния подход на parquet и JSON към празните полета.

Нека разгледаме типична ситуация. Например, вчера получаваме JSON:

Ден 1: {"a": {"b": 1}},

а днес същият JSON изглежда така:

Ден 2: {"a": null}

Да предположим, че имаме две различни партиции, в които има по един ред.
Когато четем изходните данни изцяло, Spark ще успее да определи типа и ще разбере, че "a" е поле от тип "структура", с вложено поле "b" от тип INT. Но ако всяка партиция е запазена поотделно, ще се получи parquet с несъвместими схеми на партициите:

df1 (a: <struct>)
df2 (a: STRING NULLABLE)

Тази ситуация е добре известна, затова е добавена специална опция — при парсинг на изходните данни да бъдат премахвани празните полета:

df = spark.read.json("...", dropFieldIfAllNull=True)

В такъв случай parquet ще се състои от партиции, които могат да бъдат прочетени заедно.
Въпреки това, тези, които са го правили на практика, ще се усмихнат горчиво. Защо? Защото вероятно ще възникнат още две ситуации. Или три. Или четири. Първата, която почти сигурно ще се случи, е, че числовите типове ще изглеждат различно в различните JSON файлове. Например, {intField: 1} и {intField: 1.1}. Ако такива полета попаднат в една партиция, мерж схемата ще ги прочете правилно, привеждайки към най-точния тип. Но ако са в различни, то в една ще има intField: int, а в друга intField: double.

За обработка на тази ситуация има следния флаг:

df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)

Сега имаме папка, в която се намират партиции, които могат да бъдат прочетени в единен dataframe и валиден parquet на цялата витрина. Да? Не.

Трябва да си спомним, че регистрирахме таблица в Hive. Hive не е чувствителен към регистъра в имената на полетата, а parquet е чувствителен. Затова партиции с схеми: field1: int, и Field1: int за Hive са одинакови, но за Spark не. Трябва да не забравим да приведем имената на полетата в долен регистър.

След това, изглежда, всичко е наред.

Въпреки това, не всичко е толкова просто. Възниква вторият, също така добре известен проблем. Тъй като всяка нова партиция се запазва отделно, в папката на партицията ще бъдат настанени служебни файлове на Spark, например, флаг за успешност на операцията _SUCCESS. Това ще доведе до грешка при опит за parquet. За да избегнем това, трябва да конфигурираме настройките, забранявайки на Spark да добавя служебни файлове в папката:

hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")

Изглежда, всеки ден в папката на целевата витрина се добавя нова партиция parquet, в която се съхраняват парсираните данни за деня. Ние предварително се погрижим да няма партиции с конфликт на типове данни.

Но пред нас стои третия проблем. Сега общата схема не е известна, освен това в Hive таблицата е с неправилна схема, тъй като всяка нова партиция вероятно е внесла изкривяване в схемата.

Необходимо е да регистрираме таблицата отново. Това може да се направи просто: отново да прочетем parquet витрините, да вземем схемата и да създадем на нейна основа DDL, с което отново да регистрираме папката в Hive като външна таблица, обновявайки схемата на целевата витрина.

Преди нас възниква четвърти проблем. Когато регистрирахме таблицата за първи път, се опирахме на Spark. Сега го правим сами и трябва да помним, че полетата parquet могат да започват с символи, недопустими за Hive. Например, Spark изхвърля редове, които не е успял да парсира в полето „corrupt_record“. Такова поле няма да може да бъде регистрирано в Hive без да бъде екранирано.

Знаейки това, получаваме схемата:

f_def = ""
for f in pf.dtypes:
  if f[0] != "date_load":
    f_def = f_def + "," + f[0].replace("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") 
table_define = "CREATE EXTERNAL TABLE jsonevolvtable (" + f_def[1:] + " ) "
table_define = table_define + "PARTITIONED BY (date_load string) STORED AS PARQUET LOCATION '\/user\/admin\/testJson\/testSchemaEvolution\/pq\/'"
hc.sql("drop table if exists jsonevolvtable")
hc.sql(table_define)

Код ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") прави безопасен DDL, тоест вместо:

create table tname (_field1 string, 1field string)

С такива имена на полета като „_field1, 1field“, се прави безопасен DDL, където имената на полета са екранирани: create table `tname` (`_field1` string, `1field` string).

Възниква въпросът: как правилно да получим dataframe с пълна схема (в кода pf)? Как да получим този pf? Това е петия проблем. Да пречетем схемата на всички партиции от папката с parquet файловете на целевата витрина? Това е най-сигурният метод, но и най-тежкия.

Схемата вече съществува в Hive. Новата схема може да бъде получена, като се комбинира схемата на цялата таблица с новата партиция. Значи трябва да вземем схемата на таблицата от Hive и да я комбинираме с новата партиция. Това може да се направи, като прочетем тестовите метаданни от Hive, запазим ги в временно папка и прочетем и двете партиции наведнъж с Spark.

По същество имаме всичко необходимо: оригиналната схема на таблицата в Hive и новата партиция. Досие също сме на разположение. Остава само да получим новата схема, в която се комбинира схемата на витрината с новите полета от създадената партиция:

from pyspark.sql import HiveContext
from pyspark.sql.functions import lit
hc = HiveContext(spark)
df = spark.read.json("...", dropFieldIfAllNull=True)
df.write.mode("overwrite").parquet(".../date_load=12-12-2019")
pe = hc.sql("select * from jsonevolvtable limit 1")
pe.write.mode("overwrite").parquet(".../fakePartiton/")
pf = spark.read.option("mergeSchema", True).parquet(".../date_load=12-12-2019/*", ".../fakePartiton/*")

След това създаваме DDL за регистрация на таблицата, както в предишния фрагмент.
Ако цялата верига работи правилно, а именно — е извършено инициализиране на зареждане и в Hive е създадена правилно оформена таблица, получаваме актуализирана схема на таблицата.

И последният проблем е, че не може просто да добавите партиция в таблицата Hive, тъй като тя ще се провали. Необходимо е да накарате Hive да поправи структурата на партициите:

from pyspark.sql import HiveContext
hc = HiveContext(spark) 
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)

Простата задача по четене на JSON и създаване на витрина на базата на него се изразява в преодоляване на редица неявни трудности, решенията на които трябва да се търсят поотделно. И въпреки че решенията са прости, намирането им отнема много време.

За да реализираме изграждането на витрина, трябваше:

  • Да добавим партиции в витрината, освобождавайки се от служебните файлове
  • Да се разберем с празните полета в изходните данни, които Spark е типизирал
  • Да приведем простите типове към низове
  • Да приведем имената на полетата към малки букви
  • Да разделим извеждането на данни и регистрацията на таблицата в Hive (създаване на DDL)
  • Да не забравим да екранираме имената на полетата, които могат да не са съвместими с Hive
  • Да научим как да актуализираме регистрацията на таблицата в Hive

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

Благодарим за прочитането на тази статия, надяваме се информацията да бъде полезна.

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

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