Este tópico descreve as métricas do Flink totalmente gerenciado.
Observações
Discrepâncias de dados entre o CloudMonitor e o console do Flink
Diferenças nas dimensões exibidasO console do Flink utiliza consultas PromQL para exibir apenas a latência máxima. Em cenários de computação em tempo real, a latência média pode mascarar problemas graves, como desbalanceamento de dados ou bloqueios em partições únicas. Portanto, somente a latência máxima oferece insights operacionais valiosos.
Discrepâncias de valoresO CloudMonitor emprega um mecanismo de pré-agregação para calcular métricas. O valor "máximo" no CloudMonitor pode divergir ligeiramente do valor em tempo real no console do Flink devido a diferenças nas janelas de agregação, timestamps de amostragem ou lógica de cálculo. Para solução de problemas, utilize os dados do console do Flink como source da verdade.
Latência de dados e configuração de watermark
-
Lógica de cálculo de latênciaA métrica de monitoramento atual Emit Delay é calculada com base no tempo de evento, utilizando a seguinte fórmula:
Delay = Tempo Atual do Sistema - Campo de tempo lógico no registro de dados (ex.: PriceData.time)
Isso significa que a métrica reflete a atualização dos dados, e não a velocidade de processamento do sistema. Essa métrica apresenta valores elevados quando os dados de origem são antigos ou quando o sistema pausa a saída para alinhar watermarks.
-
Recomendações
Cenário 1: Sua lógica de negócio depende de watermarks para garantir correção, mas os dados de origem são antigos
-
Situações típicas:
A entrega de dados upstream possui atraso inerente (ex.: relatórios lentos de eventos).
Você está executando um backfill para processar dados de um dia anterior.
A lógica de negócio exige watermarks para lidar com eventos fora de ordem, portanto eles não podem ser desativados.
Fenômeno: Alertas de monitoramento indicam alta latência, mas o grupo de consumidores Kafka não apresenta lag (lag ≈ 0) e a carga de CPU é baixa.
-
Recomendações:
Ignore esta métrica de latência: Neste caso, um atraso elevado é esperado, pois reflete a idade dos dados. Isso não indica uma falha no sistema.
Adote uma métrica diferente: Monitore o Kafka consumer lag. Se o lag do consumidor não aumentar continuamente, o sistema possui capacidade de processamento suficiente e não requer intervenção.
Cenário 2: Você necessita de baixa latência e tolera pequenos eventos fora de ordem ou perda de dados
-
Situações típicas:
Em aplicações como painéis de tela grande ou controle de risco em tempo real, a espera induzida por watermark retarda a saída.
A lógica de negócio prioriza o momento de recebimento dos dados (tempo de processamento) em vez do timestamp dentro do registro de dados (tempo de evento).
Fenômeno: O fluxo de dados é em tempo real, mas como o watermark está configurado com uma janela de tolerância ampla (ex.: permissão de atraso de 10 segundos), a saída sofre um atraso de 10 segundos.
-
Recomendações:
Remova ou desative watermarks: Passe a utilizar o tempo de processamento para cálculos ou defina o limiar de espera do watermark como 0.
Resultado esperado: A métrica de latência cairá significativamente, aproximando-se do tempo real de processamento. Os dados serão processados assim que chegarem, sem espera por alinhamento.
-
Características das métricas
As métricas refletem apenas o estado atual de um componente e são insuficientes para determinar a causa raiz de um problema. Para um diagnóstico abrangente, utilize sempre o monitor de backpressure da UI do Flink e outras ferramentas.
1. Backpressure de operador
Sintoma: Operadores downstream não conseguem processar dados com rapidez suficiente, fazendo com que a source reduza sua taxa de emissão.
Como identificar: Utilize o monitor de backpressure da UI do Flink para detectar esse problema.
-
Características da métrica:
sourceIdleTimeaumenta periodicamente.currentFetchEventTimeLagecurrentEmitEventTimeLagaumentam continuamente.Caso extremo: Se um operador estiver completamente travado,
sourceIdleTimeaumentará continuamente.
2. Gargalo de desempenho na source
Sintoma: A source lê na velocidade máxima, mas não atende às demandas de processamento de dados.
Como identificar: Nenhum backpressure é detectado no job.
-
Características da métrica:
sourceIdleTimepermanece em um valor muito baixo (indicando que a source opera com capacidade total).currentFetchEventTimeLagecurrentEmitEventTimeLagsão semelhantes e permanecem elevados.
3. Desbalanceamento de dados ou partições vazias
Sintoma: A distribuição de dados é desigual entre as partições Kafka upstream ou algumas partições estão vazias.
Como identificar: Compare métricas entre diferentes subtasks da source.
-
Características da métrica:
O
sourceIdleTimede uma subtask específica da source é significativamente maior que o das demais, indicando que essa instância paralela está ociosa.
4. Latência de dados
Sintoma: A latência geral do job é alta. É necessário determinar se o gargalo está na source ou em um sistema externo.
Como identificar: Analise combinadamente o tempo ocioso, a diferença entre métricas de lag e o tamanho do backlog.
-
Características da métrica:
**Alto
sourceIdleTime:Indica que a source está ociosa, o que geralmente significa que a taxa de produção de dados do sistema externo** é baixa, e não que o Flink está processando lentamente.-
Análise da diferença de lag:Compare a diferença entre
currentEmitEventTimeLagecurrentFetchEventTimeLag. Essa diferença representa o tempo que os dados passam dentro do operador source:Pequena diferença (valores próximos): Indica capacidade insuficiente de busca. O gargalo normalmente é largura de banda de I/O de rede ou paralelismo insuficiente da source.
Grande diferença: Indica capacidade insuficiente de processamento. O gargalo geralmente é parsing ineficiente de dados ou backpressure de operadores downstream.
**
pendingRecords(se suportado pelo conector):Esta métrica reflete diretamente o backlog externo**. Um valor mais alto indica um backlog de dados mais severo no sistema externo.