Este tópico descreve como identificar e resolver problemas de latência em tarefas de sincronização em tempo real.
Identificar o gargalo: source ou destino
Para uma tarefa de sincronização em tempo real no DataStudio, acesse e clique em nome da tarefa para visualize seus detalhes. Para mais informações, consulte execute e gerencie tarefas de sincronização em tempo real.
Nos detalhes da execução, verifique a métrica Window Wait Time (5 min). Essa métrica indica o tempo que a tarefa aguardou para ler da source ou gravar no destino durante os últimos cinco minutos, ajudando a identificar o gargalo na sincronização de dados. Quando ocorre latência, o lado com o maior valor de métrica geralmente é o gargalo.
Verificar exceções do sistema
Após identificar o gargalo, acesse a aba Logs. Pesquise palavras-chave como "Error", "error", "Exception", "exception" ou "OutOfMemory" para encontrar pilhas de exceção do período de alta latência. Se encontrar uma exceção, use seus detalhes e consulte Lidar com erros comuns para verificar se a otimização da configuração da tarefa resolve o problema.
Uma tarefa de sincronização em tempo real lê dados de um sistema e os grava em outro. Se a gravação for mais lenta que a leitura, o sistema de destino pode exercer contrapressão sobre o sistema de source, causando lentidão. Isso significa que um gargalo em um sistema pode gerar exceções no outro. Priorize a investigação de exceções no sistema identificado como gargalo.
O código a seguir apresenta um exemplo típico de rastreamento de pilha de exceção:
java.lang.NullPointerException
at com.alibaba.streamx.core.util.EngineHelper.filterJobConfiguration(EngineHelper.java:31)
at com.alibaba.streamx.core.flink.trans.SinkFunctionAdaptor.open(SinkFunctionAdaptor.java:602)
at org.apache.flink.api.common.functions.util.FunctionUtils.openFunction(FunctionUtils.java:36)
at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.open(AbstractUdfStreamOperator.java:102)
at org.apache.flink.streaming.api.operators.StreamSink.open(StreamSink.java:48)
at org.apache.flink.streaming.runtime.tasks.StreamTask.openAllOperators(StreamTask.java:439)
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:288)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:747)
at java.lang.Thread.run(Thread.java:853)
Verificar erros frequentes de OOM
Nos detalhes da tarefa, acesse a aba Failover para verificar se há failovers frequentes (mais de uma ocorrência a cada 10 minutos). Caso os failovers sejam frequentes, inspecione as informações de exceção de cada evento. Se encontrar mensagens contendo a palavra-chave OutOfMemory, a tarefa tem memória insuficiente e enfrenta problemas recorrentes de falta de memória (OOM).
Para aumentar a memória, abra o editor de tarefas e incremente o valor de CU na área Running Resources.
Verificar distorção de dados na source ou necessidade de partições
Se a source for Kafka, DataHub ou LogHub e as etapas anteriores não revelarem exceções ou failovers, verifique se há distorção de dados no sistema de source. Verifique também se o tráfego de leitura das partições ou shards está atingindo o limite de taxa de sincronização.
Para sources Kafka, DataHub e LogHub, cada partição ou shard só pode ser consumido por um único thread paralelo. Se os dados estiverem concentrados em poucas partições ou shards enquanto outros estão quase vazios, essa distorção pode criar um gargalo de consumo e causar latência. Ajustar as configurações da tarefa não resolve esse problema. É necessário corrigir a distorção de dados na aplicação produtora upstream do sistema Kafka, DataHub ou LogHub. A latência será resolvida assim que a distribuição dos dados for equilibrada.
Na caixa de diálogo de detalhes da tarefa, alterne para a aba Running Information e revise a contagem total de bytes dos diferentes threads de leitura. Se um thread tiver uma contagem de bytes significativamente maior que os outros, provavelmente há distorção de dados. No entanto, a contagem total de bytes inclui todos os dados processados desde o último offset processado. Para uma tarefa de longa duração, essa métrica pode não refletir uma distorção recente. Verifique também as métricas de monitoramento no sistema de source para confirme se a distorção de dados está ocorrendo.
Caso o tráfego de dados de uma única partição ou shard atinja seu limite, aumente o número de partições ou shards no sistema de source para resolver a latência. Por exemplo, um cluster Kafka pode ter um limite de taxa configurado para leituras de partição, uma única partição do DataHub tem taxa máxima de leitura de 4 MB/s e um único shard do LogHub tem taxa máxima de leitura de 10 MB/s. Se uma tarefa de sincronização em tempo real exceder o limite de velocidade de leitura de uma única partição, expanda o número de partições ou shards no sistema de source para resolver a latência.
Se várias tarefas de sincronização em tempo real consumirem dados do mesmo tópico Kafka, tópico DataHub ou logstore LogHub, garanta que a velocidade combinada de leitura de todas as tarefas não exceda o limite da source.
Verificar transações grandes ou alterações frequentes no MySQL
Em uma tarefa de sincronização em tempo real com source MySQL, caso as etapas anteriores não revelem exceções ou failovers, verifique se o sistema de source está processando transações grandes ou sofrendo alterações frequentes, como diversas operações DML e DDL. Essas atividades podem fazer com que o log binário cresça mais rápido do que a tarefa consegue consumir, resultando em latência.
Por exemplo, atualizar um campo em toda uma tabela ou excluir uma grande quantidade de dados pode causar crescimento rápido do log binário. Na caixa de diálogo de detalhes da tarefa, alterne para a aba Running Information para visualize a velocidade de sincronização:
Uma velocidade de sincronização alta indica que o log binário está crescendo rapidamente.
Se a velocidade de sincronização não for alta, verifique as estatísticas do log binário e os logs de auditoria no servidor MySQL para confirme a taxa real de crescimento.
A velocidade de sincronização pode não refletir a taxa real na qual a tarefa consome o log binário do MySQL. Se uma transação ou alteração envolver bancos de dados ou tabelas não incluídos na configuração da tarefa, a tarefa filtrará esses dados após a leitura. Esses dados filtrados não são incluídos nas estatísticas de velocidade de sincronização ou volume de dados.
Ao confirme que transações grandes ou um aumento temporário nas alterações estão causando a latência, a tarefa eventualmente alcançará o fluxo normal após processar o acúmulo de mudanças.
Verificar alternâncias frequentes no particionamento dinâmico
Para tarefas de sincronização em tempo real que gravam no MaxCompute, caso selecione o particionamento dinâmico com base no conteúdo do campo, monitore cuidadosamente a coluna de source mapeada para a coluna de partição da tabela do MaxCompute. Dentro de um único Flush Interval (o padrão é 1 minuto) configurado no painel Basic Configurations, o número de valores distintos nessa coluna deve ser baixo.
Dentro do intervalo de flush, os dados destinados à tabela do MaxCompute são armazenados em cache em um conjunto de filas dentro da tarefa de sincronização em tempo real. Cada fila armazena dados em cache para uma operação de gravação do MaxCompute. O número máximo padrão de filas é cinco. Se o número de valores distintos da coluna de partição de source exceder esse limite dentro do intervalo de flush configurado, um flush imediato de todos os dados em cache será acionado. Operações frequentes de flush degradam severamente o desempenho de gravação.
confirme se flushes frequentes estão sendo acionados porque as filas de cache de partição da tabela do MaxCompute foram esgotadas. Na caixa de diálogo de detalhes da tarefa, alterne para a aba Logs e pesquise pela mensagem uploader map size has reached uploaderMapMaximumSize.
Aumentar a concorrência ou ative a execução distribuída
Se as etapas anteriores mostrarem que a latência se deve ao aumento do tráfego de source e não a exceções, mitigue o problema aumentando a concorrência da tarefa.
Ao aumentar a concorrência, aumente também a memória da tarefa. Como regra geral, adicione 1 GB de memória para cada quatro threads paralelos adicionais.
configure a concorrência e a memória da tarefa da seguinte forma:
Para tarefas de sincronização em tempo real ETL de tabela única para tabela única criadas no DataStudio, clique em Basic Configurations à direita para configure a concorrência e a memória da tarefa. No painel Basic Settings, configure parâmetros como Synchronization Method, a opção distributed execution mode, resource group, CUs e number of parallel threads. Expanda Advanced Settings para configure o flush interval (padrão: 60000 ms) e MaxCompute Channel Resources. Por exemplo, selecione Streaming Tunnel e defina o slot number.
Para outras tarefas do DataStudio, como migrações de banco de dados para o DataHub, configure o número de threads paralelos na etapa Configure Resource e a memória no painel Basic Configurations.
Para tarefas de solução de sincronização, configure o número de threads paralelos e a memória na etapa Configure Resource.
Se o modo de execução distribuída estiver desativado, defina o número de threads paralelos como 32 ou menos. Definir a contagem acima de 20 pode causar latência devido a gargalos de recursos em uma única máquina. Para canais específicos, ative o modo de execução distribuída para melhorar o desempenho. Os canais que suportam o modo de execução distribuída estão listados na tabela a seguir.
|
Tipo de tarefa |
source |
Destino |
|
Tarefa ETL do DataStudio |
Kafka |
MaxCompute |
|
Tarefa ETL do DataStudio |
Kafka |
Hologres |