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.
Os pipelines do Auto Loader requerem monitorização ativa para detetar problemas como acumulação de atrasos, desvio de esquema, dados corrompidos e fluxos parados antes de afetarem os consumidores a jusante. Esta página descreve como monitorizar métricas-chave, consultar o estado ao nível do ficheiro, construir painéis de observabilidade e resolver problemas comuns.
Para detalhes de configuração em produção, consulte Configurar Auto Loader para cargas de trabalho de produção. Para melhores práticas de configuração, consulte as melhores práticas do Auto Loader.
Pré-requisitos
Vários fluxos de trabalho de monitorização nesta página dependem de cloud_files_state() para observar o estado da ingestão de cada ficheiro — incluindo consultas de backlog, cálculos de latência e deteção de deriva do esquema.
cloud_files_state() é uma função com valor de tabela que retorna o estado de ingestão ao nível de ficheiro para um checkpoint do Auto Loader. Nem todos os seus campos estão disponíveis por defeito. A disponibilidade depende da versão e configuração do seu Databricks Runtime:
-
Databricks Runtime 18.2 e posteriores:
discovery_time,processed_time, ecommit_timeestão disponíveis automaticamente. No Databricks Runtime 16.4–18.1, estes campos só estão disponíveis quandocloudFiles.cleanSourceestão ativados. -
Databricks Runtime 16.4 e superiores com
cloudFiles.cleanSourceativados:archive_time,archive_mode, emove_locationestão disponíveis.
Ativar cloudFiles.cleanSource tem alguma sobrecarga de desempenho. Compare com as suas cargas de trabalho num ambiente de pré-produção antes de o ativar em produção.
Additionally:
- Anote os dados ingeridos com a
_metadatacoluna. Capturar pelo menosfile_pathefile_modification_time. Consulte Coluna de metadados do arquivo. - Ativar
_rescued_datae_corrupt_recordcolunas.
Métricas-chave do Auto Loader
A tabela seguinte resume as métricas mais importantes a monitorizar para pipelines de Auto Loader. Estas métricas estão disponíveis nos eventos de progresso StreamingQueryListener, com valores específicos do Auto Loader disponibilizados no mapa metrics de cada origem.
| Métrica | O que te diz |
|---|---|
numFilesOutstanding |
Número de ficheiros no backlog à espera de serem processados |
numBytesOutstanding |
Tamanho do backlog de ficheiros em bytes |
approximateQueueSize |
Profundidade da fila na nuvem (apenas no modo de notificação de ficheiros) |
numInputRows |
Linhas processadas por lote |
inputRowsPerSecond |
Taxa de chegada de dados |
processedRowsPerSecond |
Taxa de processamento |
durationMs Quebra |
Onde o tempo é despendido em cada lote |
O que observar
Os padrões seguintes indicam que o seu pipeline pode precisar de atenção.
-
Crescimento
numFilesOutstanding: O atraso está a acumular-se. O seu pipeline está a ficar atrasado em relação aos dados recebidos. -
processedRowsPerSecond<inputRowsPerSecond: O pipeline está a processar os dados mais devagar do que chegam. -
Grande
durationMs.latestOffset: A descoberta de ficheiros é lenta. Considere mudar para eventos de ficheiro. -
Grande:
durationMs.addBatchO processamento de dados é lento. Considere escalar computação ou otimizar transformações.
Para a referência completa das métricas, veja Métricas de origem do Auto Loader.
Consultar estado ao nível do ficheiro com cloud_files_state
A cloud_files_state() função de valores de tabela fornece informações detalhadas sobre cada ficheiro descoberto pelo Auto Loader. Os seguintes campos estão disponíveis. Os campos assinalados como exigindo o Databricks Runtime 16.4 ou superior, ou 18.2 ou superior, apenas são preenchidos nas condições descritas em Pré-requisitos.
| Campo | Tipo | Description |
|---|---|---|
path |
STRING |
O percurso do ficheiro |
size |
BIGINT |
O tamanho do arquivo em bytes |
create_time |
TIMESTAMP |
Quando o ficheiro foi criado |
discovery_time |
TIMESTAMP |
Quando o Auto Loader descobriu o ficheiro (Databricks Runtime 16.4 e superior) |
processed_time |
TIMESTAMP |
Quando o Auto Loader processava o ficheiro (Databricks Runtime 16.4 e superior) |
commit_time |
TIMESTAMP |
Quando o ficheiro foi registado no checkpoint (Databricks Runtime 16.4 e posteriores) |
archive_time |
TIMESTAMP |
Quando o ficheiro foi arquivado (requer cloudFiles.cleanSource) |
archive_mode |
STRING |
MOVE, DELETE, ou NULL (requer cloudFiles.cleanSource) |
move_location |
STRING |
Caminho de destino quando cloudFiles.cleanSource é MOVE |
ingestion_state |
STRING |
Estado atual de ingestão de ficheiros |
Verificar o estado da ingestão de ficheiros
As consultas seguintes abrangem cenários de diagnóstico comuns.
Encontre todos os ficheiros não processados (o backlog atual):
SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';
Calcule a latência média de ingestão (tempo desde a criação do ficheiro até ao commit):
SELECT avg(unix_timestamp(commit_time) - unix_timestamp(create_time)) AS avg_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL AND create_time IS NOT NULL;
Encontrar ficheiros corrompidos ou omitidos:
SELECT path, ingestion_state, size, create_time
FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state LIKE 'SKIPPED%';
Acompanhar o progresso do arquivo (requerendo cloudFiles.cleanSource):
SELECT archive_mode, count(*) AS file_count
FROM cloud_files_state('path/to/checkpoint')
GROUP BY archive_mode;
Encontre ficheiros com latência elevada da descoberta ao commit para identificar gargalos:
SELECT
path,
size,
unix_timestamp(commit_time) - unix_timestamp(discovery_time) AS processing_latency_seconds,
unix_timestamp(commit_time) - unix_timestamp(create_time) AS end_to_end_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL
ORDER BY end_to_end_latency_seconds DESC
LIMIT 20;
Para a referência completa de SQL, veja cloud_files_state função de valores de tabela.
Monitorize o Auto Loader nas canalizações do Lakeflow
A Databricks recomenda a utilização de pipelines Lakeflow para pipelines Auto Loader de produção. Para tirar partido das suas capacidades de monitorização integradas:
Armazene o registo de eventos dos pipelines Lakeflow numa tabela Delta para que possa ser consultado para dados de observabilidade. Configure isto através das definições avançadas do pipeline ou da API. Para detalhes, consulte o registo de eventos do pipeline.
Estruture o seu pipeline com vista à observabilidade. Um pipeline do Auto Loader bem estruturado nos pipelines Lakeflow inclui uma
{table}_sourcevista (a definição da origem do Auto Loader), uma{table}_bronzetabela de streaming (ingestão de dados brutos com as colunas_rescued_datae_corrupt_record), umacorrupt_records_sinkque coloca em quarentena as linhas com dados impossíveis de analisar e uma{table}vista limpa para consumo a jusante.Defina expectativas nas suas tabelas de streaming bronze para monitorizar desvios de esquema e corrupção de dados.
_rescued_data IS NULLdeteta alterações inesperadas no esquema e_corrupt_record IS NULLdeteta dados não analisáveis. Os oleodutos de fluxo de lago avaliam estas expectativas à medida que os dados chegam e geram um rasto de observabilidade. Pode configurar expectativas para emitir um aviso, eliminar linhas ou fazer com que o pipeline falhe.
Depois de criar a event_log_raw vista para o seu pipeline, use as seguintes consultas para métricas específicas do Auto Loader.
Monitorizar a taxa de ingestão por fluxo:
SELECT
origin.flow_name,
origin.update_id,
timestamp,
TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS rows_written
FROM event_log_raw
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC;
Monitorizar o atraso de dados por fluxo:
SELECT
origin.flow_name,
timestamp,
DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
ORDER BY timestamp DESC;
Resuma as violações de expectativas para detetar desvio do esquema e dados corrompidos:
SELECT
origin.flow_name,
explode(from_json(
details:flow_progress.data_quality.expectations,
'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
)) AS expectation
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.data_quality.expectations IS NOT NULL;
Para orientações gerais sobre monitorização de pipelines Lakeflow, consulte Monitorizar pipelines e Registo de eventos de pipelines.
Monitor Auto Loader com Streaming Estruturado
Ao executar o Auto Loader fora dos pipelines do Lakeflow, utilize as seguintes abordagens de monitorização de streaming estruturado.
- Implemente um
StreamingQueryListenerpara capturar métricas específicas do Auto Loader de cada lote, lendo a partir desource.metrics.
from pyspark.sql.streaming import StreamingQueryListener
class AutoLoaderMonitor(StreamingQueryListener):
def onQueryStarted(self, event):
pass
def onQueryProgress(self, event):
for source in event.progress.sources:
if "CloudFilesSource" in source.description:
metrics = source.metrics
files_outstanding = metrics.get("numFilesOutstanding", "0")
bytes_outstanding = metrics.get("numBytesOutstanding", "0")
rows_per_sec = source.processedRowsPerSecond
# Push metrics to your monitoring system (for example, write to a Delta table)
def onQueryIdle(self, event):
pass
def onQueryTerminated(self, event):
pass
spark.streams.addListener(AutoLoaderMonitor())
Observação
A lógica de processamento nos ouvintes pode abrandar o processamento das consultas. Limite o processamento nas rotinas de retorno do recetor e evite operações de escrita externas síncronas nessas rotinas; em vez disso, emita telemetria de baixo custo de forma assíncrona ou encaminhe as métricas para uma tarefa separada para persistência.
Utilize
numInputRows,inputRowsPerSecondeprocessedRowsPerSeconddo progresso da origem para calcular a taxa de processamento — ficheiros por segundo e linhas por segundo para cada lote.Para calcular a latência de ingestão, compare
create_timeecommit_timedecloud_files_state()para obter a latência de ponta a ponta. Relativamente à latência de processamento, use a decomposiçãodurationMs(por exemplo,latestOffset,addBatche outras fases de lote indicadas) para identificar qual é a fase que constitui o estrangulamento.Utilize
df.observe()para definir métricas de qualidade de dados em linha diretamente no DataFrame de streaming. As métricas são visíveis emStreamingQueryListenereventos em progresso sobobservedMetrics.
from pyspark.sql.functions import count, lit, col
observed_df = df.observe(
"auto_loader_quality",
count(lit(1)).alias("total_rows"),
count(col("_rescued_data")).alias("rescued_rows"),
count(col("_corrupt_record")).alias("corrupt_rows")
)
- Use
.queryName()para atribuir um nome único a cada fluxo, facilitando a distinção dos fluxos do Auto Loader no separador de Streaming da Spark UI e nos painéis de monitorização.
Para a referência completa de monitorização do Structured Streaming, consulte Monitorizar consultas do Structured Streaming no Azure Databricks.
Construir um painel de observabilidade
Combine dados de múltiplas fontes para construir um painel de observabilidade abrangente para os seus pipelines do Auto Loader. Esta tabela apresenta algumas fontes sugeridas que pode usar para estruturar o seu painel de observabilidade.
| Fonte de dados | Dados de observabilidade |
|---|---|
cloud_files_state() |
Estado de ingestão por ficheiro: marcas temporais de descoberta, processamento, confirmação e arquivamento por ficheiro |
| Registo de eventos dos pipelines do Lakeflow | Histórico de execução do pipeline, métricas de fluxo por lote e resultados das expectativas de qualidade dos dados |
| Tabelas de saída do pipeline | Contagens de linhas e volume de dados escritos para cada tabela ingerida |
Pode então agregar dados de observabilidade em tabelas dedicadas que servem de base para painéis de controlo e alertas:
- Resumir os estados da execução do pipeline (sucesso ou falha) ao longo do tempo, derivados de
event_type = 'update_progress'eventos. - Métricas agregadas de ingestão de ficheiros (tamanho do backlog, throughput, latência por lote), derivadas de
cloud_files_state()eevent_type = 'flow_progress'eventos. - Criar estatísticas das tabelas com base nas contagens de linhas e no volume de dados por tabela, obtidos a partir de
num_output_rowsno registo de eventos. - Recolha informações de depuração a partir de registos detalhados de erros e violações de expectativas por atualização, derivadas de
event_type = 'flow_progress'eventos comdata_qualitypopulado.
Estas tabelas agregadas podem alimentar um painel de controlo de IA/BI e alertas SQL. Os painéis de dashboard recomendados incluem a linha temporal do estado da execução do pipeline, tendência do backlog de ingestão, tendência de throughput, distribuição de latência de ingestão, métricas de qualidade de dados, eventos de evolução do esquema e estado de arquivo de ficheiros.
Monitorizar eventos de evolução de esquemas
Use as seguintes abordagens para detetar alterações no esquema à medida que ocorrem.
- Valores não-NULL em
_rescued_datanas contagens de violações esperadas indicam deriva do esquema. Consultar o registo de eventos defailed_records > 0na expectativano rescued data. - Alterações ao
_schemasdiretório dentro do configuradocloudFiles.schemaLocation(ou dentro do checkpoint apenas quando a localização do esquema não está definida separadamente) indicam que ocorreu evolução do esquema. Pode consultar este diretório a partir de um trabalho de monitorização separado. - Não trate um evento
onQueryTerminatedseguido deonQueryStartedcom o mesmo nome de fluxo como prova suficiente, por si só, de evolução do esquema. Os fluxos reiniciam por várias razões (reinício do cluster, implementação de código, erros de armazenamento transitório). Correlaciona os reinícios com sinais independentes —_schemasalterações no diretório ou_rescued_dataviolações de expectativas — antes de concluir que ocorreu evolução do esquema. - Use
_metadata.file_pathpara identificar quais os ficheiros que introduziram alterações no esquema. Junte isto aocloud_files_state()nopathcampo para correlacionar alterações de esquema com ficheiros e lotes específicos.
Use esta consulta de exemplo para detetar desvios recentes do esquema através de violações de expectativas:
SELECT
timestamp,
origin.flow_name,
exp.name AS expectation_name,
exp.failed_records
FROM (
SELECT
timestamp,
origin,
explode(from_json(
details:flow_progress.data_quality.expectations,
'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
)) AS exp
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.data_quality.expectations IS NOT NULL
)
WHERE exp.name = '<rescued-data expectation name>'
AND exp.failed_records > 0
ORDER BY timestamp DESC;
Configurar alertas para problemas comuns
Use alertas SQL do Databricks ou notificações de pipeline para detetar problemas antes que afetem os consumidores a jusante.
O SQL seguinte deteta um atraso crescente e pode ser usado como base para um alerta SQL do Databricks. Programe-o para correr periodicamente (por exemplo, a cada 5 minutos) e alerte quando o resultado não estiver vazio.
-- Alert when backlog exceeds threshold or trends upward across recent batches
WITH recent_backlog AS (
SELECT
origin.flow_name,
timestamp,
DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
)
SELECT flow_name, backlog_bytes, timestamp
FROM recent_backlog
WHERE rn = 1
AND backlog_bytes > 1073741824 -- alert when backlog exceeds 1 GB
A tabela seguinte resume as condições de alerta recomendadas:
| O que detetar | Como detetá-lo | Quando alertar |
|---|---|---|
| Crescente atraso acumulado |
numFilesOutstanding tendência ascendente |
Aumento sustentado ao longo de vários lotes |
| Fluxo de água parado | Eventos sem progresso | Sem eventos durante N minutos (com base no intervalo de gatilho esperado) |
| Alta latência de ingestão | commit_time - create_time |
Ultrapassa o teu limiar de SLA |
| Degradação da qualidade dos dados | Taxa de falha esperada | Percentagem crescente de filas que falham nas expectativas |
| Evento de evolução do esquema | _rescued_data IS NOT NULL |
Quaisquer valores não NULL na violação de expectativa contam |
| Descoberta lenta de ficheiros | durationMs.latestOffset |
Significativamente acima da linha de base |
Resolver problemas comuns
A tabela seguinte descreve problemas comuns no pipeline do Auto Loader, as suas prováveis causas e as ações recomendadas para os resolver.
| Issue | Causa possível | Ação recomendada |
|---|---|---|
| O backlog cresce mais rápido do que o processamento | Capacidade de computação insuficiente, distribuição desigual dos dados ou limites de taxa restringidos | Escala o cálculo, verifica o desfasamento com a interface do Spark e revê maxFilesPerTrigger as definições para controlar o tamanho do lote |
| Ficheiros não descobertos | Eventos de ficheiros mal configurados, problema de permissões ou fluxo não executado nos últimos 7 dias | Verifique permissões de localização externa, verifique os eventos de ficheiros configurados na interface do Unity Catalog e assegure que o fluxo corre pelo menos a cada 7 dias para evitar a expiração do estado do RocksDB |
| O arranque da transmissão está a demorar demasiado tempo | Grande transferência do estado do ponto de verificação (RocksDB) | Atualize para Databricks Runtime 15.3 e superiores para carregamento de estado assíncrono, o que reduz o tempo de arranque em ~90% |
| Processamento de ficheiros duplicados | Definições cloudFiles.maxFileAge agressivas ou corrupção do ponto de verificação |
Utilize um maxFileAge conservador (mínimo de 90 dias ou mais), verifique a integridade dos pontos de controlo e evite políticas de gestão do ciclo de vida no armazenamento de pontos de controlo |
| Evolução do esquema a provocar reinícios do pipeline | Alterações frequentes ou incompatíveis de esquemas | Reveja schemaEvolutionMode, mude para addNewColumnsWithTypeWidening para promoções de tipo ou use o tipo Variant para esquemas altamente dinâmicos |
| Dados corrompidos a acumular-se no sumidouro | Questões de qualidade dos dados de origem | Verifique o _corrupt_record sumidouro de quarentena para padrões, reveja a geração de dados de origem e considere adicionar validação a montante |
discovery_time e commit_time não preenchidos |
Em execução no Databricks Runtime inferior a 18.2 sem cleanSource |
Atualize para Databricks Runtime 18.2 e superiores ou ative cloudFiles.cleanSource no Databricks Runtime 16.4–18.1 |
Para resolução de problemas adicionais, consulte o FAQ do Auto Loader.