Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Use uma mesa de controlo para gerir um
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_segments → process_segments → run_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 CATALOGprivilégios e)CREATE SCHEMApara 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
samplescatálogo, que está disponível em todos os espaços de trabalho habilitados pelo Unity Catalog. O tutorial lê desamples.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.
No seu espaço de trabalho Azure Databricks, clique no
Novo>
Consulta para abrir o editor SQL.
Insira a seguinte consulta. As tarefas SQL referenciam parâmetros com a
:param_namesintaxe, pelo que a consulta lê o seu segmento e o 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;Clica no título
New Query <date>no título do separador do teu ficheiro SQL e dá-lhe o nomerun_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.
No seu espaço de trabalho Azure Databricks, clique no
Novo>
Consulta para abrir o editor SQL.
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.
Clica no título
New Query <date>no título do separador do teu ficheiro SQL e dá-lhe o nomeread_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 Novo>
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.
- 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.
- Defina o nome da Tarefa para
read_segments. - Se necessário, selecione consulta SQL no menu suspenso de Tipo .
- No campo de consulta SQL , selecione a
read_segmentsconsulta que guardou no Passo 3. - Definir SQL warehouse como um armazém no teu espaço de trabalho.
- 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.
Adicionar tarefa e selecionar Para cada um.
Defina o nome da Tarefa para
process_segments.Verifique que Depende de está definido para
read_segments.No campo Inputs , introduza o array de linhas capturado pela tarefa SQL:
{{tasks.read_segments.output.rows}}Defina a Concorrência para
2correr duas iterações em paralelo. Aumente este valor quando a sua tarefa aninhada suportar maior paralelismo.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
Defina o nome da Tarefa para
run_segment_analysis.Definir Tipo como Portátil.
Define o Caminho para o caderno que criaste no Passo 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:
Clica em Criar tarefa para criar a
For eachtarefa e a sua tarefa aninhada em conjunto.
Tarefa SQL
Esta tarefa executa a run_segment_analysis consulta que guardaste no Passo 2.
Defina o nome da Tarefa para
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 guardou no Passo 2.Definir SQL warehouse como um armazém no teu espaço de trabalho.
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:
Clica em Criar tarefa para criar a
For eachtarefa 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
- Clique em Executar agora para ativar a tarefa.
- 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.
- Clique no
process_segmentsnó para expandir aFor eachtarefa. - 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.
- 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
-
Usar uma
For eachtarefa para executar outra tarefa num ciclo: ReferênciaFor eachcompleta para configurar tarefas, incluindo tipos de parâmetros e opções de concorrência -
Use uma tabela de consulta para grandes arrays de parâmetros numa
For eachtarefa: Como lidar com grandes arrays de parâmetros que excedam o limite de 48 KB de valores de tarefa - Aceder aos valores dos parâmetros a partir de uma tarefa: Todos os métodos para aceder aos valores dos parâmetros em notebooks, scripts em Python e tarefas SQL
- Conjunto de dados Wanderbricks: O conjunto de dados de exemplo utilizado neste tutorial