Użyj tabeli sterującej do sterowania zadaniem For each

Gdy uruchamiasz to samo przetwarzanie na wielu wejściach, takich jak rynki, tabele źródłowe, klienci czy partycje dat, twarde kodowanie tej listy w zadaniu oznacza edytowanie kodu i ponowne wdrażanie za każdym razem, gdy lista się zmienia. Zamiast tego przechowuj listę w tabeli kontrolnej , którą zadanie odczytuje w czasie działania. Aby dodać lub usunąć pracę, aktualizujesz wiersz w tabeli, a następne uruchomienie zadania odbiera zmianę bez zmian w samym zadaniu. To wzorzec oparty na metadanych : dane, a nie kod, kontrolują to, co praca przetwarza.

Ten samouczek tworzy zadanie, które wykorzystuje ten wzorzec na preinstalowanym zbiorze danych Wanderbricks, dzięki czemu możesz uruchomić go end to end to end bez tworzenia żadnych danych źródłowych. Scenariusz polega na platformie wynajmu wakacyjnego, która przeprowadza tę samą analizę cen dla każdego segmentu nieruchomości (np. Ski Resort lub Urban Year-Round). Tabela kontrolna wymienia segmenty do analizy, zadanie SQL odczytuje tę tabelę, a zadanie wykonuje For each analizę raz na segment, równolegle.

Jak to działa

Zadanie łączy trzy zadania w kolejności:

Zadanie Typ Do czego służy
read_segments SQL Odczytuje tablicę sterującą i rejestruje wiersze jako tablicę JSON
process_segments Dla każdego Iteruje przez tablicę wierszy, uruchamiając zagnieżdżone zadanie raz na wiersz
run_segment_analysis Notebook lub SQL (zagnieżdżony wewnątrz For each) Uruchamia się raz na wiersz, wykorzystując wartości tego wiersza do analizy jednego segmentu właściwości

Przepływ wynosi read_segments → process_segments → run_segment_analysis (raz na rząd). Wyjście zadania SQL, czyli tablica JSON obiektów wierszowych, trafia do pola For each zadania przez odniesienie {{tasks.read_segments.output.rows}}wartości dynamicznej . Następnie zadanie przekazuje For each pola każdego wiersza do zagnieżdżonego zadania jako parametry, dostępne jako {{input.property_type}} i {{input.min_price}}.

Wymagania wstępne

  • Przestrzeń robocza Azure Databricks z uprawnieniami do tworzenia zadań i notatników.
  • Uprawnienia do tworzenia tabel w Unity Catalog oraz uprawnień do tworzenia schematu w katalogu (uprawnienia USE CATALOG i CREATE SCHEMA do przechowywania tabeli kontrolnej.
  • Magazyn SQL do uruchamiania zadań SQL. Jeśli nie masz takiego magazynu, zobacz: Utwórz magazyn SQL.
  • Katalog, samples który jest dostępny we wszystkich przestrzeniach roboczych obsługujących Unity Catalog. Samouczek czyta od samples.wanderbricks.properties, więc nie ma danych źródłowych do ustawienia.

Krok 1: Stwórz tablicę kontrolną

Tabela kontrolna jest źródłem prawdy dla listy segmentów, które przetwarza twoja praca. Aby zmienić działanie zadania, aktualizujesz tę tabelę, a nie zadanie.

Uruchom poniższy SQL w notatniku Azure Databricks lub edytorze SQL. Pierwsze sformułowanie tworzy schemat do przechowywania tabeli kontrolnej, a drugie tworzy tabelę z jednym wierszem na segment nieruchomości oraz minimalną ceną wystawienia, którą należy uwzględnić w analizie tego segmentu:

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);

Zamień <catalog-name> go na katalog, w którym możesz tworzyć schematy, na przykład katalog w Workspace. Używaj tego samego katalogu wszędzie, gdzie odwołuje config.property_segmentssię samouczek, w tym zapytania wyszukiwania w kroku 3.

