Verwende eine Steuertabelle, um einen For each Auftrag zu steuern

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_segmentsprocess_segmentsrun_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 CATALOG und CREATE SCHEMA Privilegien), 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 von samples.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.

  1. Klicken Sie in Ihrem Azure Databricks-Arbeitsbereich auf das Plus-Symbol.Neu>Abfrage-Symbol.Abfrage, um den SQL-Editor zu öffnen.

  2. 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_type und-Parametern :min_price ablest:

    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;
    
  3. Klicken Sie auf den Titel New Query <date> in der Registerkarte Ihrer SQL-Datei und geben Sie ihr den Namen run_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.

  1. Klicken Sie in Ihrem Azure Databricks-Arbeitsbereich auf das Plus-Symbol.Neu>Abfrage-Symbol.Abfrage, um den SQL-Editor zu öffnen.

  2. 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.

  3. Klicken Sie auf den Titel New Query <date> in der Registerkarte Ihrer SQL-Datei und geben Sie ihr den Namen read_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 auf das Plus-Symbol.Neu>Workflows-Symbol.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.

  1. 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.
  2. Legen Sie den Aufgabennamen auf read_segments.
  3. Falls nötig, wählen Sie die SQL-Abfrage im Dropdown-Menü " Typ " aus.
  4. Wählen Sie im SQL-Abfragefeld die read_segments in Schritt 3 gespeicherte Abfrage aus.
  5. Stellen Sie das SQL-Warehouse auf ein Warehouse in Ihrem Arbeitsbereich ein.
  6. 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.

  1. Klicken Sie auf das Plus-Symbol.Aufgabe hinzufügen und für jede auswählen.

  2. Legen Sie den Aufgabennamen auf process_segments.

  3. Überprüfen Sie, dass Depends auf gesetzt ist read_segments.

  4. Im Eingabefeld geben Sie das von der SQL-Aufgabe erfasste Zeilenarray ein:

    {{tasks.read_segments.output.rows}}
    
  5. 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.

  6. 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

  1. Legen Sie den Aufgabennamen auf run_segment_analysis.

  2. Legen Sie "Typ " auf " Notizbuch" fest.

  3. Setze den Pfad für das Notizbuch, das du in Schritt 2 erstellt hast.

  4. 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.

  5. Klicken Sie auf Task erstellen , um die For each Aufgabe und ihre verschachtelte Aufgabe gemeinsam zu erstellen.

SQL-Aufgabe

Diese Aufgabe führt die run_segment_analysis in Schritt 2 gespeicherte Abfrage aus.

  1. Legen Sie den Aufgabennamen auf run_segment_analysis.

  2. Setze Type auf SQL, dann setze SQL-Aufgabe auf Abfrage.

  3. Wählen Sie im SQL-Abfragefeld die run_segment_analysis in Schritt 2 gespeicherte Abfrage aus.

  4. Stellen Sie das SQL-Warehouse auf ein Warehouse in Ihrem Arbeitsbereich ein.

  5. 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.

  6. Klicken Sie auf Task erstellen , um die For each Aufgabe 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

  1. Klicken Sie auf "Jetzt ausführen" , um den Auftrag auszulösen.
  2. 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.
  3. Klicken Sie auf den process_segments Knoten, um die For each Aufgabe zu erweitern.
  4. Die Laufseite zeigt eine Tabelle der Iterationen, eine Zeile pro Segment, jede mit Status, Startzeit und Dauer.
  5. 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