Einführung in Delta Lake: Schema-Überprüfung und Evolution

Hallo, Habr! Ich präsentiere Ihnen die Übersetzung des Artikels «Einführung in Delta Lake: Schema-Überprüfung & Evolution» von Burak Yavuz, Brenner Heintz und Denny Lee, vorbereitet im Vorfeld des Kursstarts „Data Engineer“ von OTUS.

Einführung in Delta Lake: Schema-Überprüfung und Evolution

Daten, ebenso wie unsere Erfahrungen, sammeln sich kontinuierlich an und entwickeln sich weiter. Um nicht hinterherzuhinken, müssen sich unsere mentalen Modelle der Welt an neue Daten anpassen, von denen einige neue Dimensionen beinhalten — neue Ansätze, um Dinge zu beobachten, von denen wir vorher nichts wussten. Diese mentalen Modelle sind wenig anders als die Schemas von Tabellen, die festlegen, wie wir neue Informationen klassifizieren und verarbeiten.

Das führt uns zur Frage des Managements von Schemata. Während sich die geschäftlichen Anforderungen und Aufgaben im Laufe der Zeit ändern, verändert sich auch die Struktur Ihrer Daten. Delta Lake ermöglicht es, neue Dimensionen beim Ändern von Daten einfach zu implementieren. Benutzer haben Zugriff auf eine einfache Semantik zur Verwaltung der Schemata ihrer Tabellen. Zu diesen Werkzeugen gehören die Schema Enforcement-Funktion, die Benutzer vor unbeabsichtigter Verschmutzung ihrer Tabellen durch Fehler oder unerwünschte Daten schützt, sowie die Schema Evolution, die es ermöglicht, neue Spalten mit wertvollen Daten automatisch an den entsprechenden Stellen hinzuzufügen. In diesem Artikel werden wir uns eingehender mit der Nutzung dieser Werkzeuge befassen.

Verständnis der Tabellenschemata

Jeder DataFrame in Apache Spark enthält ein Schema, das die Struktur der Daten definiert, einschließlich Datentypen, Spalten und Metadaten. Mit Delta Lake wird das Tabellenschema im JSON-Format im Transaktionsprotokoll gespeichert.

Was ist Schema Enforcement?

Die erzwungene Anwendung des Schemas (Schema Enforcement), auch bekannt als Schema-Überprüfung (Schema Validation), ist ein Schutzmechanismus in Delta Lake, der die Datenqualität gewährleistet, indem er Datensätze, die nicht mit dem Tabellenschema übereinstimmen, ablehnt. Ähnlich wie eine Empfangsdame in einem beliebten Restaurant, die nur Reservierungen annimmt, überprüft es, ob jede Datenspalte, die in die Tabelle eingegeben wird, in der entsprechenden Liste der erwarteten Spalten vorhanden ist (mit anderen Worten, ob für jede von ihnen eine "Reservierung" besteht) und lehnt alle Datensätze mit Spalten ab, die nicht in dieser Liste stehen.

Wie funktioniert die erzwungene Anwendung des Schemas?

Delta Lake verwendet die Schema-Überprüfung beim Schreiben, was bedeutet, dass alle neuen Datensätze in die Tabelle zur Kompatibilität mit dem Schema der Ziel-Tabelle während des Schreibvorgangs überprüft werden. Wenn das Schema inkompatibel ist, hebt Delta Lake die Transaktion vollständig auf (die Daten werden nicht geschrieben) und generiert eine Ausnahme, um den Benutzer über die Diskrepanz zu informieren.
Um die Kompatibilität eines Datensatzes mit der Tabelle zu bestimmen, verwendet Delta Lake die folgenden Regeln. Der zu schreibende DataFrame:

  • kann keine zusätzlichen Spalten enthalten, die nicht im Schema der Zieltabelle vorhanden sind. Umgekehrt ist es in Ordnung, wenn die Eingabedaten nicht alle Spalten aus der Tabelle enthalten – diesen Spalten werden einfach Nullwerte zugewiesen.
  • kann keine Datentypen für Spalten haben, die sich von den Datentypen der Spalten in der Zieltabelle unterscheiden. Wenn eine Spalte der Zieltabelle Daten des Typs StringType enthält, die entsprechende Spalte im DataFrame jedoch Daten des Typs IntegerType hat, wird beim erzwungenen Schemaanwendung eine Ausnahme ausgelöst und die Schreiboperation verhindert.
  • kann keine Spaltennamen enthalten, die sich nur durch die Groß- und Kleinschreibung unterscheiden. Das bedeutet, dass Sie nicht gleichzeitig Spalten mit den Namen ‘Foo’ und ‘foo’ in einer Tabelle haben können. Obwohl Spark im kapitalsensitiven oder kapitalsensitiven (standardmäßig) Modus verwendet werden kann, bewahrt Delta Lake die Großschreibung, ist jedoch beim Speichern des Schemas kapitalsensitiv. Parquet ist beim Speichern und Abrufen von Spalteninformationen kapitalsensitiv. Um mögliche Fehler, Datenbeschädigungen oder -verluste zu vermeiden (mit denen wir persönlich bei Databricks konfrontiert waren), haben wir beschlossen, diese Einschränkung hinzuzufügen.

