Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Problemas de desempenho de jobs

Última atualização: Jun 27, 2026

Este tópico aborda problemas comuns de desempenho em jobs.

Como dividir nós de operador?

Na página Operation Center > Job O&M, clique em o nome do job desejado. Na aba Deployment Details, na seção Runtime Parameter Settings, adicione o código abaixo em Other Settings e salve para aplicar.

pipeline.operator-chaining: 'false'

Quais são as técnicas de otimização do Group Aggregate?

  • Ative o MiniBatch (melhora o throughput)

    O MiniBatch armazena dados recebidos em buffer antes de acionar o processamento. Isso reduz a frequência de acesso ao State, aumenta o throughput e diminui o volume de saída.

    O MiniBatch dispara o processamento de microlotes com base em mensagens de evento inseridas no intervalo especificado na source.

    • Cenários

      O uso de microlotes troca uma latência ligeiramente maior por um throughput significativamente superior. Não ative essa opção se você precisar de latência ultrabaixa. Para a maioria dos cenários de agregação, ativar o MiniBatch melhora bastante o desempenho do sistema.

    • Como ativar

      O MiniBatch vem desativado por padrão. Para ativá-lo, acesse a aba Deployment Details do job desejado e, na seção Runtime Parameter Settings, em Other Settings, adicione o seguinte código.

      table.exec.mini-batch.enabled: true
      table.exec.mini-batch.allow-latency: 5s

      A tabela a seguir explica os parâmetros.

      Parâmetro

      Descrição

      table.exec.mini-batch.enabled

      Define se o mini-batch deve ser ativado.

      table.exec.mini-batch.allow-latency

      Intervalo de tempo entre as saídas dos lotes.

  • Ative o LocalGlobal (resolve problemas comuns de hotspot de dados)

    O mecanismo LocalGlobal usa o LocalAgg para pré-agregar dados enviesados, reduzindo a pressão de hotspots no GlobalAgg e melhorando o desempenho geral.

    O LocalGlobal divide uma única agregação em duas fases: Local e Global — semelhante às fases Combine e Reduce no MapReduce. Na primeira fase, os nós upstream armazenam e agregam dados localmente (localAgg), emitindo acumuladores incrementais. Na segunda fase, esses acumuladores são mesclados (Merge) para produzir o resultado final (GlobalAgg).

    • Cenários

      Melhora o desempenho de agregações padrão (como SUM, COUNT, MAX, MIN e AVG) e resolve problemas de hotspot de dados nesses cenários.

    • Limitações

      O LocalGlobal é ativado por padrão, mas apresenta as seguintes limitações:

      • O MiniBatch precisa estar ativado.

      • Sua AggregateFunction deve implementar Merge.

    • Verifique se a alteração entrou em vigor

      Confira se a topologia gerada contém nós chamados GlobalGroupAggregate ou LocalGroupAggregate.

  • Ative o PartialFinal (resolve problemas de hotspot de COUNT DISTINCT)

    Para lidar com hotspots de COUNT DISTINCT, tradicionalmente era necessário reescrever manualmente as consultas em uma agregação de dois estágios (adicionando uma camada de shuffling baseada em módulo). O Realtime Compute for Apache Flink agora oferece shuffling automático de COUNT DISTINCT por meio da otimização PartialFinal — sem necessidade de reescrita manual.

    O LocalGlobal funciona bem para agregações padrão, mas traz pouco benefício para COUNT DISTINCT. Durante a agregação local, as taxas de deduplicação de chaves distintas permanecem baixas, então os hotspots persistem no nó global.

    • Cenários

      Use quando o COUNT DISTINCT não atender aos requisitos de desempenho do nó de agregação.

      Importante
      • Não use a otimização PartialFinal em Flink SQL que inclua UDAFs.

      • Evite o PartialFinal quando o volume de dados for pequeno — ele introduz shuffling de rede desnecessário e desperdiça recursos.

    • Como ativar

      Esse recurso vem desativado por padrão. Para ativá-lo, na aba Deployment Details do job desejado, na seção Runtime Parameter Settings, insira o seguinte código no campo Other Configurations.

      table.optimizer.distinct-agg.split.enabled: true
    • Verifique se as alterações surtiram efeito.

      Confira se a topologia gerada mudou de um único estágio de agregação para dois estágios.

  • Reescreva AGG WITH CASE WHEN como sintaxe AGG WITH FILTER (aumenta o desempenho em cenários com múltiplos COUNT DISTINCT)

    Se seu job calcula UV em várias dimensões — como UV total, UV de cliente móvel e UV de PC — use a sintaxe padrão AGG WITH FILTER em vez de CASE WHEN. O otimizador SQL do Realtime Compute reconhece parâmetros Filter, permitindo que múltiplas operações COUNT DISTINCT no mesmo campo compartilhem State e reduzam E/S de State. Testes de desempenho mostram que essa reescrita pode dobrar o desempenho.

    • Cenários

      Ganhos significativos de desempenho ocorrem ao calcular múltiplos resultados COUNT DISTINCT no mesmo campo sob condições diferentes.

    • Texto original

      COUNT(distinct visitor_id) as UV1 , COUNT(distinct case when is_wireless='y' then visitor_id else null end) as UV2
    • Sintaxe otimizada

      COUNT(distinct visitor_id) as UV1 , COUNT(distinct visitor_id) filter (where is_wireless='y') as UV2

