Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Data validity FAQ

Última atualização: Jun 27, 2026

Este tópico responde às perguntas mais comuns sobre a validade de dados no Realtime Compute for Apache Flink.

Por que não há saída na tabela sink?

Se nenhum dado aparecer na tabela sink após o início do job, siga as verificações abaixo em ordem.

作业排错流程图

  1. Verifique a ocorrência de failovers. Se o job sofreu failover, analise a mensagem de erro para identificar a causa raiz e resolva o problema, garantindo a execução correta do job.

  2. Confirme se os dados estão chegando ao Realtime Compute for Apache Flink. Caso não tenha ocorrido failover, mas a latência dos dados esteja muito alta, verifique a métrica numRecordsInOfSource na página de monitoramento e alertas. Se essa métrica apresentar valor zero para um source, significa que a tabela source não está enviando dados ao Flink — investigue a fonte de dados upstream.

  3. Analise se algum operador está filtrando todos os registros. Adicione pipeline.operator-chaining: 'false' no campo Other Configuration (consulte Como configurar parâmetros de execução personalizados para um job?). Essa configuração separa a cadeia de operadores e permite inspecionar individualmente as métricas Bytes Received e Bytes Sent de cada um. O operador com entrada de dados, mas sem saída, é o responsável pelo problema — geralmente são JOIN, WINDOW ou WHERE.

  4. Avalie se o banco de dados downstream está retendo dados no buffer de gravação. Reduza o tamanho do lote do conector downstream para liberar os dados mais rapidamente.

    Importante

    Evite definir um tamanho de lote excessivamente pequeno. Um tamanho de lote igual a 1 faz com que o Flink envie uma requisição separada para cada registro processado, o que pode sobrecarregar o banco de dados downstream em cenários de alto volume de dados.

  5. Investigue possíveis deadlocks no ApsaraDB RDS for MySQL. Consulte Deadlock ao gravar no MySQL por meio do conector ApsaraDB RDS ou TDDL.

Para isolar o problema, imprima resultados intermediários nos logs usando uma tabela sink do tipo print. Consulte Visualizar saída do conector print .

Saída vazia causada pela ativação do MiniBatch com consumo parcial de binlog no CDC

  • Sintoma

    Um job CDC inicia o consumo de binlog a partir do meio, utilizando latest ou um offset específico, com o MiniBatch ativado. Consequentemente, a tabela sink recebe apenas dados parciais ou nenhuma saída.

  • Causa

    O MiniBatch mescla e cancela mensagens de changelog para a mesma chave primária dentro de um lote. Quando o job começa no meio do binlog, ele pode receber um UPDATE_AFTER sem o UPDATE_BEFORE correspondente (ou vice-versa). As mensagens opostas se cancelam dentro do lote e não geram saída, resultando em dados incorretos no downstream.

  • Soluções

    1. Consuma o binlog desde o início. Evite modos de consumo parcial, como latest ou um offset especificado, para preservar a sequência completa de alterações de cada registro.

    2. Desative o MiniBatch para que o job processe os registros um a um. Essa opção troca parte do throughput pela garantia de correção dos dados.

Como solucionar problemas de leitura de source no Flink?

Caso o Realtime Compute for Apache Flink não consiga ler de um source, realize as verificações a seguir.

Conectividade de rede

Por padrão, o Realtime Compute for Apache Flink acessa apenas serviços na mesma região e Virtual Private Cloud (VPC). Para acesso entre redes diferentes:

Lista de permissões do serviço upstream

Para ler de serviços como Kafka e Elasticsearch, adicione seu workspace do Flink às listas de permissões deles:

  1. Obtenha o bloco CIDR do vSwitch do seu workspace do Flink. Consulte Como configurar uma lista de permissões?

  2. Adicione esse bloco CIDR à lista de permissões do serviço upstream. Consulte a seção "Prerequisites" na documentação do conector, por exemplo Kafka.

Consistência de campos entre a tabela Flink e a tabela física

Incompatibilidades nas definições de campos são uma causa frequente de falhas de leitura. Ao escrever a DDL da sua tabela source do Flink:

  • Ordem dos campos: Respeite exatamente a ordem dos campos da tabela física.

  • Capitalização dos nomes: Utilize a mesma capitalização (maiúsculas e minúsculas) definida na tabela física.

  • Tipo de campo: Aplique os tipos equivalentes mapeados. Verifique a seção "Data type mappings" na documentação do conector relevante, como em Simple Log Service.

Exceções nos logs do TaskManager

Verifique se o log do TaskManager da tabela source contém mensagens de exceção:

  1. No painel de navegação à esquerda, acesse O&M > Deployment.

  2. Clique em no nome do deployment.

  3. Clique em na aba Status e selecione o vértice source no DAG.

  4. No painel direito, clique em na aba SubTasks.

  5. Na coluna More, clique em no ícone image e escolha Open TaskManager Log Page. TM日志

  6. Na aba Logs, localize a entrada mais antiga contendo "Caused by" — ela geralmente indica a causa raiz do problema.

