Eintauchen in Delta Lake: Schema-Zwang und -Entwicklung

Hallo, Habr! Ich präsentiere Ihnen eine Übersetzung des Artikels. «Eintauchen in Delta Lake: Schema-Durchsetzung & Evolution» der Autoren Burak Yavuz, Brenner Heintz und Denny Lee, der im Vorfeld des Kursstarts vorbereitet wurde „Data Engineer“ von OTUS berichten werden.

Eintauchen in Delta Lake: Schema-Zwang und -Entwicklung

Daten, wie auch unsere Erfahrung, sammeln sich ständig an und entwickeln sich weiter. Um nicht zurückzufallen, müssen unsere mentalen Modelle der Welt an die neuen Daten angepasst werden, von denen einige neue Dimensionen enthalten — neue Wege, Dinge zu betrachten, von denen wir zuvor keine Vorstellung hatten. Diese mentalen Modelle unterscheiden sich kaum von den Schemas von Tabellen, die definieren, wie wir neue Informationen klassifizieren und verarbeiten.

Das führt uns zur Frage der Schema-Verwaltung. Während sich Geschäftsanforderungen und -fragen im Laufe der Zeit ändern, ändert sich auch die Struktur Ihrer Daten. Delta Lake ermöglicht es, neue Dimensionen einfach zu implementieren, wenn sich die Daten ändern. Die Nutzer haben Zugang zu einfacher Semantik zur Verwaltung der Schemata ihrer Tabellen. Diese Werkzeuge umfassen die forcierte Anwendung von Schemata (Schema Enforcement), die die Nutzer vor unbeabsichtigter Verunreinigung ihrer Tabellen durch Fehler oder unerwünschte Daten schützt, sowie die Schema-Evolution (Schema Evolution), die es ermöglicht, automatisch neue Spalten mit wertvollen Daten an den entsprechenden Stellen hinzuzufügen. In diesem Artikel werden wir uns eingehender mit der Verwendung dieser Werkzeuge befassen.

Verständnis der Tabellenschemata

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

Was ist die forcierte Anwendung von Schemata?

Die forcierte Anwendung von Schemata (Schema Enforcement), auch bekannt als Schema-Validierung (Schema Validation), ist ein Schutzmechanismus in Delta Lake, der die Datenqualität gewährleistet, indem er Aufzeichnungen ablehnt, die nicht mit dem Tabellenschema übereinstimmen. Ähnlich wie eine Hostess am Empfang in einem beliebten Restaurant, die nur Reservierungen annimmt, überprüft sie, ob jede Spalte von Daten, die in die Tabelle eingegeben werden, in der entsprechenden Liste der erwarteten Spalten enthalten ist (mit anderen Worten, ob für jede von ihnen eine 'Reservierung' besteht) und lehnt alle Aufzeichnungen mit Spalten ab, die nicht auf der Liste stehen.

Wie funktioniert die forcierte Anwendung von Schemata?

Delta Lake verwendet beim Schreiben eine Schemaüberprüfung, was bedeutet, dass alle neuen Einträge in die Tabelle während des Schreibvorgangs auf die Kompatibilität mit dem Schema der Zieltabelle überprüft werden. Wenn das Schema inkompatibel ist, wird die Transaktion von Delta Lake vollständig rückgängig gemacht (Daten werden nicht geschrieben) und eine Ausnahme generiert, um den Benutzer auf die Inkompatibilität hinzuweisen.
Zur Bestimmung der Kompatibilität eines Eintrags verwendet Delta Lake die folgenden Regeln. Der zu schreibende DataFrame:

  • darf keine zusätzlichen Spalten enthalten, die im Schema der Zieltabelle nicht vorhanden sind. Umgekehrt ist es in Ordnung, wenn die Eingabedaten nicht alle Spalten aus der Tabelle enthalten – diesen Spalten werden einfach Nullwerte zugewiesen.
  • darf keine Datentypen von Spalten haben, die von den Datentypen der Spalten in der Zieltabelle abweichen. Wenn eine Spalte der Zieltabelle Daten des Typs StringType enthält, die entsprechende Spalte im DataFrame jedoch Daten des Typs IntegerType hat, löst die erzwingende Anwendung des Schemas eine Ausnahme aus und verhindert die Ausführung des Schreibvorgangs.
  • darf keine Spaltennamen enthalten, die sich nur durch die Groß- und Kleinschreibung unterscheiden. Das bedeutet, dass Sie keine Spalten mit den Namen 'Foo' und 'foo' in einer Tabelle haben können. Obwohl Spark im empfindlichen oder nicht empfindlichen (Standard) Groß-/Kleinschreibungsmodus verwendet werden kann, bewahrt Delta Lake die Großschreibung, ist jedoch im Rahmen der Schemaablage nicht empfindlich. Parquet ist bei der Speicherung und Rückgabe von Spalteninformationen empfindlich. Um mögliche Fehler, Datenbeschädigungen oder -verluste (mit denen wir persönlich in Databricks konfrontiert waren) zu vermeiden, haben wir beschlossen, diese Einschränkung hinzuzufügen.