Quais são as técnicas de otimização do TopN?

  • Algoritmos TopN

    Se a entrada do TopN for um fluxo append-only (por exemplo, do SLS), apenas um algoritmo está disponível: AppendRank. Se a entrada for um fluxo de atualização (por exemplo, após AGG ou JOIN), dois algoritmos estão disponíveis, classificados do maior para o menor desempenho: UpdateFastRank e RetractRank. Os nomes dos algoritmos aparecem nos rótulos dos nós da topologia.

    • AppendRank: Suportado apenas para fluxos append-only.

    • UpdateFastRank: Ideal para fluxos de atualização.

    • RetractRank: Algoritmo de fallback para fluxos de atualização. Apresenta menor desempenho. Em alguns casos, pode ser otimizado para UpdateFastRank.

    Para otimizar o RetractRank para UpdateFastRank, três condições devem ser atendidas:

    • O fluxo de entrada deve ser um fluxo de atualização.

    • O fluxo de entrada deve incluir informações de Primary Key — por exemplo, após uma agregação GROUP BY.

    • Os campos de ordenação devem atualizar monotonicamente na direção oposta à classificação. Por exemplo, ORDER BY COUNT, COUNT_DISTINCT ou SUM (valores positivos) DESC.

    Para garantir o uso do UpdateFastRank com ORDER BY SUM DESC, adicione uma condição de filtro que garanta que total_fee seja positivo.

    insert
      into print_test
    SELECT
      cate_id,
      seller_id,
      stat_date,
      pay_ord_amt  -- Omit rownum to reduce sink table output.
    FROM (
        SELECT
          *,
          ROW_NUMBER () OVER (
            PARTITION BY cate_id,
            stat_date  -- Include a time field to prevent data corruption from State TTL.
            ORDER
              BY pay_ord_amt DESC
          ) as rownum  -- Sort by upstream sum result.
        FROM (
            SELECT
              cate_id,
              seller_id,
              stat_date,
              -- Critical: Declare all SUM inputs as positive, ensuring monotonic increase.
              -- This allows TopN to use the optimized algorithm and fetch only top 100 records.
              sum (total_fee) filter (
                where
                  total_fee >= 0
              ) as pay_ord_amt
            FROM
              random_test
            WHERE
              total_fee >= 0
            GROUP
              BY cate_name,
              seller_id,
              stat_date,
              cate_id
          ) a
        ) WHERE
          rownum <= 100;
  • Métodos de otimização do TopN

    • Otimização sem saída de classificação

      Se a saída do seu TopN não precisar exibir valores de rownum, omita-os e ordene apenas uma vez no frontend. Isso reduz drasticamente o volume de saída da tabela sink. Para mais detalhes, consulte Top-N.

    • Aumente o tamanho do cache do TopN

      O TopN utiliza uma camada de Cache de State para melhorar a eficiência de acesso ao State. A taxa de acerto do cache é calculada da seguinte forma.

      cache_hit = cache_size*parallelism/top_n/partition_key_num

      Por exemplo, com Top100, tamanho de cache 10.000, paralelismo 50 e 100.000 chaves de partição, a taxa de acerto é de apenas 10000*50/100/100000=5%. Taxas de acerto baixas fazem com que a maioria das solicitações acesse o State baseado em disco, criando falhas nas métricas de busca de state e degradando severamente o desempenho.

      Quando a cardinalidade da chave de partição é muito alta, aumente o tamanho do cache do TopN e a memória heap adequadamente. Para mais detalhes, consulte Configurar definições de implantação de job.

      table.exec.rank.topn-cache-size: 200000

      O tamanho padrão do cache é 10.000. Aumentá-lo para 200.000 eleva a taxa teórica de acerto para 200000*50/100/100000 = 100%.

    • Inclua um campo baseado em tempo no PartitionBy

      Para classificações diárias, inclua um campo Day. Sem ele, o TTL do State pode corromper os resultados finais do TopN.

