Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Verwende eine Steuertabelle, um einen
Wenn Sie dieselbe Verarbeitung über viele Eingaben ausführen, wie Märkte, Quelltabellen, Kunden oder Datumspartitionen, bedeutet das Hardcodieren dieser Liste in Ihrem Job, Code zu bearbeiten und jedes Mal neu zu deployen, wenn sich die Liste ändert. Speichere stattdessen die Liste in einer Steuertabelle , die der Job zur Laufzeit liest. Um Arbeit hinzuzufügen oder zu entfernen, aktualisiert man eine Zeile in der Tabelle, und der nächste Job-Run erkennt die Änderung ohne Änderungen am Job selbst. Dies ist ein metadatengetriebenes Muster: Die Daten, nicht der Code, steuern, was der Job verarbeitet.
Dieses Tutorial erstellt einen Job, der dieses Muster auf dem vorinstallierten Wanderbricks-Beispieldatensatz verwendet, sodass man es von Ende zu Ende ausführen kann, ohne Quelldaten zu erstellen. Das Szenario ist eine Urlaubsvermietungsplattform, die für jedes Immobiliensegment dieselbe Preisanalyse durchführt (wie Ski Resort oder Urban Year-Round). Eine Kontrolltabelle listet die zu analysierenden Segmente auf, eine SQL-Aufgabe liest diese Tabelle, und eine For each Aufgabe führt die Analyse pro Segment einmal parallel durch.
So funktioniert es
Die Aufgabe verbindet drei Aufgaben in der Reihenfolge:
| Aufgabe | Typ | Was es tut |
|---|---|---|
read_segments |
SQL | Liest die Kontrolltabelle und erfasst die Zeilen als JSON-Array |
process_segments |
Für jede | Iteriert über das Zeilenarray und startet die verschachtelte Aufgabe einmal pro Zeile |
run_segment_analysis |
Notebook oder SQL (innen For eachverschachtelt) |
Wird einmal pro Zeile ausgeführt und verwendet die Werte dieser Zeile, um ein Eigenschaftssegment zu analysieren |
Der Durchfluss beträgt read_segments → process_segments → run_segment_analysis (einmal pro Reihe). Die Ausgabe der SQL-Aufgabe, ein JSON-Array von Zeilenobjekten, fließt über die dynamische For eachin das Eingabefeld der {{tasks.read_segments.output.rows}} Aufgabe . Die Aufgabe For each übergibt dann die Felder jeder Zeile als Parameter an die verschachtelte Aufgabe, verfügbar als {{input.property_type}} und {{input.min_price}}.
Voraussetzungen
- Ein Azure Databricks-Arbeitsbereich mit Berechtigung zum Erstellen von Jobs und Notizbüchern.
- Berechtigung, Tabellen im Unity-Katalog zu erstellen, und Berechtigung, ein Schema in einem Katalog zu erstellen (die
USE CATALOGundCREATE SCHEMAPrivilegien), um die Kontrolltabelle zu halten. - Ein SQL-Warehouse, um die SQL-Aufgaben auszuführen. Falls du keines hast, siehe Create a SQL Warehouse.
- Der Katalog
samples, der in jedem von Unity Catalog aktivierten Arbeitsbereich verfügbar ist. Das Tutorial liest vonsamples.wanderbricks.properties, also gibt es keine Quelldaten zum Einrichten.
Schritt 1: Erstellen Sie die Kontrolltabelle
Die Kontrolltabelle ist die Quelle der Wahrheit für die Liste der Segmente, die Ihre Arbeit verarbeitet. Um zu ändern, was der Job macht, aktualisiert man diese Tabelle, nicht den Job.
Führe folgendes SQL in einem Azure Databricks-Notebook oder im SQL-Editor aus. Die erste Anweisung erstellt ein Schema zur Kontrolltabelle, und die zweite erstellt die Tabelle mit einer Zeile pro Immobiliensegment und dem Mindestlistenpreis, der in die Analyse dieses Segments aufgenommen werden soll:
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);
Ersetzen Sie <catalog-name> sie durch einen Katalog, in dem Sie Schemata erstellen können, wie zum Beispiel Ihren Arbeitsbereichskatalog. Verwenden Sie überall denselben Katalog, auf den das Tutorial verweist config.property_segments, einschließlich der Suchanfrage in Schritt 3.
Nach diesem Schritt config.property_segments enthält er drei Zeilen, eine pro Segment. Jede Zeile enthält die zwei Werte, die der Job an jede Iteration weitergibt: die property_type zum Analysieren und die zum Filtern auf der min_price Etage.
Schritt 2: Schreiben Sie die Analyselogik
Die verschachtelte Aufgabe innerhalb der For each Aufgabe wird pro Zeile der Kontrolltabelle einmal ausgeführt und erhält die Parameter dieser Zeile property_type und min_price als Parameter. Du kannst diese Logik als Notizbuchaufgabe oder SQL-Aufgabe schreiben. Wählen Sie basierend auf Ihrer Geschäftslogik:
- Verwenden Sie eine Notebook-Aufgabe , wenn die Per-Iteration-Logik prozeduralen Code, mehrere Sprachen oder Bibliotheken benötigt (zum Beispiel einen Data-Science- oder Machine-Learning-Schritt).
- Verwenden Sie eine SQL-Aufgabe, wenn die Logik eine einzelne Abfrage oder Transformation ist, die Sie deklarativ ausdrücken können. Eine SQL-Aufgabe benötigt ein SQL-Warehouse.
Beide unten aufgeführten Varianten liefern dasselbe Ergebnis: Für das zu bearbeitende Segment die Anzahl der Angebote auf oder über dem Preisuntergrund und deren Durchschnittspreis.
Notebook-Aufgabe
Erstellen Sie ein neues Notizbuch unter einem Pfad wie /Workspace/Users/<username>/run_segment_analysis. Dieses Notizbuch läuft pro Iteration For each der Aufgabe einmal und erhält jedes Mal ein anderes Segment.
Fügen Sie dem Notizbuch den folgenden Code hinzu:
# 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
Rufen Sie dbutils.widgets.text() vor dbutils.widgets.get() an. Wenn du zuerst anrufst get , führt das Ausführen des Notizbuchs außerhalb eines Jobs zu einem InputWidgetNotDefined Fehler.
SQL-Aufgabe
Eine SQL-Aufgabe führt eine gespeicherte Abfrage aus, also erstellen und speichern Sie die Analyseanfrage jetzt im SQL-Editor. Du hängst sie an die verschachtelte Aufgabe an, wenn du die For each Aufgabe in Schritt 4 konfigurierst.
Klicken Sie in Ihrem Azure Databricks-Arbeitsbereich auf
Neu>
Abfrage, um den SQL-Editor zu öffnen.
Geben Sie die folgende Abfrage ein. SQL-Aufgaben beziehen sich auf Parameter mit der Syntax
:param_name, sodass die Abfrage ihr Segment und Preisboden aus den:property_typeund-Parametern:min_priceablest: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;Klicken Sie auf den Titel
New Query <date>in der Registerkarte Ihrer SQL-Datei und geben Sie ihr den Namenrun_segment_analysis. Dann klicke auf Speichern , um es in einen Ordner zu verschieben, in dem du es speichern möchtest.
Die Aufgabe For each übergibt die Werte jeder Iteration zur Laufzeit an die :property_type und :min_price benannten Parameter. Im Gegensatz zu Notebook-Widgets unterstützen SQL-benannte Parameter keine Standardwerte: Wenn ein Parameter nicht übergeben wird, schlägt die Abfrage mit einem Parameter-Resolutionsfehler fehl.
Schritt 3: Erstellen Sie die Suchanfrage
Die Suchaufgabe liest die Kontrolltabelle über eine gespeicherte Abfrage. Wie in Schritt 2 erstelle und speichere die Abfrage jetzt im SQL-Editor und hänge sie dann in Schritt 4 an die Nachschlageaufgabe an.
Klicken Sie in Ihrem Azure Databricks-Arbeitsbereich auf
Neu>
Abfrage, um den SQL-Editor zu öffnen.
Geben Sie Folgendes ein, mit demselben Katalog, den Sie in Schritt 1 ausgewählt haben:
SELECT property_type, min_price FROM <catalog-name>.config.property_segments;Der Name ist vollständig qualifiziert, weil das SQL-Warehouse, das diese Abfrage ausführt, standardmäßig auf einen anderen Katalog zurückgreifen könnte als der, in dem Sie die Tabelle erstellt haben.
Klicken Sie auf den Titel
New Query <date>in der Registerkarte Ihrer SQL-Datei und geben Sie ihr den Namenread_segments. Dann klicke auf Speichern , um es in einen Ordner zu verschieben, in dem du es speichern möchtest.
Schritt 4: Erstellen und konfigurieren Sie den Job
Mit beiden Abfragen gespeichert, erstellen Sie den Job und fügen seine zwei Aufgaben hinzu: die SQL-Abfrageaufgabe, die die Steuertabelle liest, und die For each Aufgabe, die die Analyse für jede Zeile ausführt.
Schaffen Sie den Job
In Ihrem Azure Databricks Arbeitsbereich klicken Sie in der Seitenleiste Neu>
Job. Gib dem Job einen beschreibenden Namen, wie zum Beispiel
Segment Analysis.
Konfigurieren Sie die SQL-Abfrageaufgabe
Diese Aufgabe liest die Kontrolltabelle und stellt deren Zeilen der Aufgabe zur Verfügung For each , indem sie die read_segments in Schritt 3 gespeicherte Abfrage ausführt.
- Klicken Sie auf die SQL-Abfragekachel , um die erste Aufgabe zu konfigurieren. Wenn die SQL-Abfragekachel nicht verfügbar ist, klicken Sie auf Einen weiteren Aufgabentyp hinzufügen und suchen Sie nach SQL-Abfrage.
- Legen Sie den Aufgabennamen auf
read_segments. - Falls nötig, wählen Sie die SQL-Abfrage im Dropdown-Menü " Typ " aus.
- Wählen Sie im SQL-Abfragefeld die
read_segmentsin Schritt 3 gespeicherte Abfrage aus. - Stellen Sie das SQL-Warehouse auf ein Warehouse in Ihrem Arbeitsbereich ein.
- Klicken Sie auf Aufgabe erstellen.
Wenn diese Aufgabe ausgeführt wird, erfasst Azure Databricks das Ergebnis als JSON-Array in tasks.read_segments.output.rows. Die SQL-Aufgabenausgabe wird immer als JSON-Array zurückgegeben, daher benötigen Sie keine zusätzliche Konfiguration. Die allgemeine Form der Referenz ist tasks.<task-name>.output.rows, wobei <task-name> der von dir festgelegte Aufgabenname übereinstimmt. Die Ausgabe sieht wie folgt aus:
[
{ "property_type": "Urban Year-Round", "min_price": 150 },
{ "property_type": "Summer Getaway", "min_price": 200 },
{ "property_type": "Ski Resort", "min_price": 250 }
]
Konfigurieren Sie die For each Aufgabe
Die For each Aufgabe liest die SQL-Ausgabe und startet eine geschachtelte Aufgabe pro Zeile.
Klicken Sie auf
Aufgabe hinzufügen und für jede auswählen.
Legen Sie den Aufgabennamen auf
process_segments.Überprüfen Sie, dass Depends auf gesetzt ist
read_segments.Im Eingabefeld geben Sie das von der SQL-Aufgabe erfasste Zeilenarray ein:
{{tasks.read_segments.output.rows}}Stellen Sie Concurrency so ein
2, dass zwei Iterationen parallel ausgeführt werden. Erhöhen Sie diesen Wert, wenn Ihre geschachtelte Aufgabe höhere Parallelität unterstützt.Um diese Aufgabe abzuschließen, klicken Sie auf Hinzufügen einer Aufgabe zum Loop und konfigurieren Sie die verschachtelte Aufgabe, die bei jeder Iteration ausgeführt wird.
Die For each Aufgabe und ihre verschachtelte Aufgabe werden zusammen als eine einzige Aufgabe erstellt. Konfigurieren Sie die verschachtelte Aufgabe basierend auf dem Typ, den Sie in Schritt 2 gewählt haben:
Notebook-Aufgabe
Legen Sie den Aufgabennamen auf
run_segment_analysis.Legen Sie "Typ " auf " Notizbuch" fest.
Setze den Pfad für das Notizbuch, das du in Schritt 2 erstellt hast.
Klicken Sie auf Parameter und dann auf Hinzufügen , um jeden Parameter hinzuzufügen:
-
Schlüssel: ,
property_type:{{input.property_type}} -
Schlüssel: ,
min_price:{{input.min_price}}
Jede
{{input.<key>}}Referenz löst sich auf das Matching-Feld aus der aktuellen Iterationszeile auf.-
Schlüssel: ,
Klicken Sie auf Task erstellen , um die
For eachAufgabe und ihre verschachtelte Aufgabe gemeinsam zu erstellen.
SQL-Aufgabe
Diese Aufgabe führt die run_segment_analysis in Schritt 2 gespeicherte Abfrage aus.
Legen Sie den Aufgabennamen auf
run_segment_analysis.Setze Type auf SQL, dann setze SQL-Aufgabe auf Abfrage.
Wählen Sie im SQL-Abfragefeld die
run_segment_analysisin Schritt 2 gespeicherte Abfrage aus.Stellen Sie das SQL-Warehouse auf ein Warehouse in Ihrem Arbeitsbereich ein.
Klicken Sie auf Parameter und dann auf Hinzufügen , um jeden Parameter hinzuzufügen:
-
Schlüssel: ,
property_type:{{input.property_type}} -
Schlüssel: ,
min_price:{{input.min_price}}
Jede
{{input.<key>}}Referenz löst sich auf das Matching-Feld aus der aktuellen Iterationszeile auf.-
Schlüssel: ,
Klicken Sie auf Task erstellen , um die
For eachAufgabe und ihre verschachtelte Aufgabe gemeinsam zu erstellen.
Dein Job Directed Acyclic Graph (DAG) zeigt read_segments jetzt einen Fluss in process_segments, mit der verschachtelten Aufgabe innerhalb des Knotens For each .
Schritt 5: Führe den Auftrag durch und verifiziere
- Klicken Sie auf "Jetzt ausführen" , um den Auftrag auszulösen.
- Wähle den Reiter "Runs " aus, um den Lauf zu sehen. Der erste Lauf eines Jobs benötigt einige Minuten, um die Berechnung zu starten; wenn es abgeschlossen ist, erscheint es in der Liste.
- Klicken Sie auf den
process_segmentsKnoten, um dieFor eachAufgabe zu erweitern. - Die Laufseite zeigt eine Tabelle der Iterationen, eine Zeile pro Segment, jede mit Status, Startzeit und Dauer.
- Klicken Sie auf eine Iterationszeile, um deren Ausgabe zu öffnen und zu bestätigen, dass das erwartete Segment analysiert wurde.
Du kannst die Ergebnisse jeder Iteration unabhängig davon sehen. Wenn eine bestimmte Iteration fehlschlägt, kannst du nur diese Iteration von der Job-Laufseite erneut ausführen, ohne den gesamten Job erneut auszuführen.
Erweitern des Musters
Um ein Segment zur Analyse hinzuzufügen, fügen Sie eine Zeile in die Kontrolltabelle ein:
INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);
Der nächste Job-Run beinhaltet das neue Segment, ohne Änderungen an der Job-Konfiguration oder Notebook-Bearbeitungen.
Dieses gleiche Muster gilt für jeden Fall, in dem man Daten als Iteration vorantreiben möchte:
- Kundenspezifische Verarbeitung: Eine Zeile pro Kunden-ID. Die verschachtelte Aufgabe wendet kundenspezifische Transformationen an oder liefert an kundenspezifische Ziele.
- Tabellenimport: Eine Zeile je Name der Quelltabelle. Die verschachtelte Aufgabe liest und liest jede Tabelle ein.
- Backfill-Verarbeitung: Eine Zeile pro Datumspartition. Die verschachtelte Aufgabe verarbeitet historische Daten für diese Partition erneut.
- Feature-Flag-gesteuerte Ausführung: Eine Zeile pro aktiviertem Feature oder Experiment. Die verschachtelte Aufgabe aktiviert die entsprechende Logik.
Um die Verarbeitung einer Zeile zu stoppen, ohne sie zu löschen, füge eine eigene Spalte in die Steuerungstabelle (wie eine active Flagge) hinzu und filtere sie in der SQL-Nachschlage-Aufgabe. Dies ist eine gewöhnliche Spalte, die du definierst und befüllst; Die For each Aufgabe hat kein eingebautes Konzept. Füge zuerst die Spalte hinzu und setze dann die vorhandenen Zeilen auf TRUE:
ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;
Filtern Sie dann in der read_segments Abfrage darauf, sodass nur aktive Zeilen die Iteration steuern:
SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;
Weitere Ressourcen
-
Verwenden einer Aufgabe zum Ausführen einer
For eachanderen Aufgabe in einer Schleife: Vollständige Referenz zum Konfigurieren vonFor eachAufgaben, einschließlich Parametertypen und Parallelitätsoptionen -
Verwenden Sie eine Nachschlagetabelle für große Parameterarrays in einer
For each-Aufgabe: So behandeln Sie große Parameterarrays, die den Aufgabenwertgrenzwert von 48 KB überschreiten - Zugreifen auf Parameterwerte aus einer Aufgabe: Alle Methoden für den Zugriff auf Parameterwerte in Notizbüchern, Python Skripts und SQL-Aufgaben
- Wanderbricks-Datensatz: Der Beispieldatensatz, der in diesem Tutorial verwendet wird