Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Usare una tabella di controllo per guidare un
Quando esegui la stessa elaborazione su molti input, come mercati, tabelle sorgente, clienti o partizioni di data, codificare quella lista nel tuo lavoro significa modificare il codice e ridistribuirla ogni volta che la lista cambia. Invece, memorizza la lista in una tabella di controllo che il lavoro legge durante l'esecuzione. Per aggiungere o rimuovere lavoro, aggiorni una riga nella tabella e la successiva esecuzione del lavoro riprende la modifica senza modifiche al lavoro stesso. Questo è un modello guidato dai metadati : i dati, non il codice, controllano ciò che il lavoro elabora.
Questo tutorial costruisce un lavoro che utilizza questo pattern sul dataset di esempio preinstallato di Wanderbricks, così puoi eseguirlo end-to-end senza creare alcun dato sorgente. Lo scenario è una piattaforma di affitto vacanze che esegue la stessa analisi dei prezzi per ogni segmento di proprietà (come Ski Resort o Urban Year-Round). Una tabella di controllo elenca i segmenti da analizzare, un task SQL legge quella tabella e un For each task esegue l'analisi una volta per segmento, in parallelo.
Come funziona
Il lavoro collega tre compiti in sequenza:
| Attività | Tipo | Funzionamento |
|---|---|---|
read_segments |
SQL | Legge la tabella di controllo e cattura le righe come un array JSON |
process_segments |
Per ogni | Itera sull'array delle righe, avviando il compito annidato una volta per riga |
run_segment_analysis |
Notebook o SQL (annidato all'interno For each) |
Viene eseguito una volta per riga, usando i valori di quella riga per analizzare un segmento di proprietà |
Il flusso è read_segments → process_segments → run_segment_analysis (una volta per gira). L'output del task SQL, un array JSON di oggetti riga, fluisce nel For each campo Inputs del task tramite il riferimento {{tasks.read_segments.output.rows}}dinamico del valore . Il For each compito poi passa i campi di ogni riga al compito annidato come parametri, disponibili come {{input.property_type}} e {{input.min_price}}.
Prerequisiti
- Uno spazio di lavoro Azure Databricks con permesso per creare job e notebook.
- Permesso per creare tabelle in Unity Catalog e permesso per creare uno schema in un catalogo (i
USE CATALOGprivilegi eCREATE SCHEMAdi e) per contenere la tabella di controllo. - Un SQL warehouse per eseguire i compiti SQL. Se non ne hai uno, vedi Creare un SQL warehouse.
- Il
samplescatalogo, disponibile in ogni spazio di lavoro abilitato per Unity Catalog. Il tutorial legge dasamples.wanderbricks.properties, quindi non c'è alcun dato sorgente da configurare.
Passo 1: Crea la tabella di controllo
La tabella di controllo è la fonte di verità per la lista dei segmenti che il tuo lavoro processa. Per cambiare ciò che fa il lavoro, aggiorni questa tabella, non il lavoro.
Esegui il seguente SQL in un notebook Azure Databricks o nell'editor SQL. La prima affermazione crea uno schema per contenere la tabella di controllo, la seconda crea la tabella con una riga per ogni segmento di proprietà e il prezzo minimo di listino da includere nell'analisi di quel segmento:
USE CATALOG <catalog-name>;
CREATE SCHEMA IF NOT EXISTS config;
CREATE OR REPLACE TABLE config.property_segments AS
SELECT * FROM VALUES
('Urban Year-Round', 150),
('Summer Getaway', 200),
('Ski Resort', 250)
AS t(property_type, min_price);
Sostituisci <catalog-name> con un catalogo in cui puoi creare schemi, come il catalogo del tuo spazio di lavoro. Usa lo stesso catalogo ovunque il tutorial faccia riferimento config.property_segments, inclusa la query di ricerca nel Passo 3.
Dopo questo passaggio, config.property_segments contiene tre righe, una per segmento. Ogni riga porta i due valori che il lavoro passa a ogni iterazione: il property_type to analyze e il min_price floor to filter.
Passo 2: Scrivi la logica di analisi
Il task annidato all'interno del For each task viene eseguito una volta per riga della tabella di controllo, ricevendo i parametri e come property_typemin_price di quella riga. Puoi scrivere questa logica come un compito di notebook o SQL. Scegli in base alla logica del tuo business:
- Usa un compito di notebook quando la logica per iterazione richiede codice procedurale, più linguaggi o librerie (ad esempio, un passaggio di data science o machine learning).
- Usa un compito SQL quando la logica è una singola query o trasformazione che puoi esprimere in modo dichiarativo. Un compito SQL richiede un SQL warehouse.
Entrambe le varianti qui sotto producono lo stesso risultato: per il segmento in processazione, il numero di inserzioni pari o superiori al prezzo minimo e il prezzo medio.
Attività del notebook
Creare un nuovo notebook in un percorso ad esempio /Workspace/Users/<username>/run_segment_analysis. Questo notebook viene eseguito una volta per ogni iterazione del For each compito, ricevendo un segmento diverso ogni volta.
Aggiungere il codice seguente al notebook:
# Set default values so you can run the notebook on its own while developing.
# When the notebook runs inside a For each task, the job overrides these defaults.
dbutils.widgets.text("property_type", "Ski Resort", "Property type")
dbutils.widgets.text("min_price", "250", "Minimum price")
# Read the parameters passed by the For each task.
property_type = dbutils.widgets.get("property_type")
min_price = dbutils.widgets.get("min_price")
result = spark.sql(
"""
SELECT :property_type AS property_type,
COUNT(*) AS property_count,
ROUND(AVG(base_price), 2) AS avg_price
FROM samples.wanderbricks.properties
WHERE property_type = :property_type
AND base_price >= :min_price
""",
args={"property_type": property_type, "min_price": min_price},
)
display(result)
Note
Chiamare dbutils.widgets.text() prima di dbutils.widgets.get(). Se chiami get per primo, eseguire il notebook fuori da un lavoro genera un InputWidgetNotDefined errore.
Attività SQL
Un task SQL esegue una query salvata, quindi crea e salva subito la query di analisi nell'editor SQL. Lo associ al task annidato quando configuri il For each task nel Passo 4.
Nel tuo workspace Azure Databricks, clicca
Nuovo>
Consulta per aprire l'editor SQL.
Immettere la query seguente. I compiti SQL fanno riferimento ai parametri con la
:param_namesintassi, quindi la query legge il suo segmento e il prezzo minimo dai:property_typeparametri e::min_priceSELECT :property_type AS property_type, COUNT(*) AS property_count, ROUND(AVG(base_price), 2) AS avg_price FROM samples.wanderbricks.properties WHERE property_type = :property_type AND base_price >= :min_price;Clicca sul titolo
New Query <date>nell'intestazione della scheda del tuo file SQL e dagli il nomerun_segment_analysis. Poi clicca su Salva per spostarlo in una cartella dove vuoi memorizzarlo.
Il For each compito passa i valori di ogni iterazione ai :property_type parametri e :min_price nominati in tempo reale. A differenza dei widget notebook, i parametri nominati SQL non supportano valori predefiniti: se un parametro non viene passato, la query fallisce con un errore di risoluzione dei parametri.
Passo 3: Crea la query di ricerca
Il compito di ricerca legge la tabella di controllo tramite una query salvata. Come nel Passo 2, crea e salva la query nell'editor SQL ora, poi allegala al compito di ricerca nel Passo 4.
Nel tuo workspace Azure Databricks, clicca
Nuovo>
Consulta per aprire l'editor SQL.
Inserisci quanto segue, utilizzando lo stesso catalogo scelto nel Passo 1:
SELECT property_type, min_price FROM <catalog-name>.config.property_segments;Il nome è pienamente qualificato perché il SQL warehouse che esegue questa query potrebbe riportare di default un catalogo diverso da quello in cui hai creato la tabella.
Clicca sul titolo
New Query <date>nell'intestazione della scheda del tuo file SQL e dagli il nomeread_segments. Poi clicca su Salva per spostarlo in una cartella dove vuoi memorizzarlo.
Passo 4: Crea e configura il lavoro
Con entrambe le query salvate, crea il lavoro e aggiungi i suoi due compiti: il compito di ricerca SQL che legge la tabella di controllo e il For each compito che esegue l'analisi per ogni riga.
Crea il lavoro
Nel tuo workspace Azure Databricks, nella barra laterale clicca Nuovo>
Lavoro. Assegna al lavoro un nome descrittivo, come
Segment Analysis.
Configura il compito di ricerca SQL
Questo compito legge la tabella di controllo e rende disponibili le sue righe eseguendo For each la read_segments query che hai salvato nel Passo 3.
- Clicca sulla tessera di query SQL per configurare il primo task. Se la tessera di query SQL non è disponibile, clicca su Aggiungi un altro tipo di attività e cerca query SQL.
- Impostare Nome attività su
read_segments. - Se necessario, seleziona query SQL dal menu a tendina Type .
- Nel campo query SQL , seleziona la
read_segmentsquery che hai salvato nel Passo 3. - Impostare SQL Warehouse su un warehouse nell'area di lavoro.
- Fare clic su Crea attività.
Quando questo compito viene eseguito, Azure Databricks cattura il risultato come un array JSON in tasks.read_segments.output.rows. L'output dei compiti SQL viene sempre restituito come array JSON, quindi non hai bisogno di configurazioni aggiuntive. La forma generale del riferimento è tasks.<task-name>.output.rows, dove <task-name> corrisponde al nome del compito che hai impostato. L'output è simile al seguente:
[
{ "property_type": "Urban Year-Round", "min_price": 150 },
{ "property_type": "Summer Getaway", "min_price": 200 },
{ "property_type": "Ski Resort", "min_price": 250 }
]
Configura il For each compito
L'attività For each legge l'output SQL e avvia un'esecuzione di attività nidificata per ogni riga.
Aggiungi compito e seleziona Per ciascuno.
Impostare Nome attività su
process_segments.Verifica che Depends su sia impostato su
read_segments.Nel campo Inputs , inserisci l'array di righe catturati dal compito SQL:
{{tasks.read_segments.output.rows}}Imposta la Concorrenza su
2per eseguire due iterazioni in parallelo. Aumentare questo valore quando l'attività nidificata supporta un livello di parallelismo maggiore.Per completare questo compito, clicca su Aggiungi un compito da eseguire in loop e configura il compito annidato che viene eseguito su ogni iterazione.
Il For each compito e il suo compito annidato sono creati insieme come un unico compito. Configura il compito annidato in base al tipo scelto nel Passo 2:
Attività del notebook
Impostare Nome attività su
run_segment_analysis.Imposta Tipo su Notebook.
Imposta il percorso del quaderno che hai creato nel Passo 2.
Clicca su Parametri, poi su Aggiungi per aggiungere ogni parametro:
-
Chiave:
property_type, Valore:{{input.property_type}} -
Chiave:
min_price, Valore:{{input.min_price}}
Ogni
{{input.<key>}}riferimento si risolve nel campo corrispondente della riga dell'iterazione corrente.-
Chiave:
Clicca su Crea compito per creare insieme il
For eachcompito e il suo compito annidato.
Attività SQL
Questo compito esegue la run_segment_analysis query che hai salvato nel Passo 2.
Impostare Nome attività su
run_segment_analysis.Imposta il Tipo su SQL, poi imposta il task SQL su Query.
Nel campo query SQL , seleziona la
run_segment_analysisquery che hai salvato nel Passo 2.Impostare SQL Warehouse su un warehouse nell'area di lavoro.
Clicca su Parametri, poi su Aggiungi per aggiungere ogni parametro:
-
Chiave:
property_type, Valore:{{input.property_type}} -
Chiave:
min_price, Valore:{{input.min_price}}
Ogni
{{input.<key>}}riferimento si risolve nel campo corrispondente della riga dell'iterazione corrente.-
Chiave:
Clicca su Crea compito per creare insieme il
For eachcompito e il suo compito annidato.
Il tuo lavoro Directed Acyclic Graph (DAG) ora mostra read_segments il flusso in process_segments, con il compito annidato all'interno del For each nodo.
Passo 5: Esegui il lavoro e verifica
- Fare clic su Esegui ora per attivare il processo.
- Seleziona la scheda Run per vedere la corsa. La prima esecuzione di un lavoro richiede alcuni minuti per iniziare il calcolo; Quando si completa, appare nella lista.
- Clicca sul
process_segmentsnodo per espandere ilFor eachcompito. - La pagina della run mostra una tabella delle iterazioni, una riga per segmento, ciascuna con il proprio stato, orario di inizio e durata.
- Clicca su qualsiasi riga di iterazione per aprire il suo output e confermare che ha analizzato il segmento atteso.
Puoi vedere i risultati di ogni iterazione in modo indipendente. Se una specifica iterazione fallisce, puoi rieseguire solo quell'iterazione dalla pagina di esecuzione del lavoro senza rieseguire l'intero lavoro.
Estendere il modello
Per aggiungere un segmento all'analisi, inserisci una riga nella tabella di controllo:
INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);
La successiva esecuzione del lavoro include il nuovo segmento, senza cambiamenti nella configurazione dei lavori o modifiche al notebook.
Lo stesso schema funziona in qualsiasi caso in cui si voglia che i dati guidino l'iterazione:
- Elaborazione per cliente: una riga per ID cliente. Il compito annidato applica trasformazioni specifiche per il cliente o consegna destinazioni specifiche per il cliente.
- Inserimento tabelle: una riga per ogni nome di tabella di origine. Il compito annidato legge e ingerisce ogni tabella.
- Elaborazione di riempimento retroattivo: una riga per ogni partizione di data. Il compito annidato rielabora i dati storici di quella partizione.
- Esecuzione guidata da feature flag: Una riga per ogni funzionalità o esperimento abilitato. Il compito annidato attiva la logica corrispondente.
Per fermare l'elaborazione di una riga senza cancellarla, aggiungi la tua colonna alla tabella di controllo (come una active flag) e filtra su di essa nel compito di ricerca SQL. Questa è una colonna ordinaria che definisci e popoli; Il For each compito non ha un concetto intrinseco. Prima aggiungi la colonna, poi imposta le righe esistenti a TRUE:
ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;
Poi filtra su di esso nella read_segments query in modo che solo le righe attive guidino l'iterazione:
SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;
Risorse aggiuntive
-
Usare un'attività
For eachper eseguire un'altra attività in un ciclo: riferimento completo per la configurazione delleFor eachattività, inclusi i tipi di parametri e le opzioni di concorrenza -
Usare una tabella di ricerca per matrici di parametri di grandi dimensioni in un'attività
For each: Come gestire matrici di parametri di grandi dimensioni che superano il limite di 48 KB del valore dell'attività - Accedere ai valori dei parametri da un'attività: tutti i metodi per l'accesso ai valori dei parametri nei notebook, negli script Python e nelle attività SQL
- Dataset Wanderbricks: Il dataset di esempio utilizzato in questo tutorial