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.
Para melhores práticas abrangentes na configuração do Auto Loader, incluindo seleção de modos de descoberta de ficheiros, gestão de esquemas e tratamento da qualidade dos dados, consulte as melhores práticas do Auto Loader.
A Databricks recomenda o uso do Auto Loader nos pipelines do Lakeflow para ingestão incremental de dados. Os pipelines do Lakeflow ampliam as funcionalidades do Structured Streaming do Apache Spark e permitem-lhe escrever apenas algumas linhas de código declarativo em Python ou SQL para implementar um pipeline de dados pronto para produção com:
- Dimensionamento automático da infraestrutura de computação para reduzir custos:Otimizar a utilização do cluster do pipeline Lakeflow com dimensionamento automático
- Verificações de qualidade dos dados com expectativas:Gerir a qualidade dos dados com as expectativas do fluxo de trabalho
- Gestão automática da evolução de esquemas:Configurar inferência e evolução de esquemas no Auto Loader
- Monitorização através de métricas no registo de eventos:Registo de eventos do pipeline
A Databricks também recomenda que sigas as melhores práticas de streaming para correr o Auto Loader em produção. Consulte Considerações sobre produção para Streaming estruturado.
Note
Os oleodutos Lakeflow são a forma recomendada de usar o Auto Loader para a maioria das ingestão de produção. Se a sua carga de trabalho não tiver requisitos de baixa latência e a sua prioridade for minimizar o custo de computação, pode em vez disso agendar o Auto Loader como um trabalho em lote desencadeado que utiliza Trigger.AvailableNow.
Veja Considerações de custo.
Monitorização Auto Loader
As secções seguintes descrevem como monitorizar o Auto Loader em produção, incluindo métricas, registos, alertas e fluxos de trabalho comuns de resolução de problemas. Para uma referência completa sobre padrões de painel, análise de latência e deteção de alterações ao esquema, consulte Monitorizar e observar o Auto Loader.
Consultando arquivos descobertos pelo Auto Loader
O Auto Loader fornece uma API SQL para inspecionar o estado de um fluxo. Usando a cloud_files_state função, você pode encontrar metadados sobre arquivos que foram descobertos por um fluxo Auto Loader. Consulta cloud_files_state, fornecendo a localização do ponto de verificação associada a um fluxo do Auto Loader.
Note
A cloud_files_state função está disponível no Databricks Runtime 11.3 LTS e superior.
SELECT * FROM cloud_files_state('path/to/checkpoint');
Ouvir atualizações de streaming
Para monitorizar ainda mais os fluxos do Auto Loader, a Databricks recomenda a utilização da interface Streaming Query Listener do Apache Spark. Consulte Monitorização de consultas de streaming estruturado no Azure Databricks.
O Auto Loader reporta métricas ao Streaming Query Listener em cada lote. Você pode visualizar quantos arquivos existem nas pendências e qual é o tamanho das pendências nas métricas numFilesOutstanding e numBytesOutstanding na guia de Dados Brutos no painel de progresso da consulta de streaming.
{
"sources": [
{
"description": "CloudFilesSource[/path/to/source]",
"metrics": {
"numFilesOutstanding": "238",
"numBytesOutstanding": "163939124006"
}
}
]
}
Ao usar o modo de notificação de ficheiros no Databricks Runtime 10.4 LTS ou superior, as métricas incluem também o número aproximado de eventos de ficheiros na fila na nuvem como approximateQueueSize para AWS e Azure.
Considerações de custo
Ao executar o Auto Loader, as suas principais fontes de custo são os recursos computacionais e a descoberta de ficheiros.
Se a sua carga de trabalho não tiver requisitos de baixa latência, pode reduzir os custos de computação usando Lakeflow Jobs para agendar o Auto Loader como tarefas em lote, em Trigger.AvailableNow vez de o executar continuamente. Consulte Configurar intervalos de ativação do Streaming Estruturado. Estes trabalhos em lote podem ser acionados usando disparadores de chegada de ficheiros para reduzir ainda mais a latência entre a chegada e o processamento dos ficheiros.
Os custos de descoberta de ficheiros podem surgir na forma de operações LIST nas suas contas de armazenamento no modo de listagem de diretório, bem como de solicitações de API no serviço de subscrição e no serviço de fila no modo de notificação de ficheiros. Acionadores contínuos como Trigger.ProcessingTime são especialmente dispendiosos no modo de listagem de diretórios, uma vez que o Auto Loader lista continuamente todo o diretório para encontrar novos ficheiros. Se a sua carga de trabalho exigir gatilhos contínuos, a Databricks recomenda escolher um modo de descoberta de ficheiros com base nos requisitos de latência:
- Baixa latência e simplicidade: Use o Auto Loader com eventos de ficheiros. Os eventos de ficheiro requerem apenas uma fila por bucket e utilizam descobertas incrementais em execuções subsequentes. Para mais informações, consulte Auto Loader com visão geral dos eventos de ficheiros.
- Aplicações muito sensíveis à latência: Use o modo clássico de notificação de ficheiros. O modo Clássico lê diretamente da fila na cloud sem o salto adicional de cache introduzido pelos eventos do ficheiro. Neste modo, pode marcar recursos criados pelo Auto Loader para acompanhar os seus custos usando etiquetas de recursos. Para mais detalhes, consulte notificação de ficheiro.
Retenção de dados de origem
Note
Disponível em Databricks Runtime 16.4 LTS e superiores.
À medida que os ficheiros se acumulam na sua pasta de origem, os custos de armazenamento aumentam e a descoberta de ficheiros torna-se mais lenta, especialmente no modo de listagem. O Auto Loader oferece a cloudFiles.cleanSource opção de gerir automaticamente a retenção de ficheiros, arquivando ou eliminando ficheiros após serem processados.
Arquivamento de ficheiros no diretório de origem para reduzir custos
Warning
- A definição
cloudFiles.cleanSourceapaga ou move ficheiros na pasta de origem. - Se usares
foreachBatchpara o processamento de dados, os teus ficheiros tornam-se candidatos a mover ou eliminar assim que a operaçãoforeachBatchregressa com sucesso, mesmo que a tua operação tenha consumido apenas um subconjunto dos ficheiros do lote.
A Databricks recomenda usar o Auto Loader com eventos de ficheiros para reduzir os custos de descoberta. Isto também reduz os custos de computação porque a descoberta é incremental.
Se não puder usar eventos de ficheiros e tiver de usar a listagem de diretórios para descobrir ficheiros, pode usar a cloudFiles.cleanSource opção de arquivar ou eliminar automaticamente os ficheiros após o Auto Loader processá-los para reduzir os custos de descoberta. Como o Auto Loader limpa ficheiros do seu diretório de origem após o processamento, é necessário listar menos ficheiros durante a descoberta.
Ao utilizar cloudFiles.cleanSource com a opção MOVE, considere os seguintes requisitos:
- Tanto o diretório de origem como o diretório de destino da operação de movimentação têm de estar na mesma localização externa, volume ou ponto de montagem do DBFS. Os movimentos entre buckets e entre contentores não são suportados e resultam em erro.
- O destino da movimentação pode ser um caminho de volume (por exemplo,
/Volumes/my_catalog/my_schema/my_volume/archive/). - Se o seu diretório de origem e destino estiverem na mesma localização externa, não devem ter diretórios irmãos que contenham armazenamento gerido (por exemplo, um volume ou catálogo gerido). Nestes casos, o Auto Loader não consegue obter as permissões necessárias para escrever no diretório de destino.
O Databricks recomenda usar esta opção quando:
- O seu diretório fonte acumula um grande número de ficheiros ao longo do tempo.
- Deve manter os ficheiros processados para conformidade ou auditoria (defina
cloudFiles.cleanSourcecomoMOVE). - Quer reduzir os custos de armazenamento removendo ficheiros após a ingestão (definido
cloudFiles.cleanSourceparaDELETE). Ao usar oDELETEmodo, a Databricks recomenda ativar a versão no bucket para que as eliminações do Auto Loader funcionem como eliminações suaves e estejam disponíveis em caso de má configuração. Além disso, a Databricks recomenda configurar políticas de ciclo de vida na cloud para eliminar versões antigas e apagadas suavemente após um período de carência especificado (como 60 ou 90 dias), com base nas suas necessidades de recuperação.
Para obter a referência completa sobre as opções de cleanSource e os respetivos valores predefinidos, consulte Limpar ficheiros processados com o Auto Loader.
Mover ficheiros processados para um caminho de armazenamento a frio
O exemplo seguinte configura o Auto Loader para mover ficheiros processados para um diretório de arquivo dentro do mesmo bucket após 14 dias. Pode aplicar uma política de ciclo de vida cloud no caminho do arquivo para transferir ficheiros para níveis de armazenamento mais baratos (por exemplo, AWS S3 Glacier, Azure Cool/Archive ou GCS Coldline/Archive).
Python
# Step 1: Configure Auto Loader to move processed files to an archive path.
checkpoint = "/Volumes/my_catalog/my_schema/my_volume/checkpoints/ingest_stream"
archive_path = "s3://my-bucket/archive/landing/"
df = (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.cleanSource", "MOVE")
.option("cloudFiles.cleanSource.moveDestination", archive_path)
.option("cloudFiles.cleanSource.retentionDuration", "14 days")
.option("cloudFiles.schemaLocation", checkpoint)
.load("s3://my-bucket/landing/")
)
# Step 2: Write to a Delta table.
(df.writeStream
.option("checkpointLocation", checkpoint)
.trigger(availableNow=True)
.toTable("my_catalog.my_schema.raw_events")
)
# Step 3 (outside Databricks): Set up a cloud lifecycle policy on the
# archive path to transition files to cold storage after a grace period.
# For example, in AWS you can configure an S3 Lifecycle rule to move
# objects under s3://my-bucket/archive/landing/ to S3 Glacier after
# 30 days.
SQL
-- Step 1: Configure Auto Loader to move processed files to an archive path
-- using a Lakeflow Declarative Pipeline.
CREATE OR REFRESH STREAMING TABLE raw_events
AS SELECT * FROM STREAM read_files(
's3://my-bucket/landing/',
format => 'json',
cleanSource => 'MOVE',
cleanSourceMoveDestination => 's3://my-bucket/archive/landing/',
cleanSourceRetentionDuration => '14 days'
);
-- Step 2 (outside Databricks): Set up a cloud lifecycle policy on the
-- archive path to transition files to cold storage.
-- For example, in AWS configure an S3 Lifecycle rule to move objects
-- under s3://my-bucket/archive/landing/ to S3 Glacier after 30 days.
Usando Trigger.AvailableNow e limitação de velocidade
Note
Disponível em Databricks Runtime 10.4 LTS e superior.
O Auto Loader pode ser programado para ser executado nos Lakeflow Jobs como um trabalho em lote usando Trigger.AvailableNow. O AvailableNow disparador instrui o Auto Loader a processar todos os ficheiros que chegaram antes do início da consulta. Os ficheiros novos que chegam após o início da transmissão são ignorados até ao próximo disparador.
Com Trigger.AvailableNow, a descoberta de arquivos ocorre de forma assíncrona com o processamento de dados, e os dados podem ser processados em vários microlotes com limitação de taxa. O Auto Loader por padrão processa um máximo de 1000 arquivos a cada microlote. Você pode configurar cloudFiles.maxFilesPerTrigger e cloudFiles.maxBytesPerTrigger para definir quantos arquivos ou quantos bytes devem ser processados em um microlote. O limite de arquivo é um limite rígido, mas o limite de bytes é um limite flexível, o que significa que mais bytes podem ser processados do que o maxBytesPerTriggerfornecido. Quando as duas opções são fornecidas em conjunto, o Auto Loader processa quantos arquivos são necessários para atingir um dos limites.
Localização do ponto de verificação
O local do ponto de verificação é usado para armazenar as informações de estado e progresso do fluxo. O Databricks recomenda definir o local do ponto de verificação para um local sem uma política de ciclo de vida do objeto na nuvem. Se os arquivos no local do ponto de verificação forem limpos de acordo com a política, o estado do fluxo de dados ficará corrompido. Se isso acontecer, você deve reiniciar o fluxo do zero.
Monitorização de eventos de ficheiro
O Auto Loader acompanha os arquivos descobertos no local do ponto de verificação usando o RocksDB para fornecer garantias de ingestão exatamente uma vez. Para fluxos de ingestão de alto volume ou longa duração, o Databricks recomenda a atualização para o Databricks Runtime 15.4 LTS ou superior. Nessas versões, o Auto Loader não espera que todo o estado RocksDB seja baixado antes que o fluxo seja iniciado, o que pode acelerar o tempo de inicialização do fluxo.
Se quiser impedir que os estados de arquivo cresçam sem limites, você também pode considerar usar a cloudFiles.maxFileAge opção para expirar eventos de arquivo com mais de uma determinada idade. O valor mínimo que você pode definir para cloudFiles.maxFileAge é "14 days". As exclusões no RocksDB aparecem como entradas de lápide. Assim, poderá observar o uso do armazenamento aumentar temporariamente à medida que os eventos expiram antes de se estabilizar.
Warning
cloudFiles.maxFileAge é fornecido como um mecanismo de controlo de custos para conjuntos de dados de grande volume. Ajustar cloudFiles.maxFileAge de forma muito agressiva pode causar problemas na qualidade dos dados, como ingestão duplicada ou arquivos ausentes. Portanto, a Databricks recomenda uma configuração conservadora para cloudFiles.maxFileAge, como 90 dias, que é semelhante ao que soluções comparáveis de ingestão de dados recomendam.
Tentar ajustar a cloudFiles.maxFileAge opção pode levar a arquivos não processados sendo ignorados pelo Auto Loader ou arquivos já processados expirando e, em seguida, sendo reprocessados causando dados duplicados. Aqui estão algumas coisas a considerar ao escolher um cloudFiles.maxFileAge:
- Se o fluxo for reiniciado após um longo tempo, os eventos de notificação retirados da fila que são mais antigos do que
cloudFiles.maxFileAgesão ignorados. Da mesma forma, se utilizares a listagem de diretórios, os ficheiros que podem ter surgido durante o tempo de inatividade e que são mais antigos quecloudFiles.maxFileAgeserão ignorados. - Se você usar o modo de listagem de diretório e usar
cloudFiles.maxFileAge, por exemplo, definido como"1 month", interrompe o fluxo e reinicia o fluxo comcloudFiles.maxFileAgedefinido como"2 months", os arquivos com mais de 1 mês, mas mais recentes que 2 meses, são reprocessados.
Se você definir essa opção na primeira vez que iniciar o fluxo, você não ingerirá dados mais antigos do que cloudFiles.maxFileAge, portanto, se você quiser ingerir dados antigos, não deve definir essa opção ao iniciar o fluxo pela primeira vez. No entanto, deve definir essa opção em corridas subsequentes.
Acione preenchimentos automáticos regulares usando cloudFiles.backfillInterval
Em casos raros, os arquivos podem ser perdidos ou atrasados quando dependem exclusivamente de sistemas de notificação, como quando os limites de retenção de mensagens de notificação são atingidos. Se tiver exigências rigorosas de integridade de dados e SLA, considere definir cloudFiles.backfillInterval para acionar backfills assíncronos em um intervalo especificado. Por exemplo, defina-o como um dia para preenchimentos diários ou uma semana para preenchimentos semanais. Acionar preenchimentos regulares não causa duplicatas.
Ao usar eventos de arquivo, execute seu fluxo pelo menos uma vez a cada 7 dias
Ao usar eventos de arquivo, execute seus fluxos do Auto Loader pelo menos uma vez a cada 7 dias para evitar uma listagem completa de diretórios. Executar seus fluxos do Auto Loader com frequência garantirá que a descoberta de arquivos seja incremental.
Para melhores práticas abrangentes para eventos de ficheiros geridos, consulte Melhores práticas para Auto Loader com eventos de ficheiro.