Quais são as soluções eficientes de deduplicação?

Os dados de source no Realtime Compute for Apache Flink às vezes contêm duplicatas. Usuários frequentemente solicitam deduplicação. O Realtime Compute suporta duas estratégias: manter a primeira linha (Deduplicate Keep FirstRow) e manter a última linha (Deduplicate Keep LastRow).

  • Sintaxe

    O SQL não possui sintaxe direta de deduplicação, portanto usamos ROW_NUMBER OVER WINDOW para implementá-la. A deduplicação é essencialmente uma forma especial de TopN.

    SELECT *
    FROM (
       SELECT *,
        ROW_NUMBER() OVER (PARTITION BY col1[, col2..]
         ORDER BY timeAttributeCol [asc|desc]) AS rownum
       FROM table_name)
    WHERE rownum = 1

    Parâmetro

    Descrição

    ROW_NUMBER()

    Função de janela que atribui números de linha começando em 1.

    PARTITION BY col1[, col2..]

    Opcional. Colunas que definem partições (chaves de deduplicação).

    ORDER BY timeAttributeCol [asc

    desc])

    Coluna usada para ordenação. Deve ser um campo de atributo de tempo (Proctime ou Rowtime). Use ordem crescente para Keep FirstRow ou decrescente para Keep LastRow.

    rownum

    Apenas rownum=1 ou rownum<=1 é suportado.

    Conforme mostrado acima, a deduplicação requer duas camadas de consulta:

    1. Use ROW_NUMBER() para ordenar os dados por atributo de tempo e atribuir classificações.

      • Se o campo de ordenação for Proctime, o Flink deduplica pelo tempo do sistema, produzindo resultados não determinísticos.

      • Se o campo de ordenação for Rowtime, o Flink deduplica pelo tempo de negócio, produzindo resultados determinísticos.

    2. Filtre pela classificação para manter apenas a primeira linha, alcançando a deduplicação.

      Os dados podem ser ordenados em ordem crescente ou decrescente pela coluna de tempo:

      • Deduplicate Keep FirstRow: Ordem crescente, mantém a primeira linha.

      • Deduplicate Keep LastRow: Ordem decrescente, mantém a primeira linha.

  • Deduplicate Keep FirstRow

    Esta estratégia mantém a primeira ocorrência de cada chave e descarta duplicatas subsequentes. Ela armazena apenas dados de chave no State, oferecendo melhor desempenho. Exemplo:

    SELECT *
    FROM (
      SELECT *,
        ROW_NUMBER() OVER (PARTITION BY b ORDER BY proctime) as rowNum
      FROM T
    )
    WHERE rowNum = 1

    Este exemplo deduplica a tabela T pelo campo b, mantendo a primeira linha pelo tempo do sistema. Aqui, proctime é um campo de atributo Processing Time na tabela source T. Ao deduplicar pelo tempo do sistema, você pode simplificar proctime para a chamada de função proctime() e omitir a declaração explícita do campo.

  • Deduplicate Keep LastRow

    Esta estratégia mantém a última ocorrência de cada chave. Ela apresenta desempenho ligeiramente superior ao LAST_VALUE. Exemplo:

    SELECT *
    FROM (
      SELECT *,
        ROW_NUMBER() OVER (PARTITION BY b, d ORDER BY rowtime DESC) as rowNum
      FROM T
    )
    WHERE rowNum = 1

    Este exemplo deduplica a tabela T pelos campos b e d, mantendo a última linha pelo tempo de negócio. Aqui, rowtime é um campo de atributo Event Time na tabela source T.