Po tym kroku zawiera config.property_segments trzy rzędy, po jednym na segment. Każdy wiersz zawiera dwie wartości, które zadanie przekazuje każdej iteracji: do property_type analizy oraz podłogę min_price do filtrowania.

Krok 2: Napisz logikę analizy

Zagnieżdżone zadanie wewnątrz zadania For each wykonuje się raz na wiersz tabeli kontrolnej, otrzymując parametry tego wiersza property_type i min_price jako parametry. Możesz zapisać tę logikę jako zadanie notatnika lub zadanie SQL. Wybierz na podstawie logiki swojej działalności:

  • Użyj zadania notatnika, gdy logika pojedynczej iteracji wymaga kodu proceduralnego, wielu języków lub bibliotek (na przykład etap data science lub uczenia maszynowego).
  • Użyj zadania SQL , gdy logika to pojedyncze zapytanie lub transformacja, którą możesz wyrazić deklaratywnie. Zadanie SQL wymaga magazynu SQL.

Oba poniższe warianty dają ten sam wynik: dla przetwarzanego segmentu liczbę ofert na poziomie lub powyżej jego ceny minimalnej oraz ich średnią cenę.

Zadanie notatnika

Utwórz nowy notatnik w ścieżce, takiej jak /Workspace/Users/<username>/run_segment_analysis. Ten notes wykonuje się raz na iterację zadania For each , otrzymując za każdym razem inny segment.

Dodaj następujący kod do notesu:

# 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

Wywołaj dbutils.widgets.text() przed dbutils.widgets.get(). Jeśli zadzwonisz get pierwszy, uruchomienie notesu poza zadaniem powoduje InputWidgetNotDefined błąd.

Zadanie SQL

Zadanie SQL uruchamia zapisane zapytanie, więc teraz stwórz i zapisz zapytanie analityczne w edytorze SQL. Dołączasz je do zagnieżdżonego zadania podczas konfiguracji For each zadania w kroku 4.

  1. W swojej przestrzeni roboczej Azure Databricks kliknij ikonę Plus.Nowe>Ikona zapytania.Zapytanie, aby otworzyć edytor SQL.

  2. Wprowadź następujące zapytanie. Zadania SQL odnoszą się do parametrów w :param_name składni, więc zapytanie odczytuje segment i minimalną cenę z parametrów :property_type i :min_price :

    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. Kliknij tytuł New Query <date> w zakładce pliku SQL i nadaj mu nazwę run_segment_analysis. Następnie kliknij Zapisz, aby przenieść go do folderu, w którym chcesz go przechować.

Zadanie przekazuje For each wartości każdej iteracji do parametrów :property_type o nazwie w :min_price czasie działania. W przeciwieństwie do widgetów notatników, parametry nazwane SQL nie obsługują wartości domyślnych: jeśli parametr nie zostanie przekazany, zapytanie kończy się niepowodzeniem z błędem rozwiązania parametrów.

Krok 3: Utwórz zapytanie wyszukujące

Zadanie wyszukiwania odczytuje tablicę kontrolną przez zapisane zapytanie. Podobnie jak w kroku 2, utworz i zapisz zapytanie w edytorze SQL, a następnie dołącz je do zadania wyszukiwania w kroku 4.

  1. W swojej przestrzeni roboczej Azure Databricks kliknij ikonę Plus.Nowe>Ikona zapytania.Zapytanie, aby otworzyć edytor SQL.

  2. Wpisz następujące wpisy, używając tego samego katalogu, który wybrałeś w kroku 1:

    SELECT property_type, min_price FROM <catalog-name>.config.property_segments;
    

    Nazwa jest w pełni kwalifikowana, ponieważ magazyn SQL, który uruchamia to zapytanie, może domyślnie używać innego katalogu niż ten, w którym utworzyłeś tabelę.

  3. Kliknij tytuł New Query <date> w zakładce pliku SQL i nadaj mu nazwę read_segments. Następnie kliknij Zapisz, aby przenieść go do folderu, w którym chcesz go przechować.

