Usar uma tabela de controle para executar um For each trabalho

Quando você executa o mesmo processamento em várias entradas, como mercados, tabelas fonte, clientes ou partições de data, codificar essa lista no seu trabalho significa editar código e reimplantar toda vez que a lista muda. Em vez disso, armazene a lista em uma tabela de controle que o trabalho lê em tempo de execução. Para adicionar ou remover trabalho, você atualiza uma linha na tabela, e a próxima execução do trabalho retoma a mudança sem editar o próprio trabalho. Esse é 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 esse padrão no conjunto de dados de exemplo pré-instalado do Wanderbricks, para que você possa executá-lo de ponta a ponta sem criar nenhum dado de origem. O cenário é uma plataforma de aluguel 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 controle lista os segmentos a serem analisados, 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 conecta três tarefas em sequência:

Tarefa Tipo O que faz
read_segments SQL Lê a tabela de controle e captura as linhas como um array JSON
process_segments Para cada Itera sobre o array de linhas, iniciando a tarefa aninhada uma vez por linha
run_segment_analysis Caderno ou SQL (aninhado dentro For eachde mim) Executa 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 linhas, 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 então passa 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 do 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 em um catálogo (os USE CATALOG privilégios e CREATE SCHEMA de e) para armazenar a tabela de controle.
  • Um warehouse SQL para executar as tarefas SQL. Se você não tem um, veja Criar um warehouse 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, então não há dados fonte para configurar.

Passo 1: Criar a tabela de controle

A tabela de controle é a fonte de verdade para a lista de segmentos que seu trabalho processa. Para mudar o que o trabalho faz, você atualiza essa tabela, não o trabalho.

Execute o seguinte SQL em um notebook do Azure Databricks ou no editor SQL. A primeira instrução cria um esquema para armazenar a tabela de controle, e a segunda cria a tabela com uma linha por segmento de propriedade e o preço mínimo de listagem a ser incluído 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);

Substitua <catalog-name> por um catálogo no qual você possa criar esquemas, como o catálogo do seu espaço de trabalho. Use o mesmo catálogo em todos os lugares que o tutorial menciona config.property_segments, inclusive na consulta de busca no Passo 3.

Após essa etapa, config.property_segments contém três linhas, uma por segmento. Cada linha carrega os dois valores que o trabalho passa para cada iteração: o property_type para analisar e o min_price piso para filtrar.

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

A tarefa aninhada dentro da For each tarefa roda uma vez por linha da tabela de controle, recebendo os parâmetros dessa linha property_type e min_price como as. Você pode escrever essa lógica como uma tarefa de notebook ou SQL. Escolha com base na lógica do seu negócio:

  • Use uma tarefa de notebook quando a lógica de iteração por necessidade de código procedimental, múltiplas linguagens ou bibliotecas (por exemplo, uma etapa de ciência de dados ou aprendizado de máquina).
  • Use uma tarefa SQL quando a lógica for uma única consulta ou transformação que você possa expressar declarativamente. Uma tarefa SQL precisa de um warehouse 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 preço médio deles.

Tarefa do notebook

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

Adicione o seguinte código ao 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

Chamar dbutils.widgets.text() antes dbutils.widgets.get(). Se você ligar get primeiro, rodar o notebook fora de um job gera um InputWidgetNotDefined erro.

Tarefa SQL

Uma tarefa SQL executa uma consulta salva, então crie e salve a consulta de análise no editor SQL agora. Você a anexa à tarefa aninhada quando configura a For each tarefa no Passo 4.

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

  2. Insira a consulta a seguir. As tarefas SQL referenciam parâmetros com a :param_name sintaxe, então a consulta lê seu segmento e 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. Clique no título New Query <date> no cabeçalho da aba do seu arquivo SQL e dê o nome run_segment_analysisa ele. Depois clique em Salvar para mover para uma pasta onde você quer armazená-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. Diferente dos widgets de notebook, parâmetros nomeados SQL não suportam valores padrão: 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 busca