Um dies zu veranschaulichen, werfen wir einen Blick auf den folgenden Code, der versucht, einige neu generierte Spalten in eine Delta Lake-Tabelle hinzuzufügen, die noch nicht für deren Annahme konfiguriert ist.

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

Anstatt neue Spalten automatisch hinzuzufügen, erzwingt Delta Lake das Schema und stoppt die Aufzeichnung. Um festzustellen, welche Spalte (oder mehrere) die Ursache für die Diskrepanz ist, gibt Spark beide Schemata im Stack-Trace zur Vergleich an.

Was sind die Vorteile des erzwungenen Schemas?

Da das Erzwingen eines Schemas eine recht strenge Prüfung darstellt, ist es ein hervorragendes Werkzeug, um den Zugang zu einem sauberen, vollständig transformierten Datensatz zu ermöglichen, der bereit für die Produktion oder den Verbrauch ist. Es wird in der Regel auf Tabellen angewendet, die direkt Daten liefern:

  • Maschinenlernalgorithmen
  • BI-Dashboards
  • Datenanalysen und Visualisierungstools
  • Jeder Produktionssystem, das streng strukturierte, stark typisierte semantische Schemata erfordert.

Um Ihre Daten auf diese letzte Hürde vorzubereiten, verwenden viele Benutzer eine einfache "Multi-Hop"-Architektur, die nach und nach Struktur in ihre Tabellen bringt. Um mehr darüber zu erfahren, können Sie den Artikel lesen. Produktionsreifes maschinelles Lernen mit Delta Lake.

Natürlich kann die erzwungene Anwendung des Schemas überall in Ihrem Pipeline verwendet werden, aber denken Sie daran, dass das Streamen in eine Tabelle in diesem Fall frustrierend sein kann, insbesondere wenn Sie vergessen haben, eine zusätzliche Spalte zu den Eingabedaten hinzuzufügen.

Vermeidung von Datenverdünnung.

An diesem Punkt könnten Sie sich fragen, warum so viel Aufregung? Schließlich kann ein unerwarteter Fehler "Schema-Mismatch" Ihnen in Ihrem Workflow Probleme bereiten, insbesondere wenn Sie neu bei Delta Lake sind. Warum also nicht einfach das Schema so ändern, wie es benötigt wird, damit ich meinen DataFrame ohne weiteres speichern kann?

Wie das Sprichwort sagt: „Eine Unze Prävention ist ein Pfund Heilung wert“. Irgendwann, wenn Sie sich nicht um die Anwendung Ihres Schemas kümmern, werden sich die Probleme mit der Datenkompatibilität bemerkbar machen – vermeintlich homogene Quellen unstrukturierter Daten können Grenzfälle, beschädigte Spalten, fehlerhaft formatierte Darstellungen oder andere schreckliche Dinge enthalten, die einem in den Albträumen erscheinen. Der beste Ansatz ist, diese Feinde an der Pforte abzuwehren – durch die konsequente Durchsetzung des Schemas – und sich mit ihnen im Hellen auseinanderzusetzen, anstatt später, wenn sie beginnen, in den dunklen Tiefen Ihres Arbeitscodes herumzustochern.

Die erzwungene Anwendung des Schemas sorgt dafür, dass sich das Schema Ihrer Tabelle nicht ändert, es sei denn, Sie bestätigen selbst eine Änderungsoption. Dies verhindert die "Verdünnung" (dilution) von Daten, die auftreten kann, wenn neue Spalten so häufig hinzugefügt werden, dass zuvor wertvolle, komprimierte Tabellen an Bedeutung und Nützlichkeit durch die Datenflut verlieren. Indem Sie dazu ermutigt werden, absichtlich zu handeln, hohe Standards zu setzen und hohe Qualität zu erwarten, erfüllt die erzwungene Anwendung des Schemas genau ihren Zweck – sie hilft Ihnen, gewissenhaft zu bleiben und Ihre Tabellen sauber zu halten.

Sollten Sie bei weiterer Überlegung feststellen, dass Sie tatsächlich man muss eine neue Spalte hinzufügen möchten, ist das kein Problem, hier ist eine einzeilige Lösung. Die Lösung ist die Evolution des Schemas!

Was ist die Evolution des Schemas?

Die Schema-Evolution ist eine Funktion, die es Benutzern ermöglicht, das aktuelle Tabellenschema entsprechend den sich im Laufe der Zeit ändernden Daten einfach zu ändern. Sie wird am häufigsten bei Additions- oder Überschreibvorgängen verwendet, um das Schema automatisch anzupassen, sodass ein oder mehrere neue Spalten eingebunden werden.

Wie funktioniert die Schema-Evolution?

Entwickler können der Schema-Evolution mühelos folgen, indem sie neue Spalten hinzufügen, die zuvor aufgrund von Schema-Inkompatibilitäten abgelehnt wurden. Die Schema-Evolution wird aktiviert, indem man .option('mergeSchema', 'true') zu Ihrem Spark-Befehl .write oder .writeStream.

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

