Use uma mesa de controlo para gerir um For each trabalho

Quando executa o mesmo processamento em várias entradas, como mercados, tabelas de origem, clientes ou partições de data, codificar essa lista no seu trabalho significa editar código e redistribuir sempre que a lista muda. Em vez disso, armazene a lista numa tabela de controlo que o trabalho lê em tempo de execução. Para adicionar ou remover trabalho, atualiza-se uma linha na tabela, e a próxima execução do trabalho retoma a alteração sem alterações no próprio trabalho. Este é um padrão orientado por metadados : os dados, e não o código, controlam o que o trabalho processa.

Este tutorial constrói um trabalho que usa este padrão no conjunto de dados de exemplo pré-instalado do Wanderbricks, para que possas executá-lo de ponta a ponta sem criar qualquer fonte de dados. O cenário é uma plataforma de arrendamento de férias que executa a mesma análise de preços para cada segmento de propriedade (como Ski Resort ou Urban Year-Round). Uma tabela de controlo lista os segmentos a analisar, uma tarefa SQL lê essa tabela e uma For each tarefa executa a análise uma vez por segmento, em paralelo.

Como funciona

O trabalho liga três tarefas em sequência:

Tarefa Tipo O que faz
read_segments SQL Lê a tabela de controlo e captura as linhas como um array JSON
process_segments Para cada Itera sobre o array de linhas, lançando a tarefa aninhada uma vez por linha
run_segment_analysis Notebook ou SQL (aninhado no interior For each) Executa-se uma vez por linha, usando os valores dessa linha para analisar um segmento de propriedade

O fluxo é read_segmentsprocess_segmentsrun_segment_analysis (uma vez por linha). A saída da tarefa SQL, um array JSON de objetos linha, flui para o For each campo Inputs da tarefa através da referência {{tasks.read_segments.output.rows}}dinâmica de valor . A For each tarefa passa então os campos de cada linha para a tarefa aninhada como parâmetros, disponíveis como {{input.property_type}} e {{input.min_price}}.

Pré-requisitos

  • Um espaço de trabalho Azure Databricks com permissão para criar trabalhos e cadernos.
  • Permissão para criar tabelas no Catálogo Unity, e permissão para criar um esquema num catálogo (os USE CATALOG privilégios e) CREATE SCHEMA para armazenar a tabela de controlo.
  • Um armazém SQL para executar as tarefas SQL. Se não tiveres um, vê Criar um armazém SQL.
  • O samples catálogo, que está disponível em todos os espaços de trabalho habilitados pelo Unity Catalog. O tutorial lê de samples.wanderbricks.properties, por isso não há dados de origem para configurar.

Passo 1: Criar a tabela de controlo

A tabela de controlo é a fonte de verdade para a lista de segmentos que o seu trabalho processa. Para mudar o que o trabalho faz, atualiza-se esta tabela, não o trabalho.

Execute o seguinte SQL num notebook Azure Databricks ou no editor SQL. A primeira instrução cria um esquema para conter a tabela de controlo, e a segunda cria a tabela com uma linha por segmento de propriedade e o preço mínimo de listagem a incluir na análise desse 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);

Substitui <catalog-name> por um catálogo onde possas criar esquemas, como o catálogo do teu espaço de trabalho. Use o mesmo catálogo em todo o local onde o tutorial se refere config.property_segments, incluindo a consulta de pesquisa no Passo 3.

Após este passo, config.property_segments contém três filas, uma por segmento. Cada linha transporta os dois valores que a tarefa passa a cada iteração: o property_type para analisar e o min_price chão para filtrar.

Passo 2: Escrever a lógica de análise

A tarefa aninhada dentro da For each tarefa executa-se uma vez por linha da tabela de controlo, recebendo os parâmetros dessa property_type linha e min_price as as. Podes escrever esta lógica como uma tarefa de caderno ou SQL. Escolha com base na lógica do seu negócio:

  • Use uma tarefa de caderno quando a lógica por iteração necessitar de código procedimental, várias linguagens ou bibliotecas (por exemplo, uma etapa de ciência de dados ou aprendizagem automática).
  • Use uma tarefa SQL quando a lógica for uma única consulta ou transformação que possa expressar declarativamente. Uma tarefa SQL precisa de um armazém SQL.

Ambas as variantes abaixo produzem o mesmo resultado: para o segmento em processamento, o número de anúncios iguais ou acima do seu preço mínimo e o seu preço médio.