O que devo observar ao usar funções integradas?

  • Substitua funções definidas pelo usuário por funções integradas

    O Realtime Compute otimiza continuamente as funções integradas. Prefira-as em vez de funções definidas pelo usuário. As principais otimizações incluem:

    • Redução da sobrecarga de serialização e desserialização.

    • Operações diretas no nível de byte.

  • Use separadores de caractere único em funções KEY VALUE

    Assinatura KEY VALUE: KEYVALUE(content, keyValueSplit, keySplit, keyName). Quando keyValueSplit e keySplit são caracteres únicos (como dois pontos ":" ou vírgula ","), o sistema usa um algoritmo otimizado para localizar diretamente keyName nos dados binários sem dividir todo o conteúdo — melhorando o desempenho em cerca de 30%.

  • Observações sobre a operação LIKE

    • Para StartWith, use LIKE 'xxx%'.

    • Para EndWith, use LIKE '%xxx'.

    • Para Contains, use LIKE '%xxx%'.

    • Para Equals, use LIKE 'xxx', equivalente a str = 'xxx'.

    • Para corresponder ao underscore (_), faça o escape: LIKE '%seller/_id%' ESCAPE '/'. O underscore (_) é um curinga de caractere único em SQL. Sem o escape, LIKE '%seller_id%' corresponde a seller_id, seller#id, sellerxid e seller1id, causando resultados incorretos.

  • Evite funções de expressão regular (REGEXP)

    Expressões regulares são extremamente custosas — frequentemente 100× mais lentas que aritmética básica — e podem entrar em loops infinitos sob certas condições, bloqueando jobs. Consulte Regex execution is too slow. Prefira LIKE. As funções de expressão regular incluem:

Como resolver baixa eficiência e backpressure durante a leitura de tabela completa?

O backpressure pode decorrer de processamento lento no downstream. Primeiro, verifique se há backpressure no downstream. Se houver, resolva-o usando uma das opções a seguir:

  • Aumente a concorrência.

  • Ative otimizações de agregação como minibatch (para nós de agregação downstream).

O que significam os indicadores coloridos em Status Durations das subtasks de vértice na visão geral do job?

Na página Overview do job, clique em um nó de operador e selecione a aba SubTasks para visualizar badges coloridos de duração na coluna Status Durations.

Status Durations mostra o tempo gasto pelas subtasks de vértice em cada fase. Os significados das cores são:

  • image.png: CREATED

  • image.png: SCHEDULED

  • image.png: DEPLOYING

  • image.png: INITIALIZING

  • image.png: RUNNING

O que é a thread RMI TCP Connection e por que ela consome muito mais CPU que outras threads?