Krok 4: Stwórz i skonfiguruj zadanie

Po zapisaniu obu zapytań utwórz zadanie i dodaj jego dwa zadania: zadanie wyszukiwania SQL, które odczytuje tabelę kontrolną, oraz zadanie For each wykonujące analizę dla każdego wiersza.

Stwórz pracę

W twojej przestrzeni roboczej Azure Databricks, w pasku bocznym kliknij ikonę Plus.Nowe>Ikona przepływów pracy.Praca. Nadaj temu zadaniu opisową nazwę, na przykład Segment Analysis.

Konfiguruj zadanie wyszukiwania SQL

To zadanie odczytuje tabelę kontrolną i udostępnia jej wiersze zadaniu, uruchamiając zapytanie zapisane For eachread_segments w kroku 3.

  1. Kliknij kafelek zapytania SQL , aby skonfigurować pierwsze zadanie. Jeśli kafelek zapytania SQL nie jest dostępny, kliknij Dodaj inny typ zadania i wyszukaj zapytanie SQL.
  2. Ustaw Nazwa zadania na read_segments.
  3. W razie potrzeby wybierz zapytanie SQL z rozwijanego menu Typ .
  4. W polu zapytania SQL wybierz zapytanie, które read_segments zapisałeś w kroku 3.
  5. Ustaw usługę SQL Warehouse na magazyn w obszarze roboczym.
  6. Kliknij pozycję Utwórz zadanie.

Gdy to zadanie się uruchamia, Azure Databricks rejestruje wynik jako tablicę JSON w .tasks.read_segments.output.rows Wyjście zadań SQL zawsze zwracane jest jako tablica JSON, więc nie potrzebujesz dodatkowej konfiguracji. Ogólna forma odniesienia to tasks.<task-name>.output.rows, gdzie <task-name> odpowiada nazwie zadania, które ustaliłeś. Dane wyjściowe wyglądają następująco:

[
  { "property_type": "Urban Year-Round", "min_price": 150 },
  { "property_type": "Summer Getaway", "min_price": 200 },
  { "property_type": "Ski Resort", "min_price": 250 }
]

Konfiguruj zadanie For each

Zadanie For each odczytuje dane wyjściowe SQL i uruchamia zagnieżdżone zadanie dla każdego wiersza.

  1. Kliknij ikonę Plus.Dodaj zadanie i wybierz Dla każdego z nich.

  2. Ustaw Nazwa zadania na process_segments.

  3. Sprawdź, czy Depends on jest ustawione na .read_segments

  4. W polu Inputs wpisz tablicę wierszową przechwytaną przez zadanie SQL:

    {{tasks.read_segments.output.rows}}
    
  5. Ustaw współbieżność na tak, 2 aby uruchamiać dwie iteracje równolegle. Zwiększ tę wartość, gdy zagnieżdżone zadanie obsługuje większą równoległość.

  6. Aby wykonać to zadanie, kliknij Dodaj zadanie, aby przeprowadzić pętlę i skonfigurować zagnieżdżone zadanie uruchamiane w każdej iteracji.

Zadanie For each i jego zagnieżdżone zadanie są tworzone razem jako jedno zadanie. Skonfiguruj zagnieżdżone zadanie na podstawie wybranego typu w kroku 2:

Zadanie notatnika

  1. Ustaw Nazwa zadania na run_segment_analysis.

  2. Ustaw Typ na Laptop.

  3. Ustaw ścieżkę do notatnika, który utworzyłeś w kroku 2.

  4. Kliknij Parametry, a następnie Dodaj do każdego parametru:

    • Klucz: property_type, Wartość: {{input.property_type}}
    • Klucz: min_price, Wartość: {{input.min_price}}

    Każde {{input.<key>}} odniesienie rozwiązuje się do pola dopasowania z wiersza bieżącej iteracji.

  5. Kliknij Create task , aby utworzyć For each zadanie i jego zagnieżdżone zadanie razem.