Tarefa do caderno

Crie um novo caderno num caminho como /Workspace/Users/<username>/run_segment_analysis. Este caderno executa-se uma vez por iteração da For each tarefa, recebendo um segmento diferente a cada vez.

Adicione o seguinte código ao caderno:

# 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

Ligue dbutils.widgets.text() antes de dbutils.widgets.get(). Se ligar get primeiro, executar o notebook fora de uma tarefa gera um InputWidgetNotDefined erro.

Tarefa SQL

Uma tarefa SQL executa uma consulta guardada, por isso cria e guarda a consulta de análise no editor SQL agora. Associa-se à tarefa aninhada quando configura a For each tarefa no Passo 4.

  1. No seu espaço de trabalho Azure Databricks, clique no ícone Plus.Novo>Ícone de consulta.Consulta para abrir o editor SQL.

  2. Insira a seguinte consulta. As tarefas SQL referenciam parâmetros com a :param_name sintaxe, pelo que a consulta lê o seu segmento e o preço mínimo a partir dos :property_type parâmetros e::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. Clica no título New Query <date> no título do separador do teu ficheiro SQL e dá-lhe o nome run_segment_analysis. Depois clica em Guardar para o mover para uma pasta onde queres guardá-lo.

A For each tarefa passa os valores de cada iteração para os :property_type parâmetros e :min_price nomeados em tempo de execução. Ao contrário dos widgets de caderno, os parâmetros nomeados SQL não suportam valores predefinidos: se um parâmetro não for passado, a consulta falha com um erro de resolução de parâmetro.

Passo 3: Criar a consulta de pesquisa

A tarefa de consulta lê a tabela de controlo através de uma consulta guardada. Tal como no Passo 2, crie e guarde a consulta no editor SQL agora, depois anexe-a à tarefa de pesquisa no Passo 4.

  1. No seu espaço de trabalho Azure Databricks, clique no ícone Plus.Novo>Ícone de consulta.Consulta para abrir o editor SQL.

  2. Introduza o seguinte, usando o mesmo catálogo que escolheu no Passo 1:

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

    O nome é totalmente qualificado porque o SQL warehouse que executa esta consulta pode indicar por defeito um catálogo diferente daquele onde criou a tabela.

  3. Clica no título New Query <date> no título do separador do teu ficheiro SQL e dá-lhe o nome read_segments. Depois clica em Guardar para o mover para uma pasta onde queres guardá-lo.

Passo 4: Criar e configurar o trabalho

Com ambas as consultas guardadas, crie o trabalho e adicione as suas duas tarefas: a tarefa de consulta SQL que lê a tabela de controlo e a For each tarefa que executa a análise para cada linha.

Criar o trabalho

No seu espaço de trabalho Azure Databricks, na barra lateral clique no ícone Plus.Novo>Ícone de fluxos de trabalho.Trabalho. Dê ao trabalho um nome descritivo, como Segment Analysis.

Configurar a tarefa de pesquisa SQL

Esta tarefa lê a tabela de controlo e disponibiliza as suas linhas para a For each tarefa ao executar a read_segments consulta que guardou no Passo 3.

  1. Clique no bloco de consulta SQL para configurar a primeira tarefa. Se o tile de consulta SQL não estiver disponível, clique em Adicionar outro tipo de tarefa e procure por consulta SQL.
  2. Defina o nome da Tarefa para read_segments.
  3. Se necessário, selecione consulta SQL no menu suspenso de Tipo .
  4. No campo de consulta SQL , selecione a read_segments consulta que guardou no Passo 3.
  5. Definir SQL warehouse como um armazém no teu espaço de trabalho.
  6. Clique em Criar tarefa.

Quando esta tarefa é executada, o Azure Databricks captura o resultado como um array JSON em tasks.read_segments.output.rows. A saída da tarefa SQL é sempre devolvida como um array JSON, por isso não precisa de nenhuma configuração adicional. A forma geral da referência é tasks.<task-name>.output.rows, onde <task-name> corresponde ao nome da tarefa que definiste. O resultado tem o seguinte aspeto:

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

Configurar a For each tarefa