Na lista de monitoramento de threads ordenada por uso de CPU, a thread RMI TCP Connection(62)-172.25.240.255 mostra status RUNNABLE com 82,3% de uso de CPU — muito superior às threads kafkaRequestSource (16,9%–26,4% de CPU, majoritariamente TIMED_WAITING).

As threads RMI TCP Connection pertencem ao framework RMI (Remote Method Invocation) integrado do Java e lidam com chamadas de método remoto. O uso de CPU flutua dinamicamente. Picos de curta duração não indicam carga alta sustentada. Observe o uso de CPU ao longo do tempo. A análise do gráfico de chama (abaixo) mostra que as threads RMI consomem quase nenhuma CPU.

image

Por que Low Watermark, Watermark e Task InputWatermark na topologia em execução apresentam diferença de tempo em relação ao horário atual?

  • Motivo 1: Watermark da tabela source declarada com TIMESTAMP_LTZ (TIMESTAMP(p) WITH LOCAL TIME ZONE) causa diferenças de tempo.

    Os exemplos a seguir comparam o comportamento do Watermark com tipos TIMESTAMP_LTZ versus TIMESTAMP.

    • Watermark da tabela source usa tipo TIMESTAMP_LTZ.

      CREATE TEMPORARY TABLE s1 (
        a INT,
        b INT,
        ts as CURRENT_TIMESTAMP,-- CURRENT_TIMESTAMP generates TIMESTAMP_LTZ.
        WATERMARK FOR ts AS ts - INTERVAL '5' SECOND 
      ) WITH (
        'connector'='datagen',
        'rows-per-second'='1',
        'fields.b.kind'='random','fields.b.min'='0','fields.b.max'='10'
      );
      CREATE TEMPORARY TABLE t1 (
        k INT,
        ts_ltz timestamp_ltz(3),
        cnt BIGINT
      ) WITH ('connector' = 'print');
      -- Output results.
      INSERT INTO t1
      SELECT b, window_start, COUNT(*) FROM
      TABLE(
          TUMBLE(TABLE s1, DESCRIPTOR(ts), INTERVAL '5' SECOND))
      GROUP BY b, window_start, window_end;
      Nota

      A sintaxe Legacy Window produz resultados idênticos a TVF Window (Table-Valued Function). Exemplo de sintaxe legada:

      SELECT b, TUMBLE_END(ts, INTERVAL '5' SECOND), COUNT(*) FROM s1 GROUP BY TUMBLE(ts, INTERVAL '5' SECOND), b;

      Após implantar e executar o job no console de desenvolvimento do Realtime Compute, observe uma diferença de tempo de 8 horas entre o Watermark e o horário atual (usando UTC+8 como referência).

      • Watermark & Low Watermark

        Na aba Watermarks da UI de monitoramento de job do Flink, a SubTask 0 mostra valor de Watermark 1706778525521, correspondente a Datetime of Watermark Timestamp 02-01 09:08:45. O horário de início do job foi 02-01 17:03:04 — uma diferença de ~8 horas. O painel do operador à esquerda mostra Low Watermark idêntico: 02-01 09:08:45.

      • Task InputWatermark

        image

    • Watermark da tabela source usa tipo TIMESTAMP (TIMESTAMP(p) WITHOUT TIME ZONE).

      CREATE TEMPORARY TABLE s1 (
        a INT,
        b INT,
        -- Simulate TIMESTAMP without timezone, starting at 2024-01-31 01:00:00 and incrementing by second.
        ts as TIMESTAMPADD(SECOND, a, TIMESTAMP '2024-01-31 01:00:00'),
        WATERMARK FOR ts AS ts - INTERVAL '5' SECOND 
      ) WITH (
        'connector'='datagen',
        'rows-per-second'='1',
        'fields.a.kind'='sequence','fields.a.start'='0','fields.a.end'='100000',
        'fields.b.kind'='random','fields.b.min'='0','fields.b.max'='10'
      );
      CREATE TEMPORARY TABLE t1 (
        k INT,
        ts_ltz timestamp_ltz(3),
        cnt BIGINT
      ) WITH ('connector' = 'print');
      -- Output results.
      INSERT INTO t1
      SELECT b, window_start, COUNT(*) FROM
      TABLE(
          TUMBLE(TABLE s1, DESCRIPTOR(ts), INTERVAL '5' SECOND))
      GROUP BY b, window_start, window_end;

      Após implantar e executar no console de desenvolvimento do Realtime Compute, o Watermark se alinha ao horário atual (especificamente, ao tempo simulado dos dados) — sem diferença de tempo.

      • Watermark & Low Watermark

        Nos detalhes da tarefa da Web UI do Flink, selecione um operador (ex.: GlobalWindowAggregate). O painel de informações do operador à esquerda mostra Low Watermark (ex.: 01-31 01:03:49). Mude para a aba Watermarks à direita para visualizar os valores de Watermark da SubTask e timestamps — ambos os tempos coincidem.

      • Task InputWatermark

        image

  • Motivo 2: Diferença de fuso horário entre o console de desenvolvimento do Realtime Compute e a UI do Apache Flink.

    O console de desenvolvimento do Realtime Compute exibe horários em UTC+0. A UI do Apache Flink usa o fuso horário local do navegador. Usando UTC+8 (horário de Pequim) como referência, o console de desenvolvimento do Realtime Compute mostra horários 8 horas atrás da UI do Apache Flink.

    • Console de desenvolvimento do Realtime Compute

      Na topologia da página Job O&M, os horários de Watermark são exibidos em UTC+0. Por exemplo, quando o tempo de evento é horário de Pequim 2024/1/31 09:01:34 AM, o console mostra 2024/1/31 01:01:34 AM.

      As métricas de monitoramento relacionadas a Watermark no console do produto também usam UTC+0 — 8 horas atrás do horário de Pequim.

    • UI do Apache Flink

      Na Web UI do Apache Flink, selecione um nó de operador (ex.: GlobalWindowAggregate) na topologia do job e mude para a aba Watermarks. Visualize os valores de Watermark da SubTask e os tempos de evento correspondentes. Por exemplo, Low Watermark 1706662894000 corresponde a Datetime of Watermark Timestamp 2024/1/31 09:01:34 AM. Este é o tempo de evento, não o tempo de processamento — portanto, diferenças em relação ao horário do sistema são normais.

  • Como solucionar problemas de backpressure no job?

    1. Na página Job O&M, clique em o nome do job desejado para abrir a aba Overview.

    2. Verifique Busy e BackPressure para localizar o backpressure.

      Indicadores Busy mais vermelhos significam carga de tarefa mais pesada. Indicadores BackPressure mais escuros significam impacto de backpressure mais forte.

      Por exemplo, se o operador upstream Backpressured (max) estiver em 99%, o operador do meio Busy (max) em 100% (destaque vermelho) e o operador downstream Busy (max) em apenas 7%, o gargalo é o operador do meio — otimize-o.

    3. Clique em o operador com backpressure.

    4. Na aba BackPressure, verifique o status de backpressure da SubTask.

      Se Back Pressure Status estiver verde OK, e a tabela mostrar SubTasks 0–7 com Backpressured / Idle / Busy como 0%, 0%, N/A e todos os estados OK, o job não tem backpressure.

    Como solucionar problemas de latência excessiva no job?

    Na página Job O&M, verifique a aba Monitoring and Alerts ou Data Curves para as métricas currentEmitEventTimeLag e currentFetchEventTimeLag:

    • Se currentEmitEventTimeLag estiver alto, o job tem atrasos na busca ou processamento de dados. Verifique o desempenho do operador.

    • Se currentFetchEventTimeLag estiver alto, os atrasos decorrem da busca de dados ou do processamento do sistema upstream. Investigue E/S de rede e sistemas upstream.

    Nota

    Quando fatores upstream causam alta latência, ambas as métricas aumentam simultaneamente.

    image.png

    Como otimizar um job Flink SQL quando data skew causa backpressure?

    Quando o backpressure decorre de hotspots de dados (confirmado via análise de Subtask), use estas otimizações:

    • Ative o LocalGlobal (resolve problemas comuns de hotspot de dados)

      O mecanismo LocalGlobal usa o LocalAgg para pré-agregar dados enviesados, reduzindo a pressão de hotspots no GlobalAgg e melhorando o desempenho geral.

      O LocalGlobal divide uma única agregação em duas fases: Local e Global — semelhante às fases Combine e Reduce no MapReduce. Na primeira fase, os nós upstream armazenam e agregam dados localmente (localAgg), emitindo acumuladores incrementais. Na segunda fase, esses acumuladores são mesclados (Merge) para produzir o resultado final (GlobalAgg).

      • Cenários

        Melhora o desempenho de agregações padrão (como SUM, COUNT, MAX, MIN e AVG) e resolve problemas de hotspot de dados nesses cenários.

      • Limitações

        O LocalGlobal é ativado por padrão, mas aplicam-se as seguintes limitações:

        • O MiniBatch precisa estar ativado.

        • Sua AggregateFunction deve implementar Merge.

      • Verificando o status

        Confira se a topologia gerada contém nós chamados GlobalGroupAggregate ou LocalGroupAggregate.

    • Ative o PartialFinal (resolve problemas de hotspot de COUNT DISTINCT)

      Para lidar com hotspots de COUNT DISTINCT, tradicionalmente era necessário reescrever manualmente as consultas em uma agregação de dois estágios (adicionando uma camada de shuffling baseada em módulo). O Realtime Compute for Apache Flink agora oferece shuffling automático de COUNT DISTINCT por meio da otimização PartialFinal — sem necessidade de reescrita manual.

      O LocalGlobal funciona bem para agregações padrão, mas traz pouco benefício para COUNT DISTINCT. Durante a agregação local, as taxas de deduplicação de chaves distintas permanecem baixas, então os hotspots persistem no nó global.

      • Cenários

        Use quando o COUNT DISTINCT não atender aos requisitos de desempenho do nó de agregação.

        Importante
        • Não use a otimização PartialFinal em Flink SQL que inclua UDAFs.

        • Evite o PartialFinal quando o volume de dados for pequeno — ele introduz shuffling de rede desnecessário e desperdiça recursos.

      • Como ativar

        Por padrão, este recurso vem desativado. Para ativá-lo, na aba Deployment Details do job desejado, insira o seguinte código na seção Other Configurations da área Runtime Parameter Settings.

        table.optimizer.distinct-agg.split.enabled: true
      • Verifique a eficácia

        Confira se a topologia gerada mudou de um único estágio de agregação para dois estágios.

    Como solucionar velocidade instável de consumo de dados de entrada?

    Possíveis causas e soluções:

    • O padrão de produção de dados upstream não corresponde à velocidade de processamento atual.

      Analise os padrões de geração de dados upstream para alinhar as taxas de produção e processamento.

    • O job está sofrendo backpressure.

      Verifique se há backpressure afetando o consumo upstream. Se seu job mostrar apenas um nó, adicione pipeline.operator-chaining: 'false', reinicie o job para dividir a cadeia de operadores e identifique quaisquer nós com backpressure que estejam impactando a taxa de consumo.

    • Taxa de E/S anormal.

      Revise as curvas de taxa de entrada e consumo de dados do Flink no momento relevante para determinar se a E/S é a causa.

    • Taxa de consumo anormal.

      Verifique se as flutuações na taxa de consumo coincidem com eventos de Garbage Collection (GC). Se coincidirem, inspecione o uso de memória do nó TM.

      image.png