Como investigar a ausência de saída no sistema downstream?

Siga as etapas de verificação abaixo.

Conectividade de rede

Por padrão, o Realtime Compute for Apache Flink alcança apenas serviços na mesma região e VPC. Para acesso entre redes distintas:

Lista de permissões do sistema downstream

Para gravar em serviços como ApsaraDB RDS for MySQL, Kafka, Elasticsearch, AnalyticDB for MySQL 3.0, Apache HBase, Redis e ClickHouse, adicione seu workspace do Flink às respectivas listas de permissões:

  1. Obtenha o bloco CIDR do vSwitch do seu workspace do Flink. Consulte Como configurar uma lista de permissões?

  2. Adicione esse bloco CIDR à lista de permissões do serviço downstream. Consulte a seção "Prerequisites" na documentação do conector, por exemplo ApsaraDB RDS for MySQL.

Consistência de campos entre a tabela Flink e a tabela física

Realize as mesmas verificações descritas em Como solucionar problemas de leitura de source no Flink?: valide a ordem dos campos, a capitalização dos nomes e os mapeamentos de tipos de dados.

Dados filtrados por operadores

Examine as contagens de entrada e saída de cada vértice no DAG do job. Se um vértice como WHERE apresentar entrada = 5 e saída = 0, esse operador está descartando todos os registros.

Limiares de buffer do conector sink muito altos

Quando o volume de entrada é baixo, limiares de buffer padrão elevados podem impedir que os dados sejam liberados para o sistema downstream — o buffer nunca atinge o nível necessário para disparar uma gravação. Reduza as opções relevantes conforme necessário:

Opção

Descrição

Serviço downstream relevante

batchSize

Volume de dados gravados por vez

DataHub, Tablestore, MongoDB, ApsaraDB RDS for MySQL, AnalyticDB for MySQL V3.0, ApsaraDB for ClickHouse, TSDB for InfluxDB

batchCount

Número máximo de registros gravados por vez

DataHub

flushIntervalMs

Intervalo de liberação para o buffer do MaxCompute Tunnel Writer

MaxCompute

sink.buffer-flush.max-size

Tamanho dos dados armazenados em buffer na memória antes da gravação no HBase, em bytes

ApsaraDB for HBase

sink.buffer-flush.max-rows

Quantidade de registros armazenados em buffer na memória antes da gravação no HBase

ApsaraDB for HBase

sink.buffer-flush.interval

Intervalo no qual os dados em buffer são periodicamente liberados para o HBase

ApsaraDB for HBase

jdbcWriteBatchSize

Máximo de linhas processadas por vez por um nó sink de streaming do Hologres ao usar driver JDBC

Hologres

Dados fora de ordem em janelas de tempo de evento

As watermarks determinam quais registros uma janela aceita. Se o primeiro registro tiver um timestamp de 2100 e definir a watermark como 2100, qualquer registro subsequente com timestamp inferior a 2100 (como 2021) será considerado atrasado e descartado. A janela não poderá ser fechada até que chegue um registro com timestamp superior a 2100.

Para detectar registros fora de ordem, utilize uma tabela sink do tipo print ou examine os logs do Log4j. Consulte Criar uma tabela sink print e Configurar saída de log. Caso confirme a presença de registros atrasados, filtre-os ou configure sua estratégia de watermark para permitir um período de tolerância para chegadas tardias.

Subtarefas de source sem entrada

Quando uma subtarefa de source não recebe dados, sua watermark permanece no padrão epoch (1970-01-01T00:00:00Z), que passa a ser a watermark geral do operador. Isso impede que as janelas de tempo de evento sejam fechadas.

Verifique o DAG do job e confirme se todas as subtarefas de source estão recebendo entrada. Se alguma subtarefa estiver ociosa, reduza o paralelismo do job para corresponder à contagem de shards da tabela upstream, garantindo que cada subtarefa receba dados.

Partições vazias do Kafka

Uma partição vazia do Kafka pode interromper a geração de watermarks. Consulte Por que uma janela de tempo de evento não gera saída a partir de uma tabela source do Kafka?

Como resolver perdas de dados?

Reduções no volume de dados geralmente resultam de cláusulas WHERE, JOINs ou operações com janelas. Para perdas inexplicáveis, verifique os itens a seguir.

Política de cache de tabela de dimensão

Uma política de cache incorreta pode causar falhas silenciosas em lookup joins, descartando registros. Configure uma política adequada usando as opções relacionadas a cache na documentação do conector, como a seção "Specific to dimension tables (such as Cache parameters)" em ApsaraDB for HBase.

Uso de funções