A tarefa de consulta lê a tabela de controle através de uma consulta salva. Como na Etapa 2, crie e salve a consulta no editor SQL agora, depois anexe à tarefa de busca na Etapa 4.

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

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

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

    O nome é totalmente qualificado porque o warehouse SQL que executa essa consulta pode ter por padrão um catálogo diferente daquele em que você criou a tabela.

  3. Clique no título New Query <date> no cabeçalho da aba do seu arquivo SQL e dê o nome read_segmentsa ele. Depois clique em Salvar para mover para uma pasta onde você quer armazená-lo.

Passo 4: Criar e configurar o trabalho

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

Crie o cargo

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

Configurar a tarefa de busca SQL

Essa tarefa lê a tabela de controle e disponibiliza suas linhas para a For each tarefa executando a read_segments consulta que você salvou no Passo 3.

  1. Clique no bloco de consulta SQL para configurar a primeira tarefa. Se o bloco de consulta SQL não estiver disponível, clique em Adicionar outro tipo de tarefa e pesquise por consulta SQL.
  2. Definir o nome da tarefa como 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 você salvou na Etapa 3.
  5. Defina o SQL Warehouse para um armazém em seu workspace.
  6. Clique em Criar tarefa.

Quando essa 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 é retornada como um array JSON, então você não precisa de nenhuma configuração extra. A forma geral da referência é tasks.<task-name>.output.rows, onde <task-name> corresponde ao nome da tarefa que você definiu. A saída tem esta aparência:

[
  { "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 For each tarefa lê a saída do SQL e inicia uma execução de tarefa aninhada por linha.

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

  2. Definir o nome da tarefa como process_segments.

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

  4. No campo Entradas , insira o array de linhas capturado pela tarefa SQL:

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

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

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

Tarefa do notebook

  1. Definir o nome da tarefa como run_segment_analysis.

  2. Defina Tipo como Notebook.

  3. Defina o Caminho para o caderno que você criou na Etapa 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. Clique em Criar tarefa para criar a For each tarefa e sua tarefa aninhada juntos.

Tarefa SQL

Essa tarefa executa a run_segment_analysis consulta que você salvou na Etapa 2.

  1. Definir o nome da tarefa como 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 você salvou na Etapa 2.

  4. Defina o SQL Warehouse para um armazém em seu workspace.

  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. Clique em Criar tarefa para criar a For each tarefa e sua tarefa aninhada juntos.

Seu Grafo Acíclico Dirigido (DAG) agora mostra read_segments 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 disparar o trabalho.
  2. Selecione a aba Corridas para ver a corrida. A primeira execução de um trabalho leva alguns minutos para 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 execução mostra uma tabela de iterações, uma linha por segmento, cada uma com seu status, horário de início e duração.
  5. Clique em qualquer linha de iteração para abrir a saída e confirmar que analisou o segmento esperado.

Você pode ver os resultados de cada iteração de forma independente. Se uma iteração específica falhar, você pode executar apenas essa iteração a partir da página de execução do trabalho sem reexecutar o trabalho inteiro.

Estender o padrão

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

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

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

Esse mesmo padrão funciona em qualquer caso em que você queira que dados conduzam a iteração:

  • Processamento por cliente: uma linha por ID do cliente. A tarefa aninhada aplica transformações específicas do cliente ou entrega para destinos específicos do cliente.
  • Ingestão de tabela: uma linha para cada nome de tabela de origem. A tarefa aninhada lê e ingere cada tabela.
  • Processamento de backfill: uma linha por partição de data. A tarefa aninhada reprocessa dados históricos dessa partição.
  • Execução orientada por feature flags: uma linha por recurso ou experimento habilitado. A tarefa aninhada ativa a lógica correspondente.

Para parar de processar uma linha sem excluí-la, adicione sua própria coluna à tabela de controle (como uma active flag) e filtre nela na tarefa de consulta SQL. Esta é uma coluna comum que você 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, filtrar na read_segments consulta para que apenas as linhas ativas conduzam a iteração:

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

Recursos adicionais