Ciao, Habr! Ti presento la traduzione di un articolo degli autori Burak Yavuz, Brenner Heintz e Denny Lee, preparato in vista dell'inizio del corso di OTUS.

I dati, come la nostra esperienza, continuano a crescere e a evolversi. Per non rimanere indietro, i nostri modelli mentali del mondo devono adattarsi ai nuovi dati, alcuni dei quali portano nuove dimensioni — nuovi modi di osservare cose di cui prima non avevamo idea. Questi modelli mentali sono poco diversi dagli schemi delle tabelle che definiscono come classifichiamo e gestiamo nuove informazioni.
Questo ci porta alla questione della gestione degli schemi. Con il cambiamento delle esigenze e delle richieste aziendali nel tempo, anche la struttura dei tuoi dati cambia. Delta Lake consente di implementare facilmente nuove dimensioni mentre i dati cambiano. Gli utenti hanno accesso a una semantica semplice per gestire gli schemi delle loro tabelle. Questi strumenti includono l'applicazione forzata dello schema (Schema Enforcement), che protegge gli utenti dall'inquinamento involontario delle loro tabelle per errori o dati indesiderati, e l'evoluzione dello schema (Schema Evolution), che consente di aggiungere automaticamente nuove colonne con dati preziosi nei posti giusti. In questo articolo, approfondiremo l'uso di questi strumenti.
Comprendere gli schemi delle tabelle
Ogni DataFrame in Apache Spark contiene uno schema che definisce la forma dei dati, come i tipi di dati, le colonne e i metadati. Con Delta Lake, lo schema della tabella è conservato in formato JSON all'interno del log delle transazioni.
Che cos'è l'applicazione forzata dello schema?
L'applicazione forzata dello schema (Schema Enforcement), nota anche come validazione dello schema (Schema Validation), è un meccanismo di protezione in Delta Lake che garantisce la qualità dei dati, rifiutando le registrazioni che non corrispondono allo schema della tabella. Proprio come un'inserviente alla reception di un ristorante popolare che accetta solo su prenotazione, verifica se ogni colonna di dati inseriti nella tabella è presente nell'elenco atteso delle colonne (in altre parole, se ogni colonna ha una 'prenotazione'), e rifiuta qualsiasi registrazione con colonne che non figurano nell'elenco.
Come funziona l'applicazione forzata dello schema?
Delta Lake utilizza la validazione dello schema durante la scrittura, il che significa che tutte le nuove registrazioni nella tabella vengono verificate per la compatibilità con lo schema della tabella di destinazione durante la scrittura. Se lo schema non è compatibile, Delta Lake annulla completamente la transazione (i dati non vengono scritti) e genera un'eccezione per avvisare l'utente della non conformità.
Per determinare la compatibilità di una registrazione con la tabella, Delta Lake utilizza le seguenti regole. Il DataFrame da scrivere:
- non può contenere colonne aggiuntive che non siano nello schema della tabella di destinazione. D'altra parte, va bene se i dati in ingresso non contengono tutte le colonne della tabella — a queste colonne verranno semplicemente assegnati valori null.
- non può avere tipi di dato delle colonne che differiscono dai tipi di dato delle colonne nella tabella di destinazione. Se la colonna nella tabella di destinazione contiene dati di tipo StringType, ma la colonna corrispondente nel DataFrame contiene dati di tipo IntegerType, l'applicazione forzata dello schema solleverà un'eccezione e impedirà l'esecuzione dell'operazione di scrittura.
- non può contenere nomi di colonne che differiscono solo per la maiuscola. Questo significa che non puoi avere colonne con i nomi 'Foo' e 'foo' definite nella stessa tabella. Sebbene Spark possa essere utilizzato in modalità sensibile o non sensibile (di default) alla maiuscola, Delta Lake conserva la maiuscola, ma è non sensibile nel contesto della memorizzazione dello schema. Parquet è sensibile alla maiuscola durante la memorizzazione e il recupero delle informazioni sulle colonne. Per evitare possibili errori, danneggiamenti dei dati o perdite (cosa con cui ci siamo trovati personalmente in Databricks), abbiamo deciso di aggiungere questa limitazione.
Per illustrare questo, diamo un'occhiata a cosa succede nel codice seguente quando si cerca di aggiungere alcune colonne recentemente generate a una tabella Delta Lake che non è ancora configurata per accettarle.
# Сгенерируем 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.Invece di aggiungere automaticamente le nuove colonne, Delta Lake impone lo schema e ferma la scrittura. Per aiutare a determinare quale colonna (o colonne) sta causando la non conformità, Spark stampa entrambi gli schemi dallo stack trace per il confronto.
Quali sono i vantaggi dell'applicazione forzata dello schema?
Poiché l'applicazione forzata dello schema è una verifica piuttosto rigorosa, è uno strumento eccellente da utilizzare come guardiano di un insieme di dati pulito e completamente trasformato, pronto per la produzione o il consumo. Di solito è applicato a tabelle che forniscono direttamente i dati:
- Algoritmi di apprendimento automatico
- Cruscotti BI
- Analisi dei dati e strumenti di visualizzazione
- Qualsiasi sistema produttivo che richieda schemi semantici rigorosamente strutturati e tipizzati.
Per preparare i propri dati a questa barriera finale, molti utenti utilizzano un'architettura semplice “multi-hop”, che introduce gradualmente struttura nelle loro tabelle. Per saperne di più, puoi consultare l'articolo
Naturalmente, l'applicazione forzata dello schema può essere utilizzata ovunque nel tuo pipeline, ma ricorda che la registrazione streaming in una tabella, in tal caso, può essere frustrante, poiché, ad esempio, potresti dimenticare di aver aggiunto un ulteriore campo ai dati in ingresso.
Prevenzione della diluizione dei dati
A questo punto potresti chiederti perché tale eccitazione? Dopotutto, a volte un errore inaspettato 'incoerenza dello schema' può ostacolare il tuo flusso di lavoro, specialmente se sei un novizio in Delta Lake. Perché non lasciare semplicemente che lo schema cambi come desidera, in modo che io possa registrare il mio DataFrame, a prescindere da tutto?
Come dice un vecchio proverbio, 'un'oncia di prevenzione vale un chilo di cura'. A un certo punto, se non ti prendi cura di applicare il tuo schema, i problemi di compatibilità dei tipi di dati alzeranno la loro testa brutta — fonti di dati non elaborati che sembrano omogenee possono contenere casi limite, colonne danneggiate, mapping malformati o altre cose spaventose che tormentano nei sogni. L'approccio migliore è fermare questi nemici alle porte — con l'applicazione forzata dello schema — e affrontarli alla luce del giorno, piuttosto che più tardi, quando iniziano a vagare nelle oscure profondità del tuo codice di lavoro.
L'applicazione forzata dello schema ti dà la certezza che lo schema della tua tabella non cambierà, a meno che non approvi tu stesso l'eventuale modifica. Questo previene la 'diluizione' dei dati, che può verificarsi quando nuovi campi vengono aggiunti così frequentemente che tabelle precedentemente preziose e compresse perdono il loro valore e utilità a causa di un'inondazione di dati. Incoraggiandoti a essere intenzionale, stabilire alti standard e aspettarti alta qualità, l'applicazione forzata dello schema fa esattamente ciò per cui è stata progettata: aiutarti a rimanere diligente e mantenere pulite le tue tabelle.
Se dopo ulteriori considerazioni decidi che ti piacerebbe veramente bisogno aggiungere un nuovo campo — nessun problema, di seguito troverai una soluzione in una riga. La soluzione è l'evoluzione dello schema!
Che cos'è l'evoluzione dello schema?
L'evoluzione dello schema è una funzione che consente agli utenti di modificare facilmente l'attuale schema di una tabella in base ai dati che cambiano nel tempo. Di solito viene utilizzata durante le operazioni di aggiunta o sovrascrittura, per adattare automaticamente lo schema all'inclusione di uno o più nuovi campi.
Come funziona l'evoluzione dello schema?
Seguendo l'esempio della sezione precedente, gli sviluppatori possono facilmente utilizzare l'evoluzione dello schema per aggiungere nuovi campi che erano stati precedentemente rifiutati a causa di incoerenze con lo schema. L'evoluzione dello schema viene attivata aggiungendo .option('mergeSchema', 'true') al tuo comando Spark .write o .writeStream.
# Добавьте параметр mergeSchema
loans.write.format("delta")
.option("mergeSchema", "true")
.mode("append")
.save(DELTALAKE_SILVER_PATH)Per visualizzare il diagramma, esegui la seguente query SQL di Spark
# Создайте график с новым столбцом, чтобы подтвердить, что запись прошла успешно
%sql
SELECT addr_state, sum(`amount`) AS amount
FROM loan_by_state_delta
GROUP BY addr_state
ORDER BY sum(`amount`)
DESC LIMIT 10 
In alternativa, puoi impostare questa opzione per l'intera sessione Spark aggiungendo spark.databricks.delta.schema.autoMerge = True nella configurazione di Spark. Ma usalo con cautela, poiché l'applicazione forzata dello schema non ti avviserà più di incoerenze involontarie con lo schema.
Abilitando il parametro nella query mergeSchema, tutte le colonne presenti nel DataFrame, ma assenti nella tabella di destinazione, vengono automaticamente aggiunte alla fine dello schema nell'ambito della transazione di scrittura. Anche i campi nidificati possono essere aggiunti, e verranno anch'essi aggiunti alla fine delle colonne di struttura corrispondenti.
Gli ingegneri e i ricercatori possono utilizzare questa opzione per aggiungere nuove colonne (come una nuova metrica tracciata di recente o una colonna delle vendite di questo mese) alle loro tabelle di produzione esistenti per il machine learning, senza interrompere i modelli esistenti basati sulle colonne più vecchie.
I seguenti tipi di modifiche allo schema sono ammesse nell'ambito dell'evoluzione dello schema durante l'aggiunta o la riscrittura di tabelle:
- Aggiunta di nuove colonne (questo è lo scenario più comune)
- Modifica dei tipi di dati da NullType -> qualsiasi altro tipo o elevazione da ByteType -> ShortType -> IntegerType
Altre modifiche non ammesse nell'ambito dell'evoluzione dello schema richiedono che lo schema e i dati siano riscritti tramite l'aggiunta .option("overwriteSchema", "true"). Ad esempio, nel caso in cui la colonna "Foo" fosse inizialmente un tipo intero e il nuovo schema sarebbe stato un tipo di dati stringa, allora tutti i file Parquet (dati) avrebbero dovuto essere riscritti. Tali modifiche includono:
- rimozione di una colonna
- modifica del tipo di dati di una colonna esistente (in loco)
- rinominare colonne che differiscono solo per maiuscole e minuscole (ad esempio, "Foo" e "foo")
Infine, con il rilascio successivo di Spark 3.0 sarà completamente supportato il DDL esplicito (utilizzando ALTER TABLE), permettendo agli utenti di eseguire le seguenti azioni sugli schemi delle tabelle:
- aggiunta di colonne
- modifica dei commenti sulle colonne
- configurazione delle proprietà della tabella che determinano il comportamento della tabella, come impostare la durata di conservazione del log delle transazioni.
Qual è il vantaggio dell'evoluzione dello schema?
L'evoluzione dello schema può essere utilizzata in qualsiasi momento in cui tu intenda modificare lo schema della tua tabella (a differenza dei casi in cui hai accidentalmente aggiunto colonne al tuo DataFrame che non dovrebbero esserci). Questo è il modo più semplice per migrare il tuo schema, poiché aggiunge automaticamente i nomi delle colonne e i tipi di dati corretti senza la necessità di dichiararli esplicitamente.
Conclusione
L'applicazione forzata dello schema rifiuta tutte le nuove colonne o altre modifiche allo schema che non sono compatibili con la tua tabella. Stabilendo e mantenendo questi elevati standard, analisti e ingegneri possono fare affidamento sul fatto che i loro dati hanno il massimo livello di integrità, riflettendo su di esso in modo chiaro e preciso, permettendo loro di prendere decisioni aziendali più efficaci.
D'altra parte, l'evoluzione dello schema integra l'applicazione forzata, semplificando modifiche automatiche presunte dello schema. In fin dei conti, non dovrebbe essere una complicazione — aggiungere una colonna.
L'applicazione forzata dello schema è lo yang, mentre l'evoluzione dello schema è lo yin. Usati insieme, queste funzionalità semplificano più che mai il filtraggio del rumore e la sintonizzazione del segnale.
Desideriamo anche ringraziare Mukula Murti e Pranav Anand per il loro contributo a questo articolo.
Altri articoli di questa serie:

Articoli correlati
Fonte: habr.com