O uso incorreto de funções como to_timestamp_tz e date_format pode provocar falhas na conversão de dados, descartando registros silenciosamente. Valide o comportamento das funções usando uma tabela sink print ou logs do Log4j. Consulte Print e Configurar saída de log.

Dados fora de ordem

Eventos atrasados são descartados quando seu timestamp fica fora do intervalo aceito pela janela atual. Por exemplo, um evento com timestamp 11s entrando em uma janela de 15–20s é descartado porque sua watermark é 11 — abaixo do limite inferior da janela.

乱序

Perdas por esse motivo costumam se concentrar em uma única janela. Use uma tabela sink print ou o Log4j para confirmar a existência de dados fora de ordem.

Para minimizar perdas por dados fora de ordem, defina uma estratégia de geração de watermark com um período de tolerância (por exemplo, Watermark = Event time - 5s). Alinhe as janelas a limites exatos de dia, hora ou minuto — isso torna o comportamento da janela previsível e reduz descartes em casos extremos quando combinado com um período de tolerância adequado.

Por que obtenho resultados imprecisos ao usar ROW_NUMBER para deduplicar dados ingeridos do Hologres no modo CDC?

image.png

O downstream utiliza um operador de retração (por exemplo, ROW_NUMBER OVER WINDOW para deduplicação), mas o source do Hologres não está configurado para emitir dados no modo upsert. Sem o modo upsert, o source emite apenas eventos de inserção, que o operador de retração não consegue processar corretamente.

Adicione 'upsertSource' = 'true' à cláusula WITH da instrução DDL da tabela source.

image.png

Como corrigir resultados imprecisos?

  1. Altere o nível de log para INFO.

  2. Ative o profiling de operadores para inspecionar resultados intermediários sem modificar a lógica do job.

  3. Analise os logs de execução: image

    1. Clique em no nome do deployment e, em seguida, na aba Status.

    2. No DAG, copie o nome do operador que está produzindo resultados incorretos.

    3. Na lista de logs, clique em em inspect-taskmanager_0.out sob Log Name e pesquise pelo nome do operador.

  4. Após identificar a causa raiz, revise a lógica do operador, reinicie o job e verifique a precisão dos dados.

Como corrigir o erro "doesn't support consuming update and delete changes which is produced by node TableSourceScan"?

A mensagem de erro tem o seguinte formato:

Table sink 'vvp.default.***' doesn't support consuming update and delete changes which is produced by node TableSourceScan(table=[[vvp, default, ***]], fields=[id,b, content])
    at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.wrapExecutor(DelegateOperationExecutor.java:286)
    at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.validate(DelegateOperationExecutor.java:211)
    at org.apache.flink.table.sqlserver.FlinkSqlServiceImpl.validate(FlinkSqlServiceImpl.java:741)
    at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$MethodHandlers.invoke(FlinkSqlServiceGrpc.java:2522)
    at io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:172)
    at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:331)
    at io.grpc.internal.ServerImpl$JumpToApplicationThreadServerStreamListener$1HalfClosed.runInContext(ServerImpl.java:820)
    at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37)
    at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:123)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1147)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:622)
    at java.lang.Thread.run(Thread.java:834)

A tabela sink está no modo append-only e não consegue consumir eventos de atualização ou exclusão provenientes do source. Substitua-a por um sink que suporte upserts, como Upsert Kafka.

Como evitar substituições ou exclusões inesperadas de dados ao usar o conector Lindorm?

Por padrão, o conector Lindorm usa o operador upsert materialize (padrão: AUTO) para gerenciar a ordem de gravação. Esse operador gera um DELETE seguido de um INSERT para a mesma chave primária. Duas características do Lindorm tornam isso problemático:

  • Precisão de timestamp em milissegundos: O Lindorm versiona dados usando timestamps em milissegundos. Múltiplos registros com a mesma chave primária gravados dentro de um único milissegundo podem chegar fora de ordem, causando conflitos de versão.

  • Ausência de suporte nativo a DELETE: O Lindorm suporta apenas semântica UPSERT — as exclusões são irreversíveis. Portanto, a lógica de manutenção de ordem do upsert materialize é ineficaz e pode causar anomalias de dados devido à sequência DELETE + INSERT.

Quando gravações concorrentes ocorrem dentro do mesmo milissegundo, as operações DELETE e INSERT resultantes podem produzir dados incorretos ou perda silenciosa de dados.

Solução: Desative explicitamente o operador upsert materialize adicionando o seguinte à configuração de parâmetros de execução do seu job ou ao código SQL:

SET 'table.exec.sink.upsert-materialize' = 'NONE';

Essa configuração aplica-se a qualquer job que grave no Lindorm por meio do Flink.

Após desativar esse operador, apenas a consistência eventual é garantida. Confirme se a consistência eventual é aceitável para o seu caso de uso antes de aplicar essa alteração.