Considerações de produção para o Streaming Estruturado

Execute cargas de trabalho de produção do Streaming Estruturado como Trabalhos do Lakeflow agendados no Azure Databricks. Consulte Trabalhos do Lakeflow.

O Databricks recomenda que você sempre configure o seguinte:

  • Remova o código desnecessário dos notebooks que retornariam resultados, como display e count.
  • Não execute cargas de trabalho de streaming estruturado usando computação de uso geral. Sempre agende fluxos como Trabalhos do Lakeflow usando a computação de trabalhos.
  • Agendar Trabalhos do Lakeflow usando Continuous modo. Isso se refere ao recurso de agendamento de Trabalhos no Azure Databricks, e não ao intervalo de gatilho do Streaming Estruturado.
  • Não habilite o dimensionamento automático para computação para trabalhos de Streaming Estruturado.

Algumas cargas de trabalho se beneficiam do seguinte:

A Databricks introduziu os pipelines Lakeflow para reduzir a complexidade de gerenciar a infraestrutura de produção para cargas de trabalho do Structured Streaming. O Databricks recomenda o uso de pipelines do Lakeflow para novos pipelines de Streaming Estruturado. Consulte Pipelines Declarativos do Spark.

Observação

O dimensionamento automático de computação tem limitações ao reduzir o tamanho do cluster para cargas de trabalho do Streaming Estruturado. A Databricks recomenda usar os Pipelines Declarativos do Spark no Lakeflow com escalonamento automático aprimorado para cargas de trabalho de streaming. Consulte Otimização da utilização do cluster de pipeline do Lakeflow com dimensionamento automático.

:::observação Computação sem servidor

Na computação sem servidor, apenas Trigger.AvailableNow() e Trigger.Once() são suportados. O Databricks recomenda Trigger.AvailableNow().

Para streaming contínuo na computação sem servidor, use o Modo de pipeline disparado versus contínuo no modo contínuo.

Consulte Limitações de streaming.

:::

Projetar cargas de trabalho de streaming para esperar falhas

O Databricks recomenda que você sempre configure trabalhos de streaming para reiniciar automaticamente em caso de falha. Algumas funcionalidades, incluindo a evolução do esquema, exigem que as cargas de trabalho do Structured Streaming tentem novamente automaticamente. Consulte Configurar trabalhos de Streaming Estruturado para reiniciar consultas de streaming em caso de falha.

Algumas operações, como foreachBatch, fornecem garantias de pelo menos uma vez em vez de exatamente uma vez. Para essas operações, verifique se o pipeline de processamento é idempotente. Consulte Usar o foreachBatch para gravar nos coletores de dados arbitrários.

Observação

Quando uma consulta é reiniciada, o micro-lote planejado durante a execução anterior é processado. Se a tarefa falhou devido a um erro de falta de memória ou você cancelou manualmente uma tarefa devido a um microlote excessivamente grande, talvez seja necessário escalar a computação para processar o microlote de forma bem-sucedida.

Se você alterar as configurações entre as execuções, essas configurações se aplicarão ao primeiro novo lote planejado. Consulte Recuperar após alterações em uma consulta de Streaming Estruturado.

Quando uma tarefa é repetida

Você pode agendar várias tarefas como parte de um trabalho Azure Databricks. Ao configurar um trabalho usando o gatilho contínuo, você não pode definir dependências entre tarefas.

Você pode optar por planejar vários fluxos em um único trabalho usando uma das seguintes abordagens:

  • Várias tarefas: defina um trabalho com várias tarefas que executem cargas de trabalho de streaming usando o gatilho contínuo.
  • Várias consultas: defina várias consultas de streaming no código-fonte para uma única tarefa.

Você também pode combinar essas estratégias. A tabela a seguir compara essas abordagens.

Estratégia Várias tarefas Várias consultas
Como a computação é compartilhada? O Databricks recomenda o dimensionamento adequado dos recursos de computação para cada tarefa de streaming. Opcionalmente, você pode compartilhar a computação entre tarefas. Todas as consultas compartilham a mesma computação. Você pode, opcionalmente, atribuir consultas a pools de agendador.
Como as repetições são tratadas? Todas as tarefas devem falhar antes que o trabalho seja repetido. A tarefa será repetida se alguma consulta falhar.

Para obter mais detalhes sobre como trabalhar com várias tarefas ou consultas, consulte Executar várias consultas de Streaming Estruturado no mesmo cluster.

Configurar processos de Streaming Estruturado para reiniciar as consultas de streaming de dados em caso de falha

O Databricks recomenda que você configure todas as cargas de trabalho de streaming usando o gatilho contínuo. Consulte Executar trabalhos continuamente.

O gatilho contínuo tem o seguinte comportamento por padrão:

  • Impede mais de uma execução simultânea do trabalho.
  • Inicia uma nova execução quando uma execução anterior falha.
  • Usa retirada exponencial para repetições.

O Databricks recomenda sempre usar a computação de trabalhos em vez da computação para todas as finalidades ao agendar fluxos de trabalho. Em caso de falha e repetição do trabalho, novos recursos de computação são implantados.

Observação

O Databricks recomenda que você não use streamingQuery.awaitTermination() ou spark.streams.awaitAnyTermination(). Veja quando usar awaitTermination().

Quando usar awaitTermination()

streamingQuery.awaitTermination() e spark.streams.awaitAnyTermination() bloqueiam o thread atual até que uma consulta de streaming seja encerrada. Se usar essas funções depende do ambiente de execução.

Para o Lakeflow Jobs, não use streamingQuery.awaitTermination() ou spark.streams.awaitAnyTermination(). Essas funções não são necessárias porque o serviço Jobs impede automaticamente a execução quando uma consulta de streaming está ativa. Ambas as funções impedem que as células do notebook sejam preenchidas e impedem que o serviço de Trabalhos acompanhe a consulta de streaming, o que interrompe as métricas de backlog e as notificações de trabalho.

Use awaitTermination() nos seguintes casos:

Caso de uso Comportamento
Notebooks interativos na computação geral awaitTermination() mantém a célula de código em execução, e permite observar o estado da consulta, garantindo que as falhas sejam exibidas na saída do notebook.
Ambientes locais e de desenvolvimento Ao executar um programa Spark localmente, o processo é encerrado quando o thread principal é concluído. Chame awaitTermination() para manter o programa ativo até que a consulta de streaming seja concluída ou falhe.
Propagação de falha para o driver Sem awaitTermination(), uma falha de consulta de streaming em um contexto que não seja de trabalho pode não se propagar para o thread de chamada. A consulta pode falhar silenciosamente, tornando as falhas mais difíceis de detectar e diagnosticar. A chamada awaitTermination() gera novamente a exceção de consulta no driver.