Zadanie SQL

To zadanie wykonuje zapytanie zapisane run_segment_analysis w kroku 2.

  1. Ustaw Nazwa zadania na run_segment_analysis.

  2. Ustaw Type na SQL, a zadanie SQL na Query.

  3. W polu zapytania SQL wybierz zapytanie, które run_segment_analysis zapisałeś w kroku 2.

  4. Ustaw usługę SQL Warehouse na magazyn w obszarze roboczym.

  5. Kliknij Parametry, a następnie Dodaj do każdego parametru:

    • Klucz: property_type, Wartość: {{input.property_type}}
    • Klucz: min_price, Wartość: {{input.min_price}}

    Każde {{input.<key>}} odniesienie rozwiązuje się do pola dopasowania z wiersza bieżącej iteracji.

  6. Kliknij Create task , aby utworzyć For each zadanie i jego zagnieżdżone zadanie razem.

Twój zadaniowy Directed Acyclic Graph (DAG) pokazuje read_segments teraz przepływ do process_segments, z zagnieżdżonym zadaniem wewnątrz For each węzła.

Krok 5: Uruchom zadanie i zweryfikowaj

  1. Kliknij Uruchom teraz, aby uruchomić zadanie.
  2. Wybierz zakładkę Biegi , aby zobaczyć bieg. Pierwszy uruchomienie zadania zajmuje kilka minut, aby rozpocząć obliczenia; Po ukończeniu pojawia się na liście.
  3. Kliknij węzeł, process_segments aby rozwinąć zadanie For each .
  4. Strona serii pokazuje tabelę iteracji, po jednym wierszu na segment, każda ze swoim statusem, czasem rozpoczęcia i czasem trwania.
  5. Kliknij dowolny wiersz iteracji, aby otworzyć jego wyjście i potwierdzić, że przeanalizował oczekiwany segment.

Wyniki każdej iteracji można zobaczyć niezależnie. Jeśli konkretna iteracja się nie powiedzie, możesz uruchomić tylko tę iterację ze strony wykonania zadania bez ponownego uruchamiania całego zadania.

Rozszerzanie wzorca

Aby dodać segment do analizy, wstaw wiersz do tabeli kontrolnej:

INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);

Kolejne zadanie obejmuje nowy segment, bez zmian konfiguracji zadań czy edycji notatnika.

Ten sam schemat działa w każdej sytuacji, gdy chcesz wywołać iterację danymi:

  • Przetwarzanie dla każdego klienta: Jeden wiersz dla każdego identyfikatora klienta. Zagnieżdżone zadanie stosuje transformacje specyficzne dla klienta lub dostarcza do specyficznych dla klienta miejsc.
  • Ładowanie tabel: Jeden wiersz dla każdej nazwy tabeli źródłowej. Zagnieżdżone zadanie odczytuje i pobiera każdą tabelę.
  • Przetwarzanie uzupełniające: jeden wiersz na partycję daty. Zagnieżdżone zadanie przetwarza ponownie dane historyczne dla tej partycji.
  • Wykonywanie oparte na flagach funkcji: jeden wiersz na włączoną funkcję lub eksperyment. Zagnieżdżone zadanie aktywuje odpowiadającą mu logikę.

Aby zatrzymać przetwarzanie wiersza bez jego usunięcia, dodaj własną kolumnę do tabeli kontrolnej (np. flagę active ) i filtruj ją w zadaniu wyszukiwania SQL. To zwykła kolumna, którą definiujesz i wypełniasz; Zadanie nie For each ma wbudowanej koncepcji. Najpierw dodaj kolumnę, a następnie ustaw istniejące wiersze na :TRUE

ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;

Następnie filtruj read_segments to w zapytaniu tak, aby tylko aktywne wiersze napędzały iterację:

SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;

Dodatkowe zasoby