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.
Os pipelines Lakeflow fornecem um framework declarativo para criar pipelines de dados em lote e em fluxo contínuo usando SQL e Python. Seus principais conceitos são pipelines, fluxos, tabelas de streaming, exibições materializadas e coletores, que trabalham juntos para processar dados com orquestração automática e atualizações incrementais.
Os pipelines do Lakeflow estendem o SDP (Pipelines Declarativos do Apache Spark™). Para saber mais sobre o SDP e como ele se compara com os pipelines do Lakeflow, consulte o Apache Spark Declarative Pipelines.
Tip
É novo em pipelines? Comece com Como usar pipelines do Lakeflow para entender como usar pipelines ao longo do ciclo de vida deles, e por quê, com links para as tarefas em cada etapa.
Note
Os pipelines do Lakeflow exigem o plano Premium. Entre em contato com sua equipe de conta do Databricks para obter mais informações.
Quais são os benefícios dos pipelines?
Em contraste com o desenvolvimento de processos de engenharia de dados com as APIs Apache Spark e Spark Structured Streaming no Databricks Runtime, com orquestração manual por meio de Lakeflow Jobs, a natureza declarativa dos pipelines oferece os seguintes benefícios:
- Orquestração automática: os pipelines executam etapas de processamento (chamadas de "fluxos") na ordem correta, com paralelismo máximo, e repetem automaticamente em caso de falhas transitórias de forma progressiva; da tarefa do Spark ao fluxo, até o pipeline inteiro.
- Processamento declarativo: as funções declarativas reduzem centenas de linhas do spark manual e do código de streaming estruturado para algumas. A API AUTO CDC manipula eventos CDC (Change Data Capture), incluindo SCD Tipo 1 e Tipo 2, sem código manual para eventos fora de ordem ou conceitos de streaming, como marcas d'água.
- Processamento incremental: um mecanismo de processamento incremental mantém as exibições materializadas atuais: você escreve a lógica de transformação com semântica em lote e o mecanismo reprocessa apenas dados de origem novos ou alterados quando possível.
Conceitos principais
O diagrama a seguir ilustra os conceitos mais importantes dos pipelines.
Datasets
Um pipeline produz três tipos de conjuntos de dados, cada um com semântica de processamento diferente:
| Tipo de conjunto de dados | Como os registros são processados |
|---|---|
| Tabela de streaming | Cada registro é processado exatamente uma vez, assumindo uma fonte de dados do tipo apenas inclusão. As tabelas de streaming são adequadas para ingestão e processamento incremental de dados em crescimento contínuo. |
| Visão materializada | Os resultados são recomputados conforme necessário para refletir o estado atual dos dados. Visões materializadas são adequadas para transformações, agregações ou resultados pré-computados consumidos por vários conjuntos de dados subsequentes. |
| View | Avaliada sob demanda; não é persistida. Use exibições para transformações intermediárias e verificações que não precisam ser publicadas em um catálogo. |
Uma tabela de streaming é uma forma de tabela gerenciada do Catálogo do Unity que também é um destino de streaming. Uma tabela de streaming pode ter um ou mais fluxos de streaming (Acréscimo, CDC AUTOMÁTICO) gravados nela. Você pode definir fluxos de streaming explicitamente e separadamente de sua tabela de streaming de destino ou implicitamente como parte de uma definição de tabela de streaming.
Uma Exibição materializada também é uma forma de tabela gerenciada do Catálogo do Unity e é um destino em lote. Uma exibição materializada pode ter um ou mais fluxos de exibição materializados gravados nele. As exibições materializadas diferem das tabelas de streaming, pois você sempre define os fluxos implicitamente como parte da definição de exibição materializada.
Para obter detalhes, consulte tabelas de fluxo e visões materializadas.
Quando usar visões, visões materializadas e tabelas de streaming
Ao implementar consultas de pipeline, escolha o tipo de conjunto de dados que melhor se ajusta ao seu caso de uso.
Considere usar uma visualização para:
- Divida uma consulta grande ou complexa em consultas mais fáceis de gerenciar.
- Valide os resultados intermediários usando as expectativas.
- Reduza os custos de armazenamento e computação para resultados que você não precisa manter. Como as tabelas são materializadas, elas requerem recursos adicionais de computação e armazenamento.
Considere o uso de uma exibição materializada quando:
- Diversas consultas downstream consomem a tabela. Como uma exibição materializada armazena em cache seus resultados, as consultas downstream leem os resultados pré-compilados em vez de computar novamente a consulta em cada acesso.
- Outros pipelines, trabalhos ou consultas consomem a tabela. Como uma visualização materializada é materializada em uma tabela do Catálogo do Unity, consumidores externos ao pipeline que a define podem consultá-la. Visualizações comuns não são materializadas; portanto, só podem ser utilizadas dentro do mesmo pipeline.
- Você deseja inspecionar os resultados de uma consulta durante o desenvolvimento. Como uma visualização materializada é materializada e pode ser consultada fora do pipeline, é possível validar a correção dos cálculos durante o desenvolvimento. Após a validação, converta as consultas que não requerem materialização em exibições.
- Sua consulta executa agregações ou junções ou os dados de origem podem ser alterados devido a atualizações e exclusões, em vez de apenas crescer. Uma visualização materializada mantém seus resultados consistentes com o estado atual dos dados de origem, enquanto uma tabela de streaming é projetada para fontes append-only e processa cada registro uma única vez.
Considere o uso de uma tabela de streaming quando:
- Uma consulta é definida em relação a uma fonte de dados que está crescendo de forma contínua ou incremental.
- Os resultados da consulta devem ser computados de forma incremental.
- O pipeline precisa de alta taxa de transferência e baixa latência.
Note
As tabelas de streaming são sempre definidas em relação a fontes de streaming. Você também pode utilizar fontes de streaming com AUTO CDC ... INTO para aplicar atualizações de feeds da CDA. Consulte As APIs AUTO CDC: Simplifique a captura de dados de alterações com pipelines.
Flows
Um fluxo é o conceito fundamental de processamento de dados em pipelines e oferece suporte à semântica de streaming e de processamento em lote. Um fluxo lê dados de uma fonte, aplica a lógica de processamento definida pelo usuário e grava o resultado em um destino. Os pipelines compartilham o mesmo tipo de fluxo de streaming (Acréscimo, Atualização, Conclusão) que o Streaming Estruturado do Spark. (Atualmente, somente os fluxos de Acréscimo e Atualização são expostos.) Para obter mais detalhes, consulte os modos de saída em Streaming Estruturado.
Os pipelines também fornecem tipos de fluxo adicionais:
- AUTO CDC é um fluxo de streaming exclusivo nos pipelines Lakeflow que lida com eventos CDC fora de ordem e oferece suporte a SCD Tipo 1 e SCD Tipo 2. O CDC automático não está disponível no SDP.
- A exibição materializada é um fluxo em lote em pipelines que processa apenas novos dados e alterações nas tabelas de origem sempre que possível.
Para obter detalhes, consulte Carregar e processar dados de forma incremental com fluxos de pipeline do Lakeflow.
Sinks
Um destino é um alvo de streaming para um pipeline e oferece suporte a tabelas Delta, tópicos do Apache Kafka, tópicos do Azure Event Hubs e fontes de dados Python personalizadas. Um destino pode receber a gravação de um ou mais fluxos de streaming (Append, Update)
Para obter detalhes, consulte Coletores em pipelines do Lakeflow.
Pipelines
Um pipeline é a unidade de desenvolvimento e execução e é o contêiner para os fluxos, tabelas de streaming, exibições materializadas e coletores que você define. Você cria um pipeline definindo esses objetos no código-fonte do pipeline e executando o pipeline. Enquanto o pipeline é executado, ele analisa as dependências de seus objetos definidos e orquestra automaticamente sua ordem de execução e paralelização.
Para obter detalhes, consulte O que são pipelines?.
Também é possível definir visões materializadas autônomas e tabelas de streaming fora de um pipeline Lakeflow, em que o Azure Databricks gerencia o pipeline para você. Para comparar as duas abordagens, consulte Pipelines autônomos vs. pipelines do Lakeflow.
Um pipeline é executado em modo acionado ou contínuo, o que determina se ele atualiza os dados disponíveis e se interrompe ou mantém as tabelas atualizadas à medida que novos dados chegam. Para comparar os dois modos, consulte Modo de pipeline acionado vs. contínuo.
Ingestão de dados
Os pipelines dão suporte a todas as fontes de dados disponíveis no Azure Databricks. O Databricks recomenda o uso de tabelas de streaming para a maioria dos casos de uso de ingestão. Para arquivos no armazenamento de objetos em nuvem, o Auto Loader oferece carregamento incremental e idempotente. Para dados em streaming, os pipelines podem captar dados diretamente de barramentos de mensagens, como Apache Kafka, Hubs de Eventos do Azure, Amazon Kinesis e Google Pub/Sub. Consulte Carregar dados em pipelines.
Qualidade dos dados
As expectativas são cláusulas opcionais em conjuntos de dados que validam os dados à medida que fluem pelo pipeline. Você define uma expectativa como uma restrição booliana do SQL e especifica o que acontece quando um registro falha: avisar, remover o registro ou falhar na atualização. Confira Gerenciar a qualidade dos dados com as expectativas do pipeline.
Integração delta
Todas as tabelas criadas e gerenciadas por pipelines são tabelas Delta. Eles oferecem as mesmas garantias do Delta Lake, incluindo transações ACID, viagem no tempo e imposição de esquema. Os pipelines adicionam propriedades de tabela extras e realizam manutenção automática usando otimização preditiva, incluindo operações OPTIMIZE e VACUUM. Veja O que é o Delta Lake em Azure Databricks?.