pipeline de execução em um fluxo de trabalho
Você pode executar um pipeline como parte de um fluxo de trabalho de processamento de dados com LakeFlow Jobs, Apache Airflow ou Azure Data Factory.
Um pipeline resolve automaticamente as dependências entre seus datasets, de modo que lida com a orquestração simples e interna do pipeline por conta própria. Para orquestração para a qual um pipeline não foi construído, como execução condicional, ramificação em resultados de tarefas, novas tentativas ou coordenação de um pipeline com outros tipos de trabalho, use um orquestrador de fluxo de trabalho dedicado em vez de integrar a lógica no pipeline.
Prepare seu pipeline para orquestração
A orquestração funciona melhor quando cada pipeline cobre uma unidade de trabalho distinta que o usuário deseja programar, validar ou executar de forma independente. Crie seus pipelines em torno desses limites para que um fluxo de trabalho possa coordená-los como tarefas separadas, incluindo o fluxo de controle apropriado entre as tarefas upstream e downstream.
Se você já tiver um pipeline grande que combina o trabalho que você deseja orquestrar separadamente, divida-o em pipelines menores movendo tabelas para um novo pipeline. Consulte Mover tabelas entre pipelines.
Lakeflow Jobs
É possível orquestrar várias tarefas no Lakeflow Jobs para implementar um fluxo de trabalho de processamento de dados. Para incluir um pipeline em um Job, use a tarefa Pipeline ao criar um Job. Consulte Tarefa de pipeline para Jobs.
Apache Airflow
Apache Airflow é uma solução de código aberto para gerenciamento e programação de fluxo de trabalho de dados. Airflow representa fluxo de trabalho como gráficos acíclicos direcionados (DAGs) de operações. Você define um fluxo de trabalho em um arquivo Python e Airflow gerencia a programação e execução. Para obter informações sobre como instalar e usar Airflow com Databricks, consulte Orquestrar trabalhos LakeFlow com Apache Airflow.
Para executar um pipeline como parte de um fluxo de trabalho Airflow , use o DatabricksSubmitRunOperator.
Requisitos
Os seguintes requisitos são necessários para usar o suporte do Airflow para Lakeflow pipelines:
- Versão 2.1.0 do Airflow ou mais tarde.
- O pacote do provedor Databricks versão 2.1.0 ou mais tarde.
Exemplo
O exemplo a seguir cria um DAG Airflow que aciona uma atualização para o pipeline com o identificador 8279d543-063c-4d63-9926-dae38e35ce8b:
from airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow'
}
with DAG('ldp',
start_date=days_ago(2),
schedule_interval="@once",
default_args=default_args
) as dag:
opr_run_now=DatabricksSubmitRunOperator(
task_id='run_now',
databricks_conn_id='CONNECTION_ID',
pipeline_task={"pipeline_id": "8279d543-063c-4d63-9926-dae38e35ce8b"}
)
Substitua CONNECTION_ID pelo identificador de uma conexãoAirflow com seu workspace.
Salve este exemplo no diretório airflow/dags e use a interface Airflow para view e acionar o DAG. Utilize a interface pipeline para view os detalhes da atualização pipeline .
Fábrica de Dados do Azure
LakeFlow Pipelines e Azure Data Factory incluem opções para configurar o número de novas tentativas quando ocorre uma falha. Se os valores de nova tentativa forem configurados no seu pipeline e na atividade do Azure Data Factory que chama o pipeline, o número de novas tentativas será o valor de nova tentativa do Azure Data Factory multiplicado pelo valor de nova tentativa do pipeline.
Por exemplo, se uma atualização de pipeline falhar, o pipeline tenta novamente a atualização até cinco vezes por default. Se a configuração de repetição do Azure Data Factory for definida como três, e seu pipeline usar o default de cinco tentativas, seu pipeline com falha poderá ser repetido até quinze vezes. Para evitar tentativas excessivas de repetição quando as atualizações de pipeline falham, a Databricks recomenda limitar o número de tentativas ao configurar o pipeline ou a atividade do Azure Data Factory que chama o pipeline.
Para alterar a configuração de nova tentativa do seu pipeline, use a configuração pipelines.numUpdateRetryAttempts ao configurar o pipeline.
Azure Data Factory é um serviço ETL baseado em cloudque permite orquestrar integração de dados e transformações de trabalho. Azure Data Factory oferece suporte direto à execução de tarefas Databricks em um fluxo de trabalho, incluindo Notebook, tarefa JAR e scripts Python . Você também pode incluir um pipeline em um fluxo de trabalho chamando a API REST do pipeline a partir de uma atividade Web do Azure Data Factory. Por exemplo, para acionar uma atualização de pipeline a partir do Azure Data Factory:
-
Crie uma fábrica de dados ou abra uma fábrica de dados existente.
-
Quando a criação for concluída, abra a página do seu data factory e clique no bloco Abrir Azure Data Factory Studio . A interface do usuário do Azure Data Factory é exibida.
-
Crie um novo pipeline Azure Data Factory selecionando pipeline no menu suspenso Novo na interface do usuário do Azure Data Factory Studio.
-
Na caixa de ferramentas Atividades , expanda Geral e arraste a atividade da Web para a tela do pipeline. Clique na tab Configurações e insira os seguintes valores:
Como prática recomendada de segurança, ao autenticar com ferramentas, sistemas, scripts e aplicativos automatizados, Databricks recomenda que você use access tokens pessoais pertencentes à entidade de serviço em vez de usuários workspace . Para criar tokens para entidade de serviço, consulte gerenciar tokens para uma entidade de serviço.
-
URL :
https://<databricks-instance>/api/2.0/pipelines/<pipeline-id>/updates.Substitua
<get-workspace-instance>.Substitua
<pipeline-id>pelo identificador do pipeline. -
Método : Selecione POST no menu suspenso.
-
Cabeçalhos : Clique em + Novo . Na caixa de texto Nome , digite
Authorization. Na caixa de texto Valor , insiraBearer <personal-access-token>.Substitua
<personal-access-token>por um access tokenpessoal Databricks . -
Corpo : Para passar parâmetros de solicitação adicionais, insira um documento JSON contendo os parâmetros. Por exemplo, para iniciar uma atualização e reprocessar todos os dados do pipeline:
{"full_refresh": "true"}. Se não houver parâmetros de solicitação adicionais, insira chaves vazias ({}).
Para testar a atividade da Web, clique em Depurar na barra de ferramentas do pipeline na interface do usuário do Data Factory. O resultado e o status da execução, incluindo erros, são exibidos na tab Saída do pipeline do Azure Data Factory. Utilize a interface do usuário do pipeline para view os detalhes da atualização pipeline .
Um requisito comum de fluxo de trabalho é iniciar uma tarefa após a conclusão de uma tarefa anterior. Como a solicitação do pipeline updates é assíncrona, retornando após iniciar a atualização, mas antes que ela seja concluída, as tarefas no seu pipeline do Azure Data Factory com uma dependência da atualização do pipeline devem aguardar a conclusão da atualização. Uma opção para aguardar a conclusão da atualização é adicionar uma atividade Until após a atividade Web que Trigger a atualização do pipeline. Na atividade Until:
- Adicione uma atividade Wait para aguardar um número configurado de segundos para a conclusão da atualização.
- Adicione uma atividade Web após a atividade de espera que utilize a solicitação de detalhes da atualização do pipeline para obter o status da atualização. O campo
statena resposta retorna o estado atual da atualização, incluindo se ela foi concluída. - Use o valor do campo
statepara definir a condição de término da atividade Until. Você também pode usar uma atividade Definir variável para adicionar uma variável de pipeline com base no valorstatee usar essa variável para a condição de término.