Um das Diagramm anzuzeigen, führen Sie die folgende Spark SQL-Abfrage aus:

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

Einführung in Delta Lake: Schema-Überprüfung und Evolution
Alternativ können Sie diese Option für die gesamte Spark-Sitzung festlegen, indem Sie spark.databricks.delta.schema.autoMerge = True in die Spark-Konfiguration einfügen. Seien Sie jedoch vorsichtig, da eine erzwungene Schemaanwendung Sie nicht mehr über unbeabsichtigte Schemainkonsistenzen informiert.

Indem Sie den Parameter mergeSchema, alle Spalten, die im DataFrame vorhanden sind, aber nicht in der Zieltabelle, werden automatisch am Ende des Schemas im Rahmen der Schreibtransaktion hinzugefügt. Auch verschachtelte Felder können hinzugefügt werden und werden ebenfalls am Ende der entsprechenden Spaltenstruktur hinzugefügt.

Dateningenieure und Wissenschaftler können diese Option nutzen, um neue Spalten (möglicherweise eine neu verfolgte Kennzahl oder einen Spaltennamen für die Verkaufszahlen diesen Monat) in ihre bestehenden Produktionsdatenbanken für maschinelles Lernen hinzuzufügen, ohne die bestehenden Modelle, die auf alten Spalten basieren, zu beeinträchtigen.

Die folgenden Arten von Schemaänderungen sind im Rahmen der evolutionären Schemaänderungen beim Hinzufügen oder Überschreiben von Tabellen zulässig:

  • Hinzufügen neuer Spalten (dies ist das häufigste Szenario)
  • Änderung von Datentypen von NullType -> jeden anderen Typ oder Hochstufung von ByteType -> ShortType -> IntegerType

Andere Änderungen, die im Rahmen der evolutionären Schemaänderungen nicht zulässig sind, erfordern, dass das Schema und die Daten durch Hinzufügen überschrieben werden .option("overwriteSchema", "true"). Zum Beispiel, wenn die Spalte „Foo“ ursprünglich vom Typ Integer war und das neue Schema einen Datentyp String hätte, müssten alle Parquet-Dateien (Daten) neu geschrieben werden. Zu solchen Änderungen gehören:

  • das Entfernen einer Spalte
  • die Änderung des Datentyps einer vorhandenen Spalte (in-place)
  • das Umbenennen von Spalten, die sich nur im Groß- und Kleinschreibungs unterscheiden (z. B. „Foo“ und „foo“)

Schließlich wird mit der nächsten Version von Spark 3.0 die explizite DDL (mittels ALTER TABLE) vollständig unterstützt, was es den Benutzern ermöglicht, folgende Aktionen an Tabellenschemata auszuführen:

  • das Hinzufügen von Spalten
  • das Ändern von Kommentaren zu Spalten
  • das Anpassen von Tabelleneigenschaften, die das Verhalten der Tabelle bestimmen, wie z. B. das Festlegen der Dauer der Transaktionsprotokollaufbewahrung.

Was sind die Vorteile der Schema-Evolution?

Schema-Evolution kann immer dann genutzt werden, wenn Sie beabsichtigen Ändern Sie das Schema Ihrer Tabelle (im Gegensatz zu den Fällen, in denen Sie versehentlich Spalten in Ihrem DataFrame hinzugefügt haben, die dort nicht sein sollten). Dies ist der einfachste Weg, Ihr Schema zu migrieren, da es automatisch die richtigen Spaltennamen und Datentypen hinzufügt, ohne dass diese explizit deklariert werden müssen.

Fazit

Die erzwungene Anwendung des Schemas weist alle neuen Spalten oder andere Schemaänderungen zurück, die nicht mit Ihrer Tabelle kompatibel sind. Indem hohe Standards gesetzt und aufrechterhalten werden, können Analysten und Ingenieure sicherstellen, dass ihre Daten die höchste Integrität aufweisen, was es ihnen ermöglicht, klar und präzise zu argumentieren, um effektivere Geschäftsentscheidungen zu treffen.

Andererseits ergänzt die Schemaentwicklung die erzwungene Anwendung und vereinfacht vorgesehene automatische Schemaänderungen. Schließlich sollte es unkompliziert sein, eine Spalte hinzuzufügen.

Die erzwungene Anwendung des Schemas ist das Yang, während die Schemaentwicklung das Yin ist. In Kombination erleichtern diese Funktionen mehr denn je die Unterdrückung von Rauschen und die Anpassung des Signals.

Wir möchten auch Mukul Murti und Pranav Anand für ihren Beitrag zu diesem Artikel danken.

Weitere Artikel aus dieser Serie:

Einführung in Delta Lake: Entpackung des Transaktionsprotokolls

Video abspielen

Artikel zu diesem Thema

Produktionsreifes maschinelles Lernen mit Delta Lake

Was ist ein Data Lake?

Erfahren Sie mehr über den Kurs

Quelle: habr.com

Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen 🔥 Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen | ProHoster