Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Usar uma tabela de controle para executar um
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_segments → process_segments → run_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 CATALOGprivilégios eCREATE SCHEMAde 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
samplescatálogo, que está disponível em todos os espaços de trabalho habilitados pelo Unity Catalog. O tutorial lê desamples.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.
No seu espaço de trabalho do Azure Databricks, clique
Novo>
Consulte para abrir o editor SQL.
Insira a consulta a seguir. As tarefas SQL referenciam parâmetros com a
:param_namesintaxe, então a consulta lê seu segmento e preço mínimo a partir dos:property_typeparâmetros e::min_priceSELECT :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;Clique no título
New Query <date>no cabeçalho da aba do seu arquivo SQL e dê o nomerun_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.
No seu espaço de trabalho do Azure Databricks, clique
Novo>
Consulte para abrir o editor SQL.
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.
Clique no título
New Query <date>no cabeçalho da aba do seu arquivo SQL e dê o nomeread_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 Novo>
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.
- 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.
- Definir o nome da tarefa como
read_segments. - Se necessário, selecione consulta SQL no menu suspenso de Tipo .
- No campo de consulta SQL , selecione a
read_segmentsconsulta que você salvou na Etapa 3. - Defina o SQL Warehouse para um armazém em seu workspace.
- 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.
Adicionar tarefa e selecionar Para cada um.
Definir o nome da tarefa como
process_segments.Verifique que Depende de está definido como
read_segments.No campo Entradas , insira o array de linhas capturado pela tarefa SQL:
{{tasks.read_segments.output.rows}}Defina a Concorrência para
2rodar duas iterações em paralelo. Aumente este valor quando sua tarefa aninhada suportar maior paralelismo.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
Definir o nome da tarefa como
run_segment_analysis.Defina Tipo como Notebook.
Defina o Caminho para o caderno que você criou na Etapa 2.
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.-
Chave:
Clique em Criar tarefa para criar a
For eachtarefa e sua tarefa aninhada juntos.
Tarefa SQL
Essa tarefa executa a run_segment_analysis consulta que você salvou na Etapa 2.
Definir o nome da tarefa como
run_segment_analysis.Defina o Tipo para SQL, depois defina a tarefa SQL para Consulta.
No campo de consulta SQL , selecione a
run_segment_analysisconsulta que você salvou na Etapa 2.Defina o SQL Warehouse para um armazém em seu workspace.
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.-
Chave:
Clique em Criar tarefa para criar a
For eachtarefa 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
- Clique em Executar agora para disparar o trabalho.
- 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.
- Clique no
process_segmentsnó para expandir aFor eachtarefa. - 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.
- 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
-
Use uma tarefa
For eachpara executar outra tarefa em um loop: referência completa para configurar tarefasFor each, incluindo tipos de parâmetros e opções de concorrência -
Use uma tabela de pesquisa para grandes matrizes de parâmetros em uma
For eachtarefa: como lidar com grandes matrizes de parâmetros que excedem o limite de valor da tarefa de 48 KB - Acessar valores de parâmetro de uma tarefa: todos os métodos para acessar valores de parâmetro em notebooks, scripts Python e tarefas SQL
- Conjunto de dados Wanderbricks: O conjunto de dados de exemplo usado neste tutorial