A tarefa For each lê a saída SQL e inicia uma execução de tarefa encadeada por linha.

  1. Clique no ícone Mais. Adicionar tarefa e selecionar Para cada um.

  2. Defina o nome da Tarefa para process_segments.

  3. Verifique que Depende de está definido para read_segments.

  4. No campo Inputs , introduza o array de linhas capturado pela tarefa SQL:

    {{tasks.read_segments.output.rows}}
    
  5. Defina a Concorrência para 2 correr duas iterações em paralelo. Aumente este valor quando a sua tarefa aninhada suportar maior paralelismo.

  6. Para completar esta tarefa, clique em Adicionar uma tarefa para fazer loop e configure a tarefa aninhada que corre em cada iteração.

A For each tarefa e a sua tarefa aninhada são criadas em conjunto como uma única tarefa. Configure a tarefa aninhada com base no tipo que escolheu no Passo 2:

Tarefa do caderno

  1. Defina o nome da Tarefa para run_segment_analysis.

  2. Definir Tipo como Portátil.

  3. Define o Caminho para o caderno que criaste no Passo 2.

  4. Clique em Parâmetros, depois clique em Adicionar para adicionar cada parâmetro:

    • Chave: property_type, Valor: {{input.property_type}}
    • Chave: min_price, Valor: {{input.min_price}}

    Cada {{input.<key>}} referência resolve para o campo correspondente da linha da iteração atual.

  5. Clica em Criar tarefa para criar a For each tarefa e a sua tarefa aninhada em conjunto.

Tarefa SQL

Esta tarefa executa a run_segment_analysis consulta que guardaste no Passo 2.

  1. Defina o nome da Tarefa para run_segment_analysis.

  2. Defina o Tipo para SQL, depois defina a tarefa SQL para Consulta.

  3. No campo de consulta SQL , selecione a run_segment_analysis consulta que guardou no Passo 2.

  4. Definir SQL warehouse como um armazém no teu espaço de trabalho.

  5. Clique em Parâmetros, depois clique em Adicionar para adicionar cada parâmetro:

    • Chave: property_type, Valor: {{input.property_type}}
    • Chave: min_price, Valor: {{input.min_price}}

    Cada {{input.<key>}} referência resolve para o campo correspondente da linha da iteração atual.

  6. Clica em Criar tarefa para criar a For each tarefa e a sua tarefa aninhada em conjunto.

O seu Grafo Acíclico Dirigido (DAG) de trabalho mostra read_segments agora o fluxo em process_segments, com a tarefa aninhada dentro do For each nó.

Passo 5: Execute o trabalho e verifique

  1. Clique em Executar agora para ativar a tarefa.
  2. Selecione o separador Runs para ver a corrida. A primeira execução de um trabalho demora alguns minutos a iniciar o cálculo; Quando termina, aparece na lista.
  3. Clique no process_segments nó para expandir a For each tarefa.
  4. A página da corrida mostra uma tabela de iterações, uma linha por segmento, cada uma com o seu estado, hora de início e duração.
  5. Clique em qualquer linha de iteração para abrir a sua saída e confirmar que analisou o segmento esperado.

Pode ver os resultados de cada iteração de forma independente. Se uma iteração específica falhar, só pode reexecutar essa iteração a partir da página de execução do trabalho sem repetir toda a tarefa.

Estende o padrão

Para adicionar um segmento à análise, insira uma linha na tabela de controlo:

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

A execução seguinte inclui o novo segmento, sem alterações na configuração do trabalho ou edições no caderno.

Este mesmo padrão funciona em qualquer caso em que se queira que os dados conduzam a iteração:

  • Processamento por cliente: Uma fila por ID de cliente. A tarefa aninhada aplica transformações específicas do cliente ou entrega a destinos específicos do cliente.
  • Ingestão de tabelas: Uma linha por nome da tabela de origem. A tarefa aninhada lê e ingere cada tabela.
  • Processamento de retropreenchimento: Uma linha por cada partição de data. A tarefa aninhada reprocessa dados históricos dessa partição.
  • Execução baseada em sinalizadores de funcionalidade: Uma linha por funcionalidade ou experimento ativado. A tarefa aninhada ativa a lógica correspondente.

Para parar de processar uma linha sem a apagar, adiciona a tua própria coluna à tabela de controlo (como uma active flag) e filtra isso na tarefa de consulta SQL. Esta é uma coluna comum que se define e preenche; A For each tarefa não tem um conceito incorporado. Primeiro adicione a coluna, depois defina as linhas existentes para TRUE:

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

Depois filtra na read_segments consulta para que só as linhas ativas conduzam a iteração:

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

Recursos adicionais