Um dies zu veranschaulichen, schauen wir uns an, was im folgenden Code passiert, wenn versucht wird, einige neu generierte Spalten zu einer 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.

Anstelle des automatischen Hinzufügens neuer Spalten erzwingt Delta Lake das Schema und stoppt das Schreiben. Um zu helfen, den Spalte (oder mehrere) zu bestimmen, die die Inkompatibilität verursacht, gibt Spark beide Schemata aus dem Stacktrace zur Vergleichszwecken aus.

Was ist der Nutzen der erzwingenden Anwendung eines Schemas?

Da die erzwungene Anwendung von Schemata eine recht strenge Überprüfung darstellt, ist sie ein ausgezeichnetes Werkzeug, um einen sauberen, vollständig transformierten Datensatz zu erhalten, der bereit für die Produktion oder den Verbrauch ist. In der Regel wird sie auf Tabellen angewendet, die Daten direkt bereitstellen:

  • Maschinenlernalgorithmen
  • BI-Dashboards
  • Datenanalyse und Visualisierungstools
  • Jedes Produktionssystem, das strikt strukturierte, streng typisierte semantische Schemata erfordert.

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

Natürlich kann die erzwungene Anwendung von Schemata überall in Ihrer Pipeline verwendet werden, aber denken Sie daran, dass das Streaming von Daten in eine Tabelle in diesem Fall frustrierend sein kann, weil Sie beispielsweise vergessen haben, dass Sie eine weitere Spalte zu den Eingangsdaten hinzugefügt haben.

Verhinderung von Datenverflüssigung

Zu diesem Zeitpunkt könnten Sie sich fragen, warum dieser ganze Aufruhr? Schließlich kann manchmal ein unerwarteter "Schema-Mismatch"-Fehler Ihrem Arbeitsprozess einen Strich durch die Rechnung machen, besonders wenn Sie neu bei Delta Lake sind. Warum sollten wir einfach nicht zulassen, dass sich das Schema so ändert, wie es notwendig ist, damit ich meinen DataFrame dennoch speichern kann?

Wie das alte Sprichwort sagt: „Eine Unze Prävention ist ein Pfund Heilung wert“. Irgendwann, wenn Sie sich nicht um die Anwendung Ihres Schemas kümmern, werden die Probleme mit der Kompatibilität der Datentypen hässliche Köpfe heben – auf den ersten Blick homogene Rohdatenquellen können Sonderfälle, beschädigte Spalten, fehlerhafte Zuordnungen oder andere schreckliche Dinge enthalten, die einem in Alpträumen erscheinen. Der beste Ansatz besteht darin, diese Feinde an den Toren abzufangen – durch die erzwungene Anwendung von Schemata – und sich mit ihnen im Licht zu befassen, anstatt später, wenn sie beginnen, in den dunklen Tiefen Ihres Arbeitscodes herumzustöbern.

Die zwangsweise Anwendung des Schemas gibt Ihnen die Gewissheit, dass sich das Schema Ihrer Tabelle nicht ändert, es sei denn, Sie bestätigen selbst eine Änderungsvariante. Dies verhindert die "Verdünnung" (dilution) von Daten, die entstehen kann, wenn neue Spalten so häufig hinzugefügt werden, dass zuvor wertvolle, komprimierte Tabellen aufgrund der Datenflut an Bedeutung und Nützlichkeit verlieren. Indem Sie dazu ermutigt werden, bewusst zu sein, hohe Standards zu setzen und eine hohe Qualität zu erwarten, sorgt die zwangsweise Anwendung des Schemas genau dafür, wozu sie gedacht war – Ihnen zu helfen, gewissenhaft zu bleiben und Ihre Tabellen sauber zu halten.

Wenn Sie bei weiterer Prüfung entscheiden, dass Sie tatsächlich wissen eine neue Spalte hinzufügen möchten - kein Problem, hier ist ein einzeiliger Fix. Die Lösung ist die Evolution des Schemas!

Was ist die Evolution des Schemas?

Die Evolution des Schemas ist eine Funktion, die es Benutzern ermöglicht, das aktuelle Tabellenschema leicht entsprechend den sich im Laufe der Zeit ändernden Daten zu ändern. Sie wird am häufigsten verwendet, um beim Hinzufügen oder Überschreiben automatisch das Schema anzupassen, um eine oder mehrere neue Spalten einzufügen.

Wie funktioniert die Evolution des Schemas?

Folgendes Beispiel aus dem vorherigen Abschnitt: Entwickler können die Evolution des Schemas einfach nutzen, um neue Spalten hinzuzufügen, die zuvor aufgrund von Schemaabweichungen abgelehnt wurden. Die Evolution des Schemas wird aktiviert, indem man .option('mergeSchema', 'true') zu Ihrem Spark-Befehl hinzufügen .write oder .writeStream.

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

Um das Schema 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

Eintauchen in Delta Lake: Schema-Zwang und -Entwicklung
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 aufnehmen. Aber verwenden Sie dies mit Vorsicht, da die zwangsweise Anwendung des Schemas Sie nicht mehr auf unbeabsichtigte Schemaabweichungen hinweist.

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

Dateningenieure und Wissenschaftler können diese Option nutzen, um neue Spalten (z. B. eine kürzlich verfolgte Metrik oder eine Spalte mit Verkaufszahlen für diesen Monat) in ihre bestehenden Produktions-Tabellen für maschinelles Lernen hinzuzufügen, ohne bestehende Modelle zu gefährden, die auf alten Spalten basieren.

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

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

Andere Änderungen, die im Rahmen der Schema-Evolution nicht zulässig sind, erfordern, dass Schema und Daten durch das Hinzufügen neu geschrieben werden .option("overwriteSchema", "true"). Zum Beispiel, wenn die Spalte „Foo“ ursprünglich ein 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:

  • Löschen einer Spalte
  • Ändern des Datentyps einer bestehenden Spalte (vor Ort)
  • Umbenennen von Spalten, die sich nur in der Groß- und Kleinschreibung unterscheiden (z. B. „Foo“ und „foo“)

Schließlich wird mit dem nächsten Release von Spark 3.0 das explizite DDL (mit Verwendung von ALTER TABLE) vollständig unterstützt, was es den Nutzern ermöglicht, die folgenden Aktionen an Tabellen-Schemata durchzuführen:

  • Hinzufügen von Spalten
  • Ändern von Kommentaren zu Spalten
  • Anpassen von Tabelleneigenschaften, die das Verhalten der Tabelle definieren, z. B. Festlegen der Aufbewahrungsdauer des Transaktionsprotokolls.

Was bringt die Schema-Evolution?

Schema-Evolution kann immer dann verwendet werden, wenn Sie beabsichtigen das Schema Ihrer Tabelle zu ändern (im Gegensatz zu Fällen, in denen Sie versehentlich Spalten zu Ihrem DataFrame hinzugefügt haben, die dort nicht sein sollten). Es ist der einfachste Weg, Ihr Schema zu migrieren, da es automatisch die richtigen Spaltennamen und Datentypen hinzufügt, ohne dass eine explizite Deklaration erforderlich ist.

Fazit

Die Zwangsanwendung des Schemas lehnt alle neuen Spalten oder andere Änderungen des Schemas ab, die nicht mit Ihrer Tabelle kompatibel sind. Durch die Festlegung und Einhaltung dieser hohen Standards können Analysten und Ingenieure sicher sein, dass ihre Daten ein Höchstmaß an Integrität aufweisen, was es ihnen ermöglicht, effizientere Geschäftsentscheidungen zu treffen.

Andererseits ergänzt die Evolution des Schemas die Zwangsanwendung und vereinfacht vorgesehene automatische Änderungen des Schemas. Letztendlich sollte es keine Komplexität sein — eine Spalte hinzuzufügen.

Die Zwangsanwendung des Schemas ist das Yang, während die Evolution des Schemas das Yin ist. Zusammen ermöglichen diese Funktionen wie nie zuvor, das Rauschen zu unterdrücken und das Signal anzupassen.

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

Weitere Artikel aus dieser Reihe:

Einführung in Delta Lake: Entschlüsselung des Transaktionsprotokolls

Video abspielen

Fachartikel

Maschinelles Lernen auf Unternehmensebene mit Delta Lake

Was ist ein Data Lake?

Erfahren Sie mehr über den Kurs

Quelle: habr.com

60GB SSD 8Gb DDR4