Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conectores

Última atualização: Jul 08, 2026

Problemas comuns e soluções para conectores no Realtime Compute for Apache Flink.

Recuperar dados JSON do Kafka usando o Flink

  • Para recuperar dados JSON padrão, consulte JSON Format.

  • Para recuperar dados JSON aninhados, defina o objeto JSON como um tipo ROW na DDL da tabela de origem. Na DDL da tabela de destino, declare as chaves a serem recuperadas. Em seguida, use uma instrução DML para acessar as chaves e extrair seus valores. O código a seguir fornece um exemplo:

    • Dados de amostra

      {
          "a":"abc",
          "b":1,
          "c":{
              "e":["1","2","3","4"],
              "f":{"m":"567"}
          }
      }
    • DDL da tabela de origem

      CREATE TEMPORARY TABLE `kafka_table` (
        `a` VARCHAR,
         b int,
        `c` ROW<e ARRAY<VARCHAR>,f ROW<m VARCHAR>>  -- 'c' is a JSON object that maps to the ROW type in Flink. 'e' is a JSON array that maps to the ARRAY type.
      ) WITH (
        'connector' = 'kafka',
        'topic' = 'xxx',
        'properties.bootstrap.servers' = 'xxx',
        'properties.group.id' = 'xxx',
        'format' = 'json',
        'scan.startup.mode' = 'xxx'
      );
    • DDL da tabela de destino

      CREATE TEMPORARY TABLE `sink` (
       `a` VARCHAR,
        b INT,
        e VARCHAR,
        `m` varchar
      ) WITH (
        'connector' = 'print',
        'logger' = 'true'
      );
    • Instrução DML

      INSERT INTO `sink`
        SELECT 
        `a`,
        b,
        c.e[1], -- Flink uses 1-based indexing for arrays. This example uses index 1 to retrieve the first element. To retrieve the entire array, omit [1].
        c.f.m
      FROM `kafka_table`;
    • Resultados

      409  2021-04-08 10:13:11,214 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] -
      410  2021-04-08 10:13:11,215 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka commitId: cb8625948210849f
      411  2021-04-08 10:13:11,215 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka startTimeMs: 1617847991214
      412  2021-04-08 10:13:11,270 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Subscribed to partition(s): lb_test-0, lb_test-1, lb_test-2, lb_test-3, lb_test-4, lb_test-5
      413  2021-04-08 10:13:11,280 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 1 for partition lb_test-0
      414  2021-04-08 10:13:11,285 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-1
      415  2021-04-08 10:13:11,285 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-2
      416  2021-04-08 10:13:11,285 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-3
      417  2021-04-08 10:13:11,290 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-4
      418  2021-04-08 10:13:11,291 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Seeking to offset 0 for partition lb_test-5
      419  2021-04-08 10:13:11,302 INFO  org.apache.flink.kafka.shaded.org.apache.kafka.clients.Metadata [] - [Consumer clientId=consumer-lb_test-2, groupId=lb_test] Cluster ID: -1flJPwnTvuGFSuyCtU1hw
      420  2021-04-08 10:15:31,597 INFO  org.apache.flink.api.common.functions.util.PrintSinkOutputWriter [] - +I(abc,1,1,567)

O Flink não consegue consumir ou gravar no Kafka

  • Causa

    Se existir um mecanismo de encaminhamento, como um proxy ou mapeamento de portas, entre o Flink e o Kafka, o cliente do Kafka recupera o endereço de rede interna do servidor Kafka, e não o endereço do proxy. Como resultado, o Flink consegue se conectar ao cluster do Kafka, mas não consegue consumir ou gravar dados, mesmo com um caminho de rede estabelecido.

    O processo de conexão entre o conector do Flink para Kafka e o servidor Kafka envolve duas etapas:

    1. O cliente do Kafka busca metadados dos brokers do Kafka. Esses metadados incluem os endereços de rede de todos os brokers no cluster.

    2. O conector do Flink usa então esses endereços de rede para consumir ou gravar dados.

  • Solução de problemas

    Siga estas etapas para determinar se existe um mecanismo de encaminhamento, como um proxy ou mapeamento de portas, entre o Flink e o Kafka:

    1. Use uma ferramenta de linha de comando do ZooKeeper (zkCli.sh ou zookeeper-shell.sh) para fazer login no cluster do ZooKeeper utilizado pelo seu cluster do Kafka.

    2. Execute o comando apropriado para o seu cluster para recuperar os metadados do broker do Kafka.

      Geralmente, você pode usar o comando get /brokers/ids/0 para recuperar metadados do broker do Kafka. O endereço de conexão está localizado no campo endpoints. Por exemplo, conecte-se usando o ZooKeeper Shell e execute get /brokers/ids/0 para visualizar as informações de registro do broker. Observe o endereço configurado no campo endpoints do JSON retornado:

      # bin/zookeeper-shell.sh localhost:2181
      Connecting to localhost:2181
      Welcome to ZooKeeper!
      JLine support is disabled
      WATCHER::
      WatchedEvent state:SyncConnected type:None path:null
      get /brokers/ids/0
      {"listener_security_protocol_map":{"PLAINTEXT":"PLAINTEXT"},"endpoints":["PLAINTEXT://116.62.xxx:9092"],"jmx_port":-1,"host":"116.62.xxx","timestamp":"1614840078030","port":9092,"version":4}
    3. Use um comando como ping ou telnet para testar a conectividade do ambiente do Flink até o endereço obtido no campo endpoints.

      Uma falha na conexão indica que existe um mecanismo de encaminhamento, como um proxy ou mapeamento de portas, entre o Flink e o Kafka.

  • Soluções

    • Não utilize mecanismos de encaminhamento. Estabeleça um caminho de rede direto entre o Flink e o Kafka. Isso permite que o Flink se conecte diretamente aos endpoints listados nos metadados do Kafka.

    • Entre em contato com o administrador do Kafka para configurar o endereço de encaminhamento na propriedade advertised.listeners nos brokers do Kafka. Isso garante que o cliente do Kafka recupere metadados que incluam o endereço de encaminhamento correto.

      Nota

      Apenas versões do Kafka 0.10.2.0 e posteriores suportam a adição de um endereço de proxy aos listeners de um broker do Kafka.

    Para mais informações sobre como isso funciona, consulte KIP-103: Separate Internal and External traffic e Kafka client cannot connect to brokers.

Se a conectividade de rede entre o Flink e o Kafka foi confirmada e o problema persistir, verifique as seguintes causas não relacionadas à rede:

Verificação 1: Estratégia de offset inicial

Verifique o parâmetro scan.startup.mode na cláusula WITH da DDL da sua tabela de origem do Kafka. Se o valor for latest-offset, o Flink lê apenas mensagens gravadas após o início do job. Se nenhuma nova mensagem chegar após a inicialização do job, ele parecerá não estar consumindo dados.

**Valor de scan.startup.mode**

Comportamento

earliest-offset

Lê a partir da mensagem mais antiga disponível em cada partição.

latest-offset

Lê apenas mensagens gravadas após o início do job. Dados produzidos antes da inicialização não são consumidos.

group-offsets

Retoma a partir do último offset confirmado do grupo de consumidores. Se nenhum offset tiver sido confirmado, reverte para latest-offset.

timestamp

Lê a partir de um timestamp especificado pelo usuário. Requer a definição de scan.startup.timestamp-millis.

Para verificar se novos dados estão sendo produzidos após o início do job, use um cliente consumidor do Kafka para monitorar o tópico em tempo real.

Verificação 2: Incompatibilidade de formato de dados

Confirme se o parâmetro format na cláusula WITH da sua tabela de origem do Kafka corresponde à codificação real das mensagens no tópico do Kafka. Uma incompatibilidade de formato causa falhas na desserialização, o que pode resultar no job ignorando mensagens silenciosamente ou não produzindo saída.

Cenário

**Valor de format**

Mensagens JSON simples

json

Mensagens Canal CDC

canal-json

Mensagens Debezium CDC

debezium-json

Mensagens Maxwell CDC

maxwell-json

Para inspecionar o formato real da mensagem, use um cliente consumidor do Kafka para ler os bytes brutos do tópico e examinar a estrutura do payload.

Nenhuma saída de dados de janelas de tempo de evento do Kafka

  • Problema

    Um job não produz saída quando utiliza uma tabela de origem do Kafka com uma janela de tempo de evento.

  • Causa

    Uma partição ociosa do Kafka pode impedir o avanço da marca d'água, o que impede a janela de tempo de evento de gerar saída.

  • Solução

    1. Garanta que todas as partições recebam dados.

    2. Para ativar a detecção de ociosidade da source, adicione o seguinte código à seção Other Configurations e salve suas alterações. Para instruções detalhadas, consulte Como configurar parâmetros de runtime personalizados para jobs?.

      table.exec.source.idle-timeout: 5

      Para mais informações sobre o parâmetro table.exec.source.idle-timeout, consulte Configuration.

Commit de offset no Kafka

Um offset confirmado no Kafka rastreia a posição dos dados processados, garantindo consistência e confiabilidade no processamento de streams ao evitar duplicação ou perda de dados. Quando um checkpoint é concluído com sucesso, o Flink confirma o offset de leitura correspondente no Kafka. Se o checkpointing não estiver habilitado, ou se o intervalo de checkpoint for muito longo, o offset confirmado no Kafka pode ficar desatualizado, levando ao reprocessamento ou perda de dados.

Analisar JSON aninhado com o conector do Kafka

Por exemplo, ao analisar os seguintes dados JSON diretamente com o formato json, eles são resolvidos em um único campo do tipo ARRAY<ROW<cola VARCHAR, colb VARCHAR>>. Esse campo é um array de linhas, onde cada linha contém dois campos VARCHAR. Você pode então analisar esse array usando uma função de tabela definida pelo usuário (UDTF).

{"data":[{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"},{"cola":"test1","colb":"test2"}]}

Conectar a um cluster seguro do Kafka

  1. Na cláusula WITH da DDL da sua tabela do Kafka, adicione as configurações de segurança para autenticação e criptografia. Para obter uma lista completa de opções, consulte SECURITY.

    Importante

    Adicione o prefixo properties. a todos os parâmetros de configuração de segurança.

    • Este exemplo mostra como configurar uma tabela do Kafka para usar o mecanismo SASL PLAIN e fornecer uma configuração JAAS.

      CREATE TABLE KafkaTable (
        `user_id` BIGINT,
        `item_id` BIGINT,
        `behavior` STRING,
        `ts` TIMESTAMP(3) METADATA FROM 'timestamp'
      ) WITH (
        'connector' = 'kafka',
        ...
        'properties.security.protocol' = 'SASL_PLAINTEXT',
        'properties.sasl.mechanism' = 'PLAIN',
        'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"username\" password=\"password\";'
      );
    • Este exemplo mostra como usar o protocolo de segurança SASL_SSL com o mecanismo SASL SCRAM-SHA-256.

      CREATE TABLE KafkaTable (
        `user_id` BIGINT,
        `item_id` BIGINT,
        `behavior` STRING,
        `ts` TIMESTAMP(3) METADATA FROM 'timestamp'
      ) WITH (
        'connector' = 'kafka',
        ...
        'properties.security.protocol' = 'SASL_SSL',
        /* SSL configuration */
        /* Path to the truststore (CA certificate) provided by the server */
        'properties.ssl.truststore.location' = '/flink/usrlib/kafka.client.truststore.jks',
        'properties.ssl.truststore.password' = 'test1234',
        /* If client-side authentication is required, configure the path to the keystore (private key) */
        'properties.ssl.keystore.location' = '/flink/usrlib/kafka.client.keystore.jks',
        'properties.ssl.keystore.password' = 'test1234',
        /* SASL configuration */
        /* Configure the SASL mechanism as SCRAM-SHA-256 */
        'properties.sasl.mechanism' = 'SCRAM-SHA-256',
        /* Configure JAAS */
        'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";'
      );
      Nota
      • Se properties.sasl.mechanism for SCRAM-SHA-256, use org.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModule para properties.sasl.jaas.config.

      • Se properties.sasl.mechanism for PLAIN, use org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule para properties.sasl.jaas.config.

  2. Na seção Additional dependency files do seu job, faça upload de todos os arquivos necessários, como certificados, chaves públicas e chaves privadas.

    A plataforma armazena os arquivos enviados no diretório /flink/usrlib. Para instruções de upload, consulte implantação de jobs.

    Importante

    Se o mecanismo de autenticação no seu broker do Kafka for SASL_SSL, mas o mecanismo do lado do cliente for SASL_PLAINTEXT, o job falhará com uma exceção OutOfMemory durante a validação. Para resolver esse problema, garanta que os mecanismos de autenticação do lado do cliente e do servidor correspondam.

Conflitos de nomenclatura de campos

  • Problema

    Uma fonte de dados do Kafka serializa mensagens em duas strings JSON separadas: uma para a chave e outra para o valor. Nesse cenário, tanto a chave quanto o valor contêm um campo com o mesmo nome, como o campo id no exemplo abaixo. Analisar esses dados diretamente em uma tabela do Flink causa um conflito de nomenclatura de campos.

    • chave

      {
         "id": 1
      }
    • valor

      {
         "id": 100,
         "name": "flink"
      }
  • Solução

    Utilize a propriedade key.fields-prefix para evitar esse problema.

    CREATE TABLE kafka_table (
      -- Define the columns for the key and value fields
      key_id INT,
      value_id INT,
      name STRING
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'test_topic',
      'properties.bootstrap.servers' = 'localhost:9092',
      'format' = 'json',
      'json.ignore-parse-errors' = 'true',
      -- Specify the fields and data types for the key
      'key.format' = 'json',
      'key.fields' = 'id',
      'value.format' = 'json',
      'value.fields' = 'id, name',
      -- Add a prefix to fields from the key
      'key.fields-prefix' = 'key_'
    );

    Definir a propriedade key.fields-prefix como key_ instrui o conector a adicionar o prefixo key_ a todos os campos da chave da mensagem. Por exemplo, o campo id da chave torna-se a coluna key_id na tabela do Flink. Isso evita conflitos com o campo id do valor, que é mapeado para a coluna value_id.

    Executar a consulta SELECT * FROM kafka_table; produz a seguinte saída:

    key_id: 1,
    value_id: 100,
    name: flink

Solucionar alta latência de uma source do Kafka

  • Problema

    Ao ler de uma tabela de origem do Kafka, a métrica currentEmitEventTimeLag apresenta um valor superior a 50 anos. Por exemplo, vários jobs Flink SQL estão no estado running, mas a coluna business latency exibe um valor anormalmente alto que excede 19.160 dias, como 19160d 1h 59m 28s.

  • Solução de problemas

    1. Primeiro, determine se o job é um job JAR ou um job SQL.

      Para um job JAR, verifique se seu arquivo pom.xml usa a dependência do Kafka fornecida pelo Realtime Compute for Apache Flink. A versão open-source do conector não reporta essas métricas.

    2. Verifique se todas as partições no tópico upstream do Kafka estão recebendo dados em tempo real.

    3. Verifique se o timestamp nos metadados da mensagem do Kafka é 0 ou null.

      A latência de uma source do Kafka é calculada subtraindo o timestamp da mensagem da hora atual. Se uma mensagem não tiver timestamp, a latência pode ser exibida como superior a 50 anos. Você pode verificar o timestamp de uma das seguintes maneiras:

      • Para jobs SQL, recupere o timestamp da mensagem definindo uma coluna de metadados. Para mais informações, consulte Tabela de origem do Kafka.

        CREATE TEMPORARY TABLE sk_flink_src_user_praise_rt (
            `timestamp` BIGINT ,
            `timestamp` TIMESTAMP METADATA,  --Metadata timestamp.
            ts as to_timestamp (
              from_unixtime (`timestamp`, 'yyyy-MM-dd HH:mm:ss')
            ),
            watermark for ts as ts - interval '5' second
          ) WITH (
            'connector' = 'kafka',
            'topic' = '',
            'properties.bootstrap.servers' = '',
            'properties.group.id' = '',
            'format' = 'json',
            'scan.startup.mode' = 'latest-offset',
            'json.fail-on-missing-field' = 'false',
            'json.ignore-parse-errors' = 'true'
          );
      • Escreva um programa Java simples que use o cliente KafkaConsumer para ler uma mensagem e inspecionar seu timestamp.

Erro: tabelas 'upsert-kafka' exigem uma PRIMARY KEY

  • Problema

    ) WITH (
        'connector' = 'upsert-kafka',
        'topic' = 'flow_stay_duration',
        'properties.bootstrap.servers' = 'xxx',
        'key.format' = 'avro',
        'value.format' = 'avro'
    );
        insert into sink_ad_data_device_info
    org.apache.flink.table.api.ValidationException: SQL validation failed. Unable to create a sink for writing table 'vvp.default.sink_ad_data_device_info'.
    The cause is following: 'upsert-kafka' tables require to define a PRIMARY KEY constraint. The PRIMARY KEY specifies which columns should be read from or write to the Kafka message key. The PRIMARY KEY also defines records in the 'upsert-kafka' table should update or delete on which keys.
    Table options are:
    'connector'='upsert-kafka'
    'key.format'='avro'
    'properties.bootstrap.servers'='xxx'
    'topic'='flow_stay_duration'
    'value.format'='avro'
        at org.apache.flink.table.sqlserver.utils.FormatValidatorExceptionUtils.newValidationException(FormatValidatorExceptionUtils.java:41)
        at org.apache.flink.table.sqlserver.utils.ErrorConverter.formatException(ErrorConverter.java:123)
        at org.apache.flink.table.sqlserver.utils.ErrorConverter.toErrorDetail(ErrorConverter.java:60)
        at org.apache.flink.table.sqlserver.utils.ErrorConverter.toGrpcException(ErrorConverter.java:54)
        at org.apache.flink.table.sqlserver.FlinkSqlServiceImpl.validateAndGeneratePlan(FlinkSqlServiceImpl.java:979)
        at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$MethodHandlers.invoke(FlinkSqlServiceGrpc.java:3283)
        at io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:172)
        at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:331)
  • Causa

    Esse erro ocorre porque a DDL não possui uma chave primária. Quando usado como tabela de destino, o conector upsert-kafka consome o stream de changelog da lógica upstream. O conector grava dados INSERT e UPDATE_AFTER no Kafka. Para operações DELETE, ele grava uma mensagem com valor nulo para indicar que a mensagem correspondente àquela chave foi excluída. O Flink usa as colunas da chave primária para particionar os dados. Isso garante que mensagens com a mesma chave sejam ordenadas e que suas respectivas mensagens de atualização ou exclusão cheguem à mesma partição.

  • Solução

    Defina uma chave primária na DDL.

Recuperar um job do Flink após divisão ou redução de tópico

Se você dividir ou reduzir um tópico do DataHub que está sendo lido por um job do Flink, o job entrará em um loop de falhas e não conseguirá se recuperar automaticamente. Para resolver isso, reinicie o job.

Exclusão de tópicos com consumidores ativos

Não é possível excluir ou recriar um tópico do DataHub que possua consumidores ativos.

Os parâmetros endPoint e tunnelEndpoint

Os parâmetros endPoint e tunnelEndpoint estão descritos em endpoint. Em um ambiente de VPC, a configuração incorreta desses parâmetros pode causar exceções na tarefa:

  • Se o parâmetro endPoint estiver configurado incorretamente, a implantação da tarefa ficará travada em 91% de progresso.

  • Caso o parâmetro tunnelEndpoint esteja configurado de forma errada, a execução da tarefa falhará.

Falha na criação de tabela DataHub com NoPermissionException: dhs:ListShard

  • Sintoma

    Quando um job do Flink implanta uma tabela de origem ou destino do DataHub, ele falha com um erro semelhante ao seguinte:

    NoPermissionException: You have no permission to perform this action. Action: dhs:ListShard
  • Causa

    Esse erro ocorre devido à nomenclatura não padrão dos parâmetros WITH na DDL do DataHub. O conector do DataHub exige que as credenciais sejam especificadas como accessId e accessKey. Se você utilizar as formas com ponto access.id e access.key, o conector não reconhecerá os campos de credencial e não conseguirá autenticar. Consequentemente, o conector tentará listar shards sem credenciais válidas, e o DataHub retornará uma NoPermissionException.

    A mensagem de erro menciona a ausência do privilégio dhs:ListShard, mas a causa raiz são os parâmetros de credencial não reconhecidos — e não uma deficiência real de permissão IAM.

  • Solução

    Na sua DDL do DataHub, renomeie access.id para accessId e access.key para accessKey. Exclua a tabela existente e recrie-a com a cláusula WITH corrigida.

    Configuração incorreta:

    CREATE TABLE datahub_source (...) WITH (
      'connector'       = 'datahub',
      'endPoint'        = 'https://dh-cn-hangzhou.aliyuncs.com',
      'project'         = 'your_project',
      'topic'           = 'your_topic',
      'access.id'       = 'your-access-key-id',      -- Incorrect: dotted form not recognized
      'access.key'      = 'your-access-key-secret'   -- Incorrect: dotted form not recognized
    );

    Configuração correta:

    CREATE TABLE datahub_source (...) WITH (
      'connector'  = 'datahub',
      'endPoint'   = 'https://dh-cn-hangzhou.aliyuncs.com',
      'project'    = 'your_project',
      'topic'      = 'your_topic',
      'accessId'   = 'your-access-key-id',      -- Correct
      'accessKey'  = 'your-access-key-secret'   -- Correct
    );

    Para obter a lista completa de parâmetros WITH suportados, consulte a documentação do conector do DataHub.

Leituras completas e incrementais de origens MaxCompute

As origens MaxCompute realizam leituras completas e incrementais por meio do MaxCompute Tunnel. O throughput de leitura é limitado pela largura de banda do MaxCompute Tunnel.

As tabelas de origem MaxCompute podem ler dados anexados?

Não. Após o início de um job do Flink, ele não lerá novos dados anexados a uma tabela ou partição de origem. Isso se aplica tanto quando a origem está lendo ativamente quanto quando já terminou a leitura. Anexar dados dessa maneira também pode causar um failover do job.

Tanto as tabelas de origem MaxCompute completas quanto as incrementais usam ODPS DOWNLOAD SESSION para ler dados de tabelas ou partições. Ao criar uma DOWNLOAD SESSION, o servidor cria um arquivo de índice. Esse arquivo é um snapshot dos dados no momento em que a DOWNLOAD SESSION é criada, e as leituras subsequentes baseiam-se nesse snapshot. Portanto, após a criação de uma DOWNLOAD SESSION, os dados anexados à tabela ou partição MaxCompute não são lidos em circunstâncias normais. No entanto, se novos dados forem gravados na tabela de origem MaxCompute, duas exceções poderão ocorrer:

  • Falha durante a leitura: Se novos dados forem gravados enquanto o Tunnel estiver lendo ativamente, a operação falhará com o erro ErrorCode=TableModified,ErrorMessage=The specified table has been modified since the download initiated..

  • Dados inconsistentes no failover: Caso novos dados sejam gravados após o fechamento do Tunnel, a execução atual do job não os lerá. Porém, se o job sofrer failover ou for retomado de um estado pausado, ele poderá reprocessar dados antigos e ler apenas parte dos novos dados.

Alterações de concorrência para jobs MaxCompute pausados

Para tabelas de origem MaxCompute com a opção useNewApi ativada (ativada por padrão), os jobs em modo streaming suportam alterações de concorrência após serem pausados e retomados. Uma tabela de origem MaxCompute lê as partições correspondentes sequencialmente. Ao ler uma partição, ela distribui os dados dessa partição entre os operadores paralelos. A alteração da concorrência não afeta a distribuição de dados da partição que estava sendo processada antes da pausa. O novo paralelismo entra em vigor somente quando o job começa a processar a próxima partição. Como resultado, se um job estiver processando uma única partição grande, aumentar a concorrência e retomar o job pode fazer com que apenas alguns dos operadores MaxCompute leiam dados.

Alterações de concorrência não são suportadas para jobs em lote ou para jobs nos quais a opção useNewApi esteja definida como false.

Por que o MaxCompute lê partições anteriores quando a posição inicial é 2019-10-11 00:00:00?

A configuração da posição inicial afeta apenas fontes de dados de fila de mensagens, como o DataHub. Essa configuração não afeta uma tabela de origem MaxCompute. Quando um job do Flink é iniciado, ele lê os dados da seguinte forma:

  • Para uma tabela particionada: Todas as partições existentes são lidas.

  • Para uma tabela não particionada: Todos os dados existentes são lidos.

Evitar leituras incompletas de dados de novas partições

Atualmente, não existe um mecanismo para verificar se os dados de uma partição estão completos. Como resultado, uma tabela de origem MaxCompute incremental começa a ler uma nova partição imediatamente após detectá-la. Suponha que você esteja usando uma tabela de origem MaxCompute incremental para ler uma tabela MaxCompute particionada T, onde ds é a coluna de partição. Nesse caso, recomendamos que você não crie a partição primeiro. Em vez disso, execute uma instrução INSERT OVERWRITE TABLE T PARTITION (ds='20191010') .... Quando o job for concluído, a partição e seus dados aparecerão simultaneamente.

Importante

Não crie a partição primeiro (por exemplo, ds=20191010) e depois grave dados nela. Se você usar esse método, a tabela de origem MaxCompute incremental detectará a nova partição ds=20191010 e começará a ler dela imediatamente. Isso leva a leituras de dados incompletas se a operação de gravação ainda estiver em andamento.

Erro de autorização do conector MaxCompute

  • Detalhes do erro

    Durante a execução do job, um erro é exibido na página de failover ou no arquivo TaskManager.log:

    ErrorMessage=Authorization Failed [4019], You have NO privilege'ODPS:***'
  • Causa

    As informações de identidade do usuário especificadas na definição DDL do MaxCompute não possuem as permissões necessárias para acessar o MaxCompute.

  • Solução

    Autentique-se usando uma conta Alibaba Cloud, um usuário RAM ou uma função RAM. Para mais informações, consulte autenticação de usuário.

Configurar o parâmetro startPartition

Etapa

Descrição

Exemplo

1

Conecte cada nome de coluna de partição ao seu valor fixo correspondente com um sinal de igual (=).

Se a coluna de partição for dt e você quiser ler dados a partir do valor de partição 20220901, o resultado será dt=20220901.

2

Ordene os resultados da Etapa 1 por nível de partição em ordem crescente e una-os com vírgulas (,) sem espaços. Essa string é o valor do parâmetro startPartition.

Nota

É possível especificar apenas os primeiros níveis de partição.

  • Para uma única partição de primeiro nível dt, para começar a ler a partir de dt=20220901, defina o parâmetro como 'startPartition' = 'dt=20220901'.

  • Para três níveis de partição (dt, hh e mm), para começar a ler a partir de dt=20220901,hh=08,mm=10, defina o parâmetro como 'startPartition' = 'dt=20220901,hh=08,mm=10'.

  • Para três níveis de partição (dt, hh e mm), para começar a ler a partir de dt=20220901,hh=08, defina o parâmetro como 'startPartition' = 'dt=20220901,hh=08'.

Ao carregar a lista de partições, o sistema compara cada partição com o valor startPartition em ordem lexicográfica. Em seguida, ele carrega todas as partições que são maiores ou iguais ao valor startPartition. Por exemplo, considere uma tabela particionada MaxCompute para leituras incrementais com uma partição de primeiro nível ds e uma partição de segundo nível type. A tabela contém as seis partições a seguir:

  • ds=20191201,type=a

  • ds=20191201,type=b

  • ds=20191202,type=a

  • ds=20191202,type=b

  • ds=20191202,type=c

  • ds=20191203,type=a

Se startPartition estiver definido como ds=20191202, o sistema lerá quatro partições: ds=20191202,type=a, ds=20191202,type=b, ds=20191202,type=c e ds=20191203,type=a. Se startPartition estiver definido como ds=20191202,type=b, o sistema lerá três partições: ds=20191202,type=b, ds=20191202,type=c e ds=20191203,type=a.

Nota

A partição especificada em startPartition não precisa existir. O sistema lê todas as partições que são lexicograficamente maiores ou iguais ao valor startPartition.

Início lento para jobs MaxCompute incrementais

O job inicia lentamente porque precisa primeiro processar metadados de todas as partições que são lexicograficamente maiores ou iguais ao valor startPartition. Esse processo sofre atrasos significativos quando há um grande número de partições ou arquivos pequenos. Para mitigar esse atraso, siga estas recomendações:

  • Evite ler muitos dados históricos.

    Nota

    Se precisar processar dados históricos, execute um job em lote com uma tabela de origem MaxCompute.

  • Reduza o número de arquivos pequenos nos seus dados históricos .

Definir o parâmetro partition

Ler de partições

  • Leitura de partições estáticas

    Ao ler de partições estáticas de uma tabela de origem ou tabela de dimensão, defina o parâmetro partition conforme descrito abaixo.

    Etapa

    Descrição

    Exemplo

    1

    • Para uma tabela de dimensão, especifique cada partição como 'partition_column_name=partition_value'. O valor da partição deve ser um valor fixo.

    • Para uma tabela de origem, especifique cada partição como 'partition_column_name=partition_value'. O valor da partição pode ser um valor fixo ou um valor que contenha um curinga (*). O curinga pode corresponder a qualquer string, incluindo uma string vazia.

    • Para ler dados da coluna de partição dt com o valor 20220901, especifique dt=20220901.

    • Para ler dados de partições na coluna dt cujos valores começam com 202209, especifique dt=202209* (aplica-se apenas a tabelas de origem).

    • Para ler dados de partições na coluna dt cujos valores começam com 2022 e terminam com 01, especifique dt=2022*01 (aplica-se apenas a tabelas de origem).

    • Para ler dados de todas as partições na coluna dt, especifique dt=* (aplica-se apenas a tabelas de origem).

    2

    Ordene as strings de partição da Etapa 1 por nível de partição em ordem crescente e, em seguida, una-as com vírgulas (sem espaços). A string resultante é o valor do parâmetro partition.

    É possível especificar apenas os primeiros níveis de partição.

    • Uma tabela possui uma única partição de primeiro nível dt. Para ler dados da partição dt=20220901, especifique 'partition' = 'dt=20220901'.

    • Uma tabela possui três níveis de partição: uma partição de primeiro nível dt, uma partição de segundo nível hh e uma partição de terceiro nível mm. Para ler dados de dt=20220901, hh=08 e mm=10, especifique 'partition' = 'dt=20220901,hh=08,mm=10'.

    • Para a mesma tabela, para ler dados de dt=20220901, hh=08 e qualquer valor para mm, especifique 'partition' = 'dt=20220901,hh=08' or 'partition' = 'dt=20220901,hh=08,mm=*'.

    • Para a mesma tabela, para ler dados de dt=20220901, qualquer valor para hh e mm=10, especifique 'partition' = 'dt=20220901,hh=*,mm=10'.

    Se essas etapas não atenderem às suas necessidades de filtragem de partições, adicione as condições de filtro à cláusula WHERE da sua instrução SQL. Isso permite que o otimizador SQL use o pushdown de partição para filtragem. Por exemplo, para ler partições de uma tabela com dois níveis de partição (dt e hh) onde dt está entre '20220901' e '20220903', e hh está entre '09' e '17', use uma instrução SQL como a seguinte.

    CREATE TABLE maxcompute_table (
      content VARCHAR,
      dt VARCHAR,
      hh VARCHAR
    ) PARTITIONED BY (dt, hh) WITH ( 
       -- You must specify the partition columns with PARTITIONED BY to enable 
       -- partition pushdown in the SQL optimizer, which improves performance.
      'connector' = 'odps',
      ... -- Fill in required parameters such as accessId. You can omit the 'partition' parameter and let the SQL optimizer filter partitions.
    );
    SELECT content, dt, hh FROM maxcompute_table
    WHERE dt >= '20220901' AND dt <= '20220903' AND hh >= '09' AND hh <= '17'; -- Specify partition filters in the WHERE clause.
  • Ler a partição com a maior ordem lexicográfica

    • Para ler a partição lexicograficamente maior de uma tabela de origem ou tabela de dimensão, defina o parâmetro partition como 'max_pt()'.

    • Para ler as duas partições lexicograficamente maiores de uma tabela de origem ou tabela de dimensão, defina o parâmetro partition como 'max_two_pt()'.

    • Para ler a partição lexicograficamente maior que também possui uma partição .done correspondente de uma tabela de origem ou tabela de dimensão, defina o parâmetro partition como 'max_pt_with_done()'.

    Geralmente, a partição lexicograficamente maior é a criada mais recentemente. A opção max_pt_with_done() é útil quando os dados na partição mais recente podem não estar prontos, e você deseja que a tabela de dimensão leia temporariamente de uma partição ligeiramente mais antiga, porém completa.

    Quando os dados de uma partição estiverem prontos, crie também uma partição vazia correspondente. Seu nome é o nome da partição de dados com .done anexado. Por exemplo, depois que os dados da partição dt=20220901 estiverem prontos, crie uma partição vazia chamada dt=20220901.done. Ao definir o parâmetro partition como max_pt_with_done(), a tabela de dimensão lerá apenas de partições que possuem uma partição .done correspondente. Partições de dados sem uma partição .done são temporariamente ignoradas. Para mais informações, consulte Qual é a diferença entre max_pt() e max_pt_with_done()?.

    Nota

    Uma tabela de origem determina a partição com a maior ordem lexicográfica apenas quando o job é iniciado. Ela para após ler todos os dados e não monitora novas partições. Se precisar ler novas partições continuamente, use o modo de tabela de origem incremental. Uma tabela de dimensão verifica e lê os dados mais recentes sempre que é atualizada.

Gravar em partições

  • Gravação em partições estáticas

    Para gravar dados em uma partição estática de uma tabela de resultados, defina o parâmetro partition usando o mesmo método utilizado para leitura de partições estáticas.

    Importante

    O parâmetro partition para uma tabela de resultados não suporta curingas (*).

  • Gravação em partições dinâmicas

    Para gravar em partições dinâmicas, onde os valores de partição são derivados dos dados, defina o parâmetro partition como uma lista separada por vírgulas de nomes de colunas de partição, ordenados por nível de partição em ordem crescente. Por exemplo, se uma tabela tiver três níveis de partição, dt, hh e mm, especifique 'partition' = 'dt,hh,mm'.

Inicialização lenta de jobs para tabelas de origem MaxCompute

As possíveis causas incluem:

  • A tabela MaxCompute contém muitos arquivos pequenos .

  • Alta latência de rede ocorre se o cluster de armazenamento MaxCompute e o cluster de computação Flink estiverem em regiões diferentes. Para resolver isso, coloque ambos os clusters na mesma região.

  • As permissões do MaxCompute estão configuradas incorretamente. A leitura de uma tabela de origem requer a permissão de download para a tabela MaxCompute.

Escolher um canal de dados

O MaxCompute fornece dois canais de dados: Batch Tunnel e Streaming Tunnel. Você pode escolher um canal de dados com base nos seus requisitos de consistência e eficiência de execução. A tabela a seguir compara os dois canais de dados.

Critério

Batch Tunnel

Streaming Tunnel

Consistência

O Batch Tunnel geralmente grava dados nas tabelas MaxCompute de forma mais confiável que o Streaming Tunnel e garante nenhuma perda de dados (semântica at-least-once).

A duplicação de dados pode ocorrer em algumas partições, mas apenas se uma exceção acontecer durante o processo de checkpoint enquanto o job estiver gravando em várias partições simultaneamente.

Garante nenhuma perda de dados (semântica at-least-once). No entanto, a duplicação de dados pode ocorrer se o job falhar por qualquer motivo.

Eficiência de execução

A eficiência geral de execução é menor que a do Streaming Tunnel porque os dados precisam ser confirmados (committed) durante o processo de checkpoint, o que envolve operações no servidor, como criação de arquivos.

Os dados não precisam ser confirmados durante o processo de checkpoint. Se você usar o Streaming Tunnel e definir o parâmetro numFlushThreads com um valor maior que 1, o destino poderá receber dados upstream continuamente enquanto descarrega (flush) dados. Isso resulta em uma eficiência geral de execução maior que a do Batch Tunnel.

Nota

Se jobs que usam o MaxCompute Batch Tunnel apresentarem checkpoints lentos ou com timeout, considere mudar para o Streaming Tunnel, desde que seu sistema downstream possa tolerar duplicação de dados.

Duplicação de dados em tabelas de resultados MaxCompute

Dados duplicados em uma tabela de resultados MaxCompute gravados por um job Flink podem resultar das seguintes causas:

  • Verifique a lógica do seu job. Mesmo que uma restrição de chave primária seja declarada na tabela de resultados MaxCompute, o Flink não realiza verificações de unicidade ao gravar em armazenamento externo. Além disso, tabelas não transacionais no MaxCompute não suportam restrições de chave primária. Portanto, se a lógica do seu job Flink gerar dados duplicados, essas duplicatas serão gravadas na tabela MaxCompute.

  • Verifique se vários jobs Flink estão gravando na mesma tabela MaxCompute simultaneamente. Como mencionado, o MaxCompute não impõe restrições de chave primária. Se vários jobs Flink produzirem os mesmos resultados, eles criarão registros duplicados na tabela.

  • O job Flink falha durante um checkpoint enquanto usa o Batch Tunnel. Quando ocorre uma falha durante um checkpoint, os dados da tabela de resultados podem já ter sido confirmados no servidor. Como resultado, quando o job se recupera do último checkpoint bem-sucedido, ele pode gravar dados duplicados referentes ao período entre o último checkpoint bem-sucedido e a falha.

  • Ocorre um failover do job Flink enquanto usa o Stream Tunnel. Ao gravar no MaxCompute com o Stream Tunnel, os dados são confirmados no servidor MaxCompute entre os checkpoints. Se o job sofrer failover e se recuperar do checkpoint mais recente, ele poderá gravar dados duplicados que foram processados após a conclusão do checkpoint, mas antes do failover. Para mais informações, consulte Escolher um canal de dados. Para evitar esse tipo de duplicação, mude para o modo Batch Tunnel.

  • Um job Flink que usa o Batch Tunnel sofre failover ou é reiniciado após ser cancelado (por exemplo, acionado pelo Autopilot). Em versões anteriores a vvr-6.0.7-flink-1.15, o job confirma dados na tabela de resultados MaxCompute ao ser desligado. Consequentemente, quando o job Flink para e depois se recupera do último checkpoint, ele pode criar dados duplicados para o período entre o último checkpoint e o desligamento. Para resolver esse problema, atualize sua versão do Flink para vvr-6.0.7-flink-1.15 ou posterior.

Job MaxCompute falha com 'Invalid partition spec'

  • Causa: Este erro ocorre quando os dados gravados no MaxCompute contêm valores inválidos na coluna de partição. Valores inválidos incluem strings vazias, valores nulos ou valores que contêm um sinal de igual (=), uma vírgula (,) ou uma barra (/).

  • Solução: Verifique se os valores na coluna de partição dos seus dados de origem são válidos.

Erro 'No more available blockId' em jobs MaxCompute

  • Causa: O número de blocos gravados na tabela de resultados MaxCompute excedeu o limite. Isso geralmente é causado pelo descarregamento (flush) de pequenas quantidades de dados com muita frequência.

  • Solução: Ajuste os parâmetros batchSize e flushIntervalMs.

Usar a dica SHUFFLE_HASH

Por padrão, cada instância paralela armazena em cache toda a tabela de dimensão. Se uma tabela de dimensão for grande, use a dica SHUFFLE_HASH para distribuir os dados da tabela uniformemente entre as instâncias paralelas e reduzir o consumo de memória heap da JVM. No exemplo a seguir, os dados das tabelas de dimensão dim_1 e dim_3 são distribuídos entre as instâncias paralelas, enquanto os dados de dim_2 permanecem totalmente armazenados em cache em cada uma delas.

-- Create a source table and three dimension tables.
CREATE TABLE source_table (k VARCHAR, v VARCHAR) WITH ( ... );
CREATE TABLE dim_1 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
CREATE TABLE dim_2 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
CREATE TABLE dim_3 (k VARCHAR, v VARCHAR) WITH ('connector' = 'odps', 'cache' = 'ALL', ... );
-- Specify the names of the dimension tables to distribute in the SHUFFLE_HASH hint.
SELECT /*+ SHUFFLE_HASH(dim_1), SHUFFLE_HASH(dim_3) */
k, s.v, d1.v, d2.v, d3.v
FROM source_table AS s
INNER JOIN dim_1 FOR SYSTEM_TIME AS OF PROCTIME() AS d1 ON s.k = d1.k
LEFT JOIN dim_2 FOR SYSTEM_TIME AS OF PROCTIME() AS d2 ON s.k = d2.k
LEFT JOIN dim_3 FOR SYSTEM_TIME AS OF PROCTIME() AS d3 ON s.k = d3.k;

Configurar CacheReloadTimeBlackList

Especifica as janelas de tempo durante as quais as atualizações da tabela de dimensão são desabilitadas.

  • tipo de dados: String

  • Use -> entre os horários de início e fim.

  • Separe múltiplas janelas de tempo com uma ,.

  • Formato de hora: AAAA-MM-DD HH:mm. Se você especificar apenas horas e minutos, a janela de tempo será aplicada diariamente por padrão.

'cacheReloadTimeBlackList' = '14:00 -> 15:00,23:00 -> 01:00'

Cenário

Valor

Janela de tempo única

14:00 -> 15:00

Múltiplas janelas de tempo

14:00 -> 15:00,23:00 -> 01:00

Janela de tempo especial

14:00 -> 15:00, 23:00 -> 01:00,2025-10-01 22:00 -> 2025-10-01 23:00

Erro: java.io.EOFException: SSL peer shut down incorrectly

  • Detalhes do erro

    Caused by: java.io.EOFException: SSL peer shut down incorrectly
        at sun.security.ssl.SSLSocketInputRecord.decodeInputRecord(SSLSocketInputRecord.java:239) ~[?:1.8.0_302]
        at sun.security.ssl.SSLSocketInputRecord.decode(SSLSocketInputRecord.java:190) ~[?:1.8.0_302]
        at sun.security.ssl.SSLTransport.decode(SSLTransport.java:109) ~[?:1.8.0_302]
        at sun.security.ssl.SSLSocketImpl.decode(SSLSocketImpl.java:1392) ~[?:1.8.0_302]
        at sun.security.ssl.SSLSocketImpl.readHandshakeRecord(SSLSocketImpl.java:1300) ~[?:1.8.0_302]
        at sun.security.ssl.SSLSocketImpl.startHandshake(SSLSocketImpl.java:435) ~[?:1.8.0_302]
        at com.mysql.cj.protocol.ExportControlled.performTlsHandshake(ExportControlled.java:347) ~[?:?]
        at com.mysql.cj.protocol.StandardSocketFactory.performTlsHandshake(StandardSocketFactory.java:194) ~[?:?]
        at com.mysql.cj.protocol.a.NativeSocketConnection.performTlsHandshake(NativeSocketConnection.java:101) ~[?:?]
        at com.mysql.cj.protocol.a.NativeProtocol.negotiateSSLConnection(NativeProtocol.java:308) ~[?:?]
        at com.mysql.cj.protocol.a.NativeAuthenticationProvider.connect(NativeAuthenticationProvider.java:204) ~[?:?]
        at com.mysql.cj.protocol.a.NativeProtocol.connect(NativeProtocol.java:1369) ~[?:?]
        at com.mysql.cj.NativeSession.connect(NativeSession.java:133) ~[?:?]
        at com.mysql.cj.jdbc.ConnectionImpl.connectOneTryOnly(ConnectionImpl.java:949) ~[?:?]
        at com.mysql.cj.jdbc.ConnectionImpl.createNewIO(ConnectionImpl.java:819) ~[?:?]
        at com.mysql.cj.jdbc.ConnectionImpl.<init>(ConnectionImpl.java:449) ~[?:?]
        at com.mysql.cj.jdbc.ConnectionImpl.getInstance(ConnectionImpl.java:242) ~[?:?]
        at com.mysql.cj.jdbc.NonRegisteringDriver.connect(NonRegisteringDriver.java:198) ~[?:?]
        at org.apache.flink.connector.jdbc.internal.connection.SimpleJdbcConnectionProvider.getOrEstablishConnection(SimpleJdbcConnectionProvider.java:128) ~[?:?]
        at org.apache.flink.connector.jdbc.internal.AbstractJdbcOutputFormat.open(AbstractJdbcOutputFormat.java:54) ~[?:?]
        ... 14 more
  • Causa

    Este erro geralmente ocorre quando um banco de dados MySQL tem o protocolo SSL habilitado, mas a conexão SSL do cliente não está configurada corretamente. Por exemplo, com o driver MySQL 8.0.27 e um banco de dados MySQL com SSL habilitado, esse erro acontece porque o método de acesso padrão do driver não usa SSL.

  • Solução

    Anexe characterEncoding=utf-8&useSSL=false ao parâmetro URL da tabela de dimensão MySQL. Por exemplo:

    'url'='jdbc:mysql://***.***.***.***:3306/test?characterEncoding=utf-8&useSSL=false'

Alterações no tipo de chave MySQL** bigint unsigned**

O Flink não suporta o tipo de dados bigint unsigned. Para evitar possível estouro de dados, o Flink mapeia a chave primária bigint unsigned para o tipo decimal. Durante a sincronização com o Hologres, o sistema converte a coluna para o tipo text porque o Hologres não suporta bigint unsigned nem permite o tipo decimal como chave primária.

Considere esse comportamento durante o design e desenvolvimento. Para manter a coluna como tipo decimal, crie manualmente a tabela no Hologres antes de iniciar a sincronização. Nessa tabela, você pode definir uma coluna diferente como chave primária ou não definir nenhuma chave primária. No entanto, essa abordagem pode levar à duplicação de dados, pois a chave primária original não estará mais impondo unicidade. Você deve então lidar com esse problema no nível da aplicação, por exemplo, tolerando a duplicação de dados ou implementando lógica de deduplicação.

Flink para RDS: Atualização vs. inserção

Se uma chave primária for definida na DDL, o conector usará uma instrução INSERT INTO tablename(field1,field2, field3, ...) VALUES(value1, value2, value3, ...) ON DUPLICATE KEY UPDATE field1=value1,field2=value2, field3=value3, ...;. Essa instrução insere um novo registro se a chave primária não existir, ou atualiza o registro existente caso ela exista. Se nenhuma chave primária for declarada na DDL, o conector insere novos registros com uma instrução insert into.

Uso de índice único com GROUP BY

  • Declare o índice único na cláusula GROUP BY do seu job.

  • Se uma tabela RDS usar uma chave primária de incremento automático, não a declare como PRIMARY KEY no job Flink.

Mapeamento INT UNSIGNED: MySQL para Flink SQL

O MySQL JDBC Driver mapeia inteiros sem sinal para tipos de dados Java maiores para preservar a precisão. Especificamente, o driver mapeia valores MySQL INT UNSIGNED para o tipo Java LONG, que o Flink SQL trata como BIGINT. Da mesma forma, ele mapeia valores MySQL BIGINT UNSIGNED para o tipo Java BigInteger, que o Flink SQL trata como DECIMAL(20, 0).

Erro: Incorrect string value

  • Detalhes do erro

    Caused by: java.sql.BatchUpdateException: Incorrect string value: '\xF0\x9F\x98\x80\xF0\x9F...' for column 'test' at row 1
    at sun.reflect.GeneratedConstructorAccessor59.newInstance(Unknown Source)
    at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
    at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
    at com.mysql.cj.util.Util.handleNewInstance(Util.java:192)
    at com.mysql.cj.util.Util.getInstance(Util.java:167)
    at com.mysql.cj.util.Util.getInstance(Util.java:174)
    at com.mysql.cj.jdbc.exceptions.SQLError.createBatchUpdateException(SQLError.java:224)
    at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchedInserts(ClientPreparedStatement.java:755)
    at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchInternal(ClientPreparedStatement.java:426)
    at com.mysql.cj.jdbc.StatementImpl.executeBatch(StatementImpl.java:796)
    at com.alibaba.druid.pool.DruidPooledPreparedStatement.executeBatch(DruidPooledPreparedStatement.java:565)
    at com.alibaba.ververica.connectors.rds.sink.RdsOutputFormat.executeSql(RdsOutputFormat.java:488)
    ... 15 more
  • Causa

    Os dados contêm caracteres especiais ou usam uma codificação de caracteres que o banco de dados não suporta.

  • Solução

    Adicione characterEncoding=UTF-8 à URL ao conectar-se a um banco de dados MySQL via JDBC. Por exemplo: jdbc:mysql://<internal address>/<databaseName>?characterEncoding=UTF-8.

Deadlocks no MySQL (TDDL/RDS)

  • Problema

    Ocorre um deadlock durante a gravação de dados no MySQL (TDDL/RDS).

    Importante

    No Realtime Compute for Apache Flink, se você utilizar um banco de dados relacional como o MySQL como sink (por meio do conector TDDL/RDS), gravações frequentes em uma tabela ou recurso podem causar deadlocks.

    Considere que uma operação INSERT precise adquirir dois locks, (A,B), em sequência. O lock A é um lock de intervalo. Existem duas transações, (T1,T2), e o schema da tabela é (id(auto-incrementing primary key),nid(unique key)). A transação T1 contém duas instruções, insert(null,2),(null,1), enquanto a T2 contém uma instrução, insert(null,2).

    1. No instante t, a T1 executa sua primeira instrução INSERT. Nesse momento, a T1 detém ambos os locks (A,B).

    2. No instante t+1, a T2 inicia uma operação de inserção e precisa aguardar o lock A para bloquear o intervalo (-inf,2]. Nesse momento, o lock A está retido pela T1 e já bloqueou o intervalo (-inf,2]. Como os intervalos possuem uma relação de inclusão, a T2 depende que a T1 libere o lock A.

    3. No instante t+2, a T1 executa sua segunda instrução INSERT, que requer o lock A no intervalo (-inf,1]. Como esse intervalo é um subconjunto de (-inf,2], a T1 precisa entrar na fila e aguardar a T2 liberar o lock. Consequentemente, a T1 passa a depender que a T2 libere o lock A.

    O deadlock ocorre porque a T1 e a T2 ficam esperando mutuamente que a outra libere seus respectivos locks.

  • O RDS/TDDL e o Tablestore utilizam mecanismos de bloqueio diferentes.

    • RDS/TDDL: O lock de linha do InnoDB é aplicado ao índice, e não a um registro individual. Por isso, mesmo ao acessar linhas diferentes, pode ocorrer um conflito de lock se as linhas compartilharem a mesma chave de índice. Isso pode impedir atualizações em todo um intervalo de dados.

    • Tablestore: Utiliza um lock de linha única, que não afeta atualizações em outros dados.

  • Soluções para deadlocks

    Para cenários de alta QPS/TPS ou gravação com alta concorrência, utilize o Tablestore como tabela de resultados para evitar deadlocks. Em geral, não é recomendável usar o TDDL ou o RDS como tabela de resultados para um job do Flink.

    Caso seja obrigatório usar um banco de dados relacional como o MySQL como nó sink, considere as seguintes recomendações:

    • Garanta que nenhuma outra carga de trabalho esteja lendo ou gravando nas mesmas tabelas.

    • Se o volume de dados do job for pequeno, tente gravar os dados em modo single-thread. No entanto, em cenários de alta QPS/TPS e alta concorrência, essa abordagem reduz o desempenho de gravação.

    • Evite usar uma unique key sempre que possível, pois gravações em tabelas com unique key podem causar deadlocks. Se os requisitos do seu negócio exigirem uma unique key, defina-a ordenando as colunas da mais seletiva para a menos seletiva. Isso reduz significativamente a probabilidade de deadlock. Por exemplo, coloque uma coluna com hash MD5 antes da coluna day_time(20171010) para garantir que a unique key seja definida com as colunas mais seletivas primeiro.

    • Utilize sharding de banco de dados e divisão de tabelas com base nas características da sua carga de trabalho para distribuir as gravações por várias tabelas. Para detalhes de implementação, entre em contato com o administrador do banco de dados.

Falha na atualização da estrutura da tabela downstream

A sincronização da estrutura da tabela não rastreia instruções DDL. Em vez disso, ela detecta alterações de schema comparando registros consecutivos. Se ocorrer uma alteração DDL sem mudanças subsequentes nos dados upstream, a estrutura da tabela downstream não será atualizada. Para mais detalhes, consulte política de sincronização para alterações na estrutura da tabela.

Erro de timeout na resposta de finish split

Esse erro ocorre quando a alta utilização da CPU em uma task impede que ela responda às solicitações RPC do coordenador. Para resolver, aumente os recursos de CPU do TaskManager na página de configuração de recursos.

Impacto de alterações de schema durante o full load

Uma alteração de schema durante a fase de full load pode fazer o job falhar ou impedir que a mudança de schema seja sincronizada. Para resolver, pare o job, exclua a tabela downstream e reinicie o job sem estado.

Alterações de schema não suportadas durante a sincronização CTAS/CDAS

Ressincronize os dados da tabela. Para isso, pare o job, remova a tabela downstream e reinicie o job de sincronização com início sem estado. Evite fazer alterações incompatíveis desse tipo, pois o job falhará novamente ao reiniciar. Para detalhes sobre as alterações de schema suportadas, consulte instrução CREATE TABLE AS (CTAS).

Atualizações por retração no ClickHouse

Atualizações por retração são suportadas para uma tabela de resultados do ClickHouse se você especificar uma primary key na DDL da tabela de resultados do Flink e definir o parâmetro ignoreDelete como false. No entanto, isso causa uma queda significativa no desempenho.

O ClickHouse é um sistema de gerenciamento de banco de dados colunar projetado para processamento analítico online (OLAP), e seu suporte a operações UPDATE e DELETE é limitado. Se você especificar uma primary key na DDL do Flink, o conector tentará usar ALTER TABLE UPDATE e ALTER TABLE DELETE para atualizar e excluir dados. Essas operações são altamente ineficientes.

Visibilidade de dados no ClickHouse

  • Para uma tabela de resultados do ClickHouse com exactly-once semantics desativado (padrão), os dados tornam-se visíveis assim que o buffer é descarregado. O sistema descarrega automaticamente esse buffer quando o número de registros atinge o valor de batchSize ou quando o tempo desde a última gravação excede flushIntervalMs. Não é necessário aguardar a conclusão de um checkpoint.

  • Para uma tabela de resultados do ClickHouse com exactly-once semantics ativado, os dados tornam-se visíveis somente após a conclusão bem-sucedida do checkpoint correspondente.

Visualizar resultados de print

Existem duas maneiras de visualizar os resultados de print:

  • No Real-time Compute Development Console:

    1. No painel de navegação à esquerda do Real-time Compute Development Console, escolha Operations Center > Job Operations.

    2. Clique no nome do job desejado.

    3. Clique na aba Job Log.

    4. Na aba Runtime Log, selecione o job em execução na lista suspensa ao lado de Job.

    5. Na aba Running Task Managers, clique em um Path, ID.

    6. Clique na aba Log para visualizar os resultados de print.

  • Na Flink UI:

    1. No painel de navegação à esquerda do Real-time Compute Development Console, escolha Operations Center > Job Operations.

    2. Clique no nome do job desejado.

    3. Na aba Status Overview, clique em Flink UI.

    4. Clique em Task Managers.

    5. Clique em um Path, ID.

    6. Na aba logs, visualize os resultados de print.

Join com tabela de dimensão não retorna dados

Certifique-se de que o schema — incluindo tipos de dados e nomes de colunas — na instrução DDL seja consistente com o da tabela física.

max_pt() e max_pt_with_done()

A função max_pt() retorna a partição lexicograficamente maior. A função max_pt_with_done() retorna a partição lexicograficamente maior que possui uma partição .done correspondente. Por exemplo, considere a seguinte lista de partições:

  • ds=20190101

  • ds=20190101.done

  • ds=20190102

  • ds=20190102.done

  • ds=20190103

Com base nessa lista, max_pt() e max_pt_with_done() comportam-se da seguinte forma:

  • partition retorna a partição ds=20190102.

  • partition retorna a partição ds=20190103.

Erro no job de gravação Paimon: "Heartbeat of TaskManager timed out"

A causa mais provável desse erro é memória heap insuficiente no TaskManager. O Paimon utiliza a memória heap principalmente das seguintes formas:

  • Cada instância paralela do operador writer para uma tabela de primary key do Paimon possui um buffer de memória para ordenação. O tamanho desse buffer é controlado pela propriedade de tabela write-buffer-size, cujo padrão é 256 MB.

  • O Paimon usa o formato de arquivo ORC por padrão, o que exige um buffer de memória adicional para converter dados em memória para um formato colunar em lotes. O tamanho desse buffer é controlado pela propriedade de tabela orc.write.batch-size, cujo padrão é 1024, significando que o buffer armazena 1024 linhas de dados.

  • Cada bucket modificado possui um objeto writer dedicado para gravar seus dados.

Com base nesses padrões de uso, aqui estão as possíveis causas de memória heap insuficiente e suas soluções:

  • O valor de write-buffer-size é muito grande.

    Tente reduzir esse parâmetro. No entanto, um buffer muito pequeno pode levar a gravações frequentes em disco e acionar compactações de arquivos pequenos com mais frequência, o que impacta o desempenho de gravação.

  • Um único registro de dados é muito grande.

    Por exemplo, se um registro contiver um campo JSON de 4 MB, o buffer ORC pode crescer para 4 MB × 1024 = 4 GB, consumindo uma quantidade significativa de memória heap. Você tem duas soluções:

    • Reduza o valor de orc.write.batch-size.

    • Se você não precisar realizar consultas ad-hoc (OLAP) na tabela de resultados Paimon e precisar apenas de consumo em batch ou streaming, defina as propriedades de tabela 'file.format' = 'avro' e 'metadata.stats-mode' = 'none' durante a criação da tabela. Isso altera a tabela para o formato Avro e desativa a coleta de estatísticas.

      Nota

      Esses parâmetros só podem ser definidos durante a criação da tabela. Eles não podem ser modificados com uma instrução ALTER TABLE ou um SQL hint após a criação da tabela.

  • Gravar em muitas partições simultaneamente ou ter muitos buckets por partição cria um número excessivo de objetos writer.

    Revise a configuração da coluna de partição para garantir que seja apropriada. Certifique-se de que SQL incorreto não esteja causando a gravação de dados inesperados na coluna de partição. Além disso, verifique se o número de buckets é razoável. Como melhor prática, o tamanho total de dados por bucket deve ficar em torno de 2 GB e não exceder 5 GB. Para detalhes sobre como ajustar o número de buckets, consulte Ajustar o número de buckets para uma tabela de bucket fixo.

Erro: "Sink materializer must not be used with Paimon sink"

O operador sink materializer lida com dados fora de ordem provenientes de joins em cascata em jobs de streaming. No entanto, em jobs que gravam em uma tabela Paimon, esse operador introduz sobrecarga e pode causar resultados incorretos quando há agregação. Não utilize o operador sink materializer com um sink Paimon.

Desative o operador sink materializer definindo o parâmetro table.exec.sink.upsert-materialize como false usando uma instrução SET ou como um parâmetro de runtime. Se você também precisar lidar com dados fora de ordem, consulte tratamento de dados fora de ordem.

Paimon: Erro File deletion conflicts detected ou LSM conflicts detected

Esse erro pode ocorrer pelos seguintes motivos:

  • Vários jobs estão gravando na mesma partição da mesma tabela Paimon. Nesse caso, o Paimon resolve o conflito por meio de failover and restart. Esse é um comportamento esperado e nenhuma ação é necessária se o erro não se repetir.

  • O job foi restaurado a partir de um estado desatualizado, o que faz o erro recorrer. Para resolver, restaure o job a partir do seu estado mais recente ou start it without state.

  • O Paimon não suporta gravações separadas de múltiplas instruções INSERT dentro de um único job. Em vez disso, use uma instrução UNION ALL para gravar múltiplos fluxos de dados na tabela Paimon.

  • A concorrência do nó Global Committer ou do nó Compaction Coordinator (ao gravar em uma tabela Append Scalable) é maior que 1. A concurrency para esses nós deve ser 1 para garantir a data consistency.

Erro "File xxx not found" no job de consumo Paimon

O consumo de tabelas Paimon depende de arquivos de snapshot. Se o período de retenção de snapshot for muito curto ou o job de consumo for ineficiente, os arquivos de snapshot podem expirar e ser excluídos antes que o job termine. Isso faz com que o job de consumo falhe.

Para resolver esse problema, você pode ajustar o período de retenção do arquivo de snapshot, especificar um consumer ID ou otimizar o job de consumo. Para verificar quais arquivos de snapshot estão disponíveis e seus timestamps de criação, consulte a tabela de sistema Snapshots.

Erro no job Paimon: No space left on device

  • Um número excessivo de arquivos de cache pode causar esse erro se o seu job executar uma consulta Paimon, como usar uma tabela Paimon como tabela de dimensão ou definir changelog-producer='lookup'. Para evitar isso, use SQL Hints para definir os seguintes parâmetros e limitar o espaço máximo em disco e o tempo de retenção do cache de consulta.

    • lookup.cache-max-disk-size: O espaço máximo em disco local que o cache de consulta pode usar. Os valores recomendados incluem 256 MB, 512 MB e 1 GB.

    • lookup.cache-file-retention: O período de retenção para arquivos de cache de consulta. Os valores recomendados incluem 30 min, 15 min ou um intervalo menor.

  • Para jobs que gravam em uma tabela Paimon, use SQL Hints para definir os seguintes parâmetros. Essas configurações limitam o tamanho dos arquivos temporários locais durante o processo de gravação, evitando falta de espaço em disco.

    • write-buffer-spillable: Controla se o buffer de gravação pode transbordar para o disco. Definir isso como false impede completamente que o buffer use qualquer espaço em disco.

    • write-buffer-spill.max-disk-size: O espaço máximo em disco que o buffer de gravação pode usar quando transborda. Os valores recomendados incluem 256 MB, 512 MB e 1 GB.

Gerenciar arquivos Paimon no OSS

  • O Paimon retém arquivos de dados históricos para permitir o acesso a versões anteriores de uma tabela. Você pode ajustar a política de retenção desses arquivos para gerenciar o armazenamento. Para instruções detalhadas, consulte Limpar dados expirados.

  • Uma configuração inadequada de coluna de partição ou um número excessivo de buckets também pode causar esse problema. Como melhor prática, busque um tamanho de dados de aproximadamente 2 GB por bucket, com um máximo de 5 GB. Para mais informações, consulte Bucketing.

  • Por padrão, os arquivos de dados são salvos no formato ORC. Para reduzir o tamanho total dos arquivos de dados, utilize o formato de compressão ZSTD definindo o parâmetro de tabela 'file.compression' = 'zstd' ao criar a tabela.

    Nota

    Este parâmetro só pode ser definido durante a criação da tabela e não pode ser modificado posteriormente com uma instrução ALTER TABLE ou um SQL hint.

A visibilidade dos dados depende do intervalo de checkpoint****

Sim. O Paimon depende de checkpoints para garantir a semântica exactly-once. Os dados são confirmados e tornam-se visíveis downstream somente após a conclusão do checkpoint. Antes dessa confirmação, os dados no buffer local são descarregados para o sistema de arquivos remoto, mas ainda não estão legíveis.

Aumento lento de memória em jobs Paimon de longa duração

  • Um aumento no uso de memória é esperado se o rps do job também estiver aumentando lentamente.

  • Se você usar um catalog de sistema de arquivos Paimon para ler ou gravar no OSS, certifique-se de configurar os parâmetros de catalog fs.oss.endpoint, fs.oss.accessKeyId e fs.oss.accessKeySecret. Caso contrário, o job do Flink pode apresentar um vazamento de memória lento, que é um problema conhecido na comunidade.

IllegalArgumentException: timeout value is negative

  • Detalhes do erro

    2021-02-24 15:14:58
    java.lang.RuntimeException: java.lang.RuntimeException: java.lang.IllegalArgumentException: timeout value is negative
        at com.alibaba.ververica.connectors.common.source.reader.ParallelReader.run(ParallelReader.java:166)
        at com.alibaba.ververica.connectors.common.source.AbstractParallelSourceBase.run(AbstractParallelSourceBase.java:205)
        at com.alibaba.ververica.connectors.metaq.source.MetaQRowDataSource.run(MetaQRowDataSource.java:84)
        at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:100)
        at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:63)
        at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:213)
    Caused by: java.lang.RuntimeException: java.lang.IllegalArgumentException: timeout value is negative
        at com.alibaba.ververica.connectors.common.source.reader.ParallelReader.runImpl(ParallelReader.java:244)
  • Causa do erro

    Quando nenhuma nova mensagem MQ é consumida, a thread MetaQSource entra em suspensão por um intervalo definido pelo parâmetro pullIntervalMs, cujo padrão é -1. O job então falha com uma IllegalArgumentException porque a duração da suspensão não pode ser negativa.

  • Solução

    Defina o parâmetro pullIntervalMs com um valor não negativo.

Descoberta de alteração de partição

  • Para versões do Realtime Compute for Apache Flink anteriores à 6.0.2, o operador source busca o número de partições a cada 5 a 10 minutos. Um failover é acionado se o número de partições diferir em três verificações consecutivas. Como resultado, o source inicia um failover dentro de 10 a 30 minutos. Após o reinício do job, ele lê do conjunto atualizado de partições.

  • Para versões 6.0.2 e posteriores do Realtime Compute for Apache Flink, o operador source busca o número de partições a cada 5 minutos por padrão. Quando novas partições são descobertas, elas são atribuídas diretamente ao operador source no TaskManager, que então começa a ler os dados. Nenhum failover de job é necessário, permitindo que o source detecte alterações de partição dentro de 1 a 5 minutos.

Erro: Backpressure exceeds reject limit

  • Detalhes do erro

    26      at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:47)
    27      at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:911)
    28      at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:879)
    29      ... 13 more
    30  Caused by: java.lang.RuntimeException: Rpc Exception failed errorCount=6 with RpcException: request niagara.table.proto.UpsertRecordBatchRequest@b53e2d74 failed on final try 4, maxAttempts=4, sn=11.117.xxx, errorCode=11, msg=BackPresure Exceed Reject Limit [method:UpsertRecordBatch,transaction_id:xxx,table_id:xxx,table_version:128, actor_id:74538 xxx,worker_address:11.117.xxx]
    31          at com.alibaba.ververica.connectors.hologres.sink.HologresOutputFormat.sync(HologresOutputFormat.java:264)
    32          at com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction.snapshotState(OutputFormatSinkFunction.java:91)
    33          at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:128)
    34          at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:101)
    35          at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.snapshotState(AbstractUdfStreamOperator.java:90)
    36          at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.snapshotState(StreamOperatorStateHandler.java:186)
    37          ... 23 more
  • Causa

    A pressão de gravação na instância Hologres está muito alta.

  • Solução

    Entre em contato com o suporte técnico do Hologres fornecendo as informações da sua instância para solicitar um upgrade.

Erro: remaining connection slots are reserved for non-replication superuser connections

  • Detalhes do erro

    Caused by: com.alibaba.hologres.client.exception.HoloClientWithDetailsException: failed records 1, first:Record{schema=org.postgresql.model.TableSchema@188365, values=[f06b41455c694d24a18d0552b8b0****, com.chot.tpfymnq.meta, 2022-04-02 19:46:40.0, 28, 1, null], bitSet={0, 1, 2, 3, 4}},first err:[106]FATAL: remaining connection slots are reserved for non-replication superuser connections
        at com.alibaba.hologres.client.impl.Worker.handlePutAction(Worker.java:406) ~[?:?]
        at com.alibaba.hologres.client.impl.Worker.run(Worker.java:118) ~[?:?]
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_302]
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_302]
        ... 1 more
    Caused by: com.alibaba.hologres.org.postgresql.util.PSQLException: FATAL: remaining connection slots are reserved for non-replication superuser connections
        at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2553) ~[?:?]
        at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.readStartupMessages(QueryExecutorImpl.java:2665) ~[?:?]
        at com.alibaba.hologres.org.postgresql.core.v3.QueryExecutorImpl.<init>(QueryExecutorImpl.java:147) ~[?:?]
        at com.alibaba.hologres.org.postgresql.core.v3.ConnectionFactoryImpl.openConnectionImpl(ConnectionFactoryImpl.java:273) ~[?:?]
        at com.alibaba.hologres.org.postgresql.core.ConnectionFactory.openConnection(ConnectionFactory.java:51) ~[?:?]
        at com.alibaba.hologres.org.postgresql.jdbc.PgConnection.<init>(PgConnection.java:240) ~[?:?]
        at com.alibaba.hologres.org.postgresql.Driver.makeConnection(Driver.java:478) ~[?:?]
        at com.alibaba.hologres.org.postgresql.Driver.connect(Driver.java:277) ~[?:?]
        at java.sql.DriverManager.getConnection(DriverManager.java:674) ~[?:1.8.0_302]
        at java.sql.DriverManager.getConnection(DriverManager.java:217) ~[?:1.8.0_302]
        at com.alibaba.hologres.client.impl.ConnectionHolder.buildConnection(ConnectionHolder.java:122) ~[?:?]
        at com.alibaba.hologres.client.impl.ConnectionHolder.retryExecute(ConnectionHolder.java:195) ~[?:?]
        at com.alibaba.hologres.client.impl.ConnectionHolder.retryExecute(ConnectionHolder.java:184) ~[?:?]
        at com.alibaba.hologres.client.impl.Worker.doHandlePutAction(Worker.java:460) ~[?:?]
        at com.alibaba.hologres.client.impl.Worker.handlePutAction(Worker.java:389) ~[?:?]
        at com.alibaba.hologres.client.impl.Worker.run(Worker.java:118) ~[?:?]
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) ~[?:1.8.0_302]
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) ~[?:1.8.0_302]
        ... 1 more
  • Causa

    O limite de conexões da instância Hologres foi excedido.

  • Solução

    • Verifique o app_name das conexões de cada Frontend (FE) para contar as conexões de cliente Hologres provenientes do flink-connector.

    • Verifique se há outros jobs conectados ao Hologres.

    • Libere conexões. Para mais informações, consulte gerenciamento de conexões.

No table is defined in publication

  • Detalhes do erro

    Excluir e recriar uma tabela com o mesmo nome pode fazer com que um job reporte no table is defined in publication.

  • Causa

    Excluir uma tabela não remove sua publicação associada.

  • Solução

    1. No Hologres, execute o comando select * from pg_publication where pubname not in (select pubname from pg_publication_tables); para consultar informações de publicação que não foram limpas quando uma tabela foi excluída.

    2. Execute a instrução drop publication xx; para remover a publicação restante.

    3. Reinicie o job.

Intervalo de checkpoint e visibilidade de dados

O intervalo de checkpoint do conector sink Flink Hologres não controla diretamente a visibilidade dos dados no Hologres. Sua função principal é definir o SLA para recuperação de falhas.

O conector Hologres não suporta transações. Ele descarrega periodicamente o buffer em memória para o banco de dados. Um checkpoint garante que todos os dados sejam descarregados até o momento de sua conclusão, mas o conector não espera que todo o intervalo passe antes de descarregar. O conector aciona um descarregamento antecipado se certas condições de buffer forem atendidas (para mais informações, consulte Hologres, Hologres e Hologres). Como um data warehouse normalmente não exige consistência transacional, o conector descarrega dados de forma assíncrona em segundo plano. Em seguida, realiza um descarregamento final forçado durante cada checkpoint para preparar a recuperação de falhas.

Implantação de job gera erro permission denied for database****

  • Causa

    A partir do Realtime Compute for Apache Flink VVR 8.0.4, o conector impõe o modo JDBC para consumir o binary log de instâncias Hologres V2.0 ou posteriores. Para essas instâncias, contas que não são superusuário requerem permissões especiais para consumir o binary log no modo JDBC.

  • Solução

    Conceda à conta que não é superusuário permissões para consumir o binary log no modo JDBC.

    user_name refere-se a um ID de conta Alibaba Cloud ou a um usuário RAM. Para mais informações, consulte visão geral de contas.

    -- For the expert permission model, grant the CREATE permission and the replication role to the user.
    GRANT CREATE ON DATABASE <db_name> TO <user_name>;
    alter role <user_name> replication;
    -- If the database uses the simple permission model (SPM), you cannot run GRANT statements. 
    -- Instead, use spm_grant to grant the user the Admin role for the database. You can also grant permissions directly in HoloWeb.
    call spm_grant('<db_name>_admin', '<user_name>');
    alter role <user_name> replication;

Falha na recuperação do job: table id parsed from checkpoint is different from the current table id

  • Causa

    Essa exceção ocorre nas versões VVR 8.0.5 a VVR 8.0.8 do Realtime Compute for Apache Flink. Quando um job que usa uma tabela source de binlog do Hologres se recupera de um checkpoint, o motor aplica uma verificação estrita no table id. Se o table id atual da tabela Hologres não corresponder ao armazenado no checkpoint, a recuperação falha. Isso indica que a tabela source foi truncada ou recriada enquanto o job estava em execução.

  • Solução

    Atualize para a versão VVR 8.0.9 ou posterior e reinicie o job. A versão VVR 8.0.9 remove a verificação estrita do table id para acomodar cenários de negócios complexos. No entanto, evite recriar uma tabela source de binlog. Quando uma tabela é recriada, todo o seu histórico de binlog é apagado. Se o Flink usar o consumer offset da tabela antiga para ler dados da nova tabela, isso pode causar inconsistências nos dados.

Precisão inesperada de dados Binlog no modo JDBC

  • Causa

    No Realtime Compute for Apache Flink 8.0.10 e anteriores, ocorre precisão inesperada de dados se a precisão do tipo DECIMAL declarada em uma DDL do Flink para uma tabela source Binlog não corresponder à precisão no Hologres.

  • Solução

    Esse problema foi corrigido no Realtime Compute for Apache Flink 8.0.11. No entanto, garanta que a precisão do tipo DECIMAL seja consistente entre o Flink e o Hologres para evitar perda de precisão.

Excluir e recriar uma tabela com o mesmo nome pode fazer com que jobs reportem exceções no table is defined in publication ou The table xxx has no slot named xxx****

  • Causa

    Isso ocorre porque excluir a tabela não remove sua publicação associada.

  • Solução

    Solução 1: No Hologres, execute a instrução select * from pg_publication where pubname not in (select pubname from pg_publication_tables); para encontrar publicações deixadas por tabelas excluídas. Em seguida, execute a instrução drop publication xx; para removê-las. Finalmente, reinicie o job.

    Solução 2: Use a versão VVR 8.0.5 ou posterior. O conector lida automaticamente com a limpeza.

ClassCastException ao ler do Hologres

  • Detalhes do erro

    A mensagem de erro é semelhante à seguinte:

    java.lang.ClassCastException: class java.lang.Integer cannot be cast to class java.lang.Long (java.lang.Integer and java.lang.Long are in module java.base of loader 'bootstrap')
  • Causa

    Esse erro ocorre quando o tipo de um campo na DDL do Flink não corresponde ao tipo do campo correspondente na tabela física do Hologres. Por exemplo, um campo é definido como BIGINT na DDL do Flink, mas o campo correspondente na tabela Hologres é INTEGER. Como a verificação de tipo é ignorada para valores NULL, o job lança a exceção somente quando lê dados reais.

  • Solução

    Consulte a documentação do Hologres sobre Resumo de tipos de dados para garantir que os tipos de campos na sua DDL do Flink correspondam aos da tabela física do Hologres.

LogSizeTooLargeException

  • Detalhes do erro

    Caused by: com.aliyun.openservices.aliyun.log.producer.errors.LogSizeTooLargeException: the logs is 8785684 bytes which is larger than MAX_BATCH_SIZE_IN_BYTES 8388608
    at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.ensureValidLogSize(LogAccumulator.java:249)
    at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.doAppend(LogAccumulator.java:103)
    at com.aliyun.openservices.aliyun.log.producer.internals.LogAccumulator.append(LogAccumulator.java:84)
    at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:385)
    at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:308)
    at com.aliyun.openservices.aliyun.log.producer.LogProducer.send(LogProducer.java:211)
    at com.alibaba.ververica.connectors.sls.sink.SLSOutputFormat.writeRecord(SLSOutputFo
    rmat.java:100)
  • Causa

    Esse erro ocorre porque um log de linha única enviado ao Log Service excede o limite de tamanho de 8 MB.

  • Solução

    Para pular a entrada de log com tamanho excessivo, altere a posição inicial. Para detalhes, consulte inicialização do job.

OOM no TaskManager: Erro Java heap space durante a recuperação

  • Causa

    Esse problema é tipicamente causado por um corpo de mensagem SLS com tamanho excessivo. O conector SLS solicita dados em lotes. O número de LogGroups por lote é controlado pelo parâmetro batchGetSize, cujo padrão é 100. Isso significa que cada solicitação pode recuperar até 100 LogGroups. Durante a operação normal, o programa Flink consome dados prontamente e raramente recupera um lote completo de 100 LogGroups. No entanto, durante um failover, uma grande quantidade de dados não consumidos pode se acumular. Se a memória necessária para um lote completo de 100 LogGroups exceder a memória disponível da JVM, o TaskManager sofre um OOM.

  • Solução

    Reduza o valor do parâmetro batchGetSize.

Definir o offset de consumo para uma tabela source Paimon

Para definir o offset de consumo de uma tabela source Paimon, utilize o parâmetro scan.mode. A tabela a seguir descreve os valores disponíveis e seus respectivos comportamentos.

Valor

Comportamento de leitura em batch

Comportamento de leitura em stream

default

Valor padrão. O comportamento efetivo depende de outros parâmetros.

  • Se scan.timestamp-millis estiver definido, o comportamento será idêntico ao do valor de parâmetro from-timestamp.

  • Se scan.snapshot-id estiver definido, o comportamento será idêntico ao do valor de parâmetro from-snapshot.

Caso nenhum dos dois parâmetros esteja configurado, o comportamento corresponderá a latest-full.

latest-full

Lê o snapshot mais recente da tabela.

Ao iniciar o job, o sistema lê primeiro o snapshot mais recente da tabela e, em seguida, passa a ler dados incrementais continuamente.

compacted-full

Lê o snapshot mais recente da tabela após a compactação mais recente.

Quando o job é iniciado, ele lê inicialmente o snapshot mais recente da tabela pós-compactação e depois continua lendo dados incrementais.

latest

Equivalente a latest-full.

No início do job, o snapshot mais recente é ignorado; apenas os dados incrementais são lidos de forma contínua.

from-timestamp

Retorna a tabela a partir do snapshot mais recente cuja data seja igual ou anterior a scan.timestamp-millis.

O job não gera um snapshot na inicialização e produz dados incrementais continuamente a partir de scan.timestamp-millis (incluindo este ponto).

from-snapshot

Gera um snapshot da tabela. O ID do snapshot é especificado por scan.snapshot-id.

Nenhum snapshot é produzido na inicialização do job. Em seguida, dados incrementais são gerados continuamente a partir de scan.snapshot-id (inclusive).

from-snapshot-full

Idêntico a from-snapshot.

Durante a inicialização do job, um snapshot da tabela é produzido com o ID especificado por scan.snapshot-id. Posteriormente, o job gera dados incrementais continuamente após o snapshot definido em scan.snapshot-id.

Configurar expiração automática de partições

Tabelas Paimon podem excluir automaticamente partições cujo tempo de vida exceda o período de expiração configurado. Esse recurso ajuda a reduzir custos de armazenamento. O processo funciona da seguinte maneira:

  • Tempo de vida: calculado subtraindo-se o timestamp derivado do valor da partição da hora atual do sistema. A conversão do valor da partição para timestamp ocorre assim:

    1. Conversão do valor da partição em uma string de tempo usando a string de formato definida pelo parâmetro partition.timestamp-pattern.

      Nessa string de formato, uma coluna de partição é representada por um cifrão ($) seguido do nome da coluna. Por exemplo, se as colunas de partição forem ano, mês, dia e hora, a string de formato $year-$month-$day $hour:00:00 converte a partição year=2023,month=04,day=21,hour=17 na string 2023-04-21 17:00:00.

    2. Conversão da string de tempo em um timestamp utilizando a string de formato especificada pelo parâmetro partition.timestamp-formatter.

      Se esse parâmetro não for definido, o sistema usará por padrão os formatos yyyy-MM-dd HH:mm:ss e yyyy-MM-dd. É possível utilizar qualquer string de formato compatível com o DateTimeFormatter do Java.

  • Tempo de expiração da partição: corresponde ao valor configurado no parâmetro partition.expiration-time.

Solucionar problemas de dados ausentes no armazenamento

  • Os dados podem não ficar visíveis imediatamente no armazenamento. Um writer do Flink descarrega dados para o disco nas seguintes condições:

    • Os dados armazenados em buffer em um bucket atingem determinado tamanho (padrão: 64 MB).

    • O tamanho total do buffer atinge um limiar (padrão: 1 GB).

    • Um checkpoint é acionado, descarregando todos os dados presentes na memória.

  • Ao utilizar escrita em stream, certifique-se de que o checkpointing esteja ativado.

Tratar dados duplicados no Hudi

  • Para escritas COW, ative o parâmetro write.insert.drop.duplicates.

    Por padrão, uma escrita COW não remove duplicatas no primeiro arquivo de cada bucket, aplicando a deduplicação apenas aos dados incrementais. Para realizar uma deduplicação global, é necessário ativar esse parâmetro. Em escritas MOR, não são necessários parâmetros adicionais; a definição de uma chave primária já habilita a deduplicação global por padrão.

    Nota

    A partir da versão 0.10.0 do Hudi, esta propriedade foi renomeada para write.precombine e seu valor padrão é true.

  • Para executar a deduplicação entre múltiplas partições, defina o parâmetro index.global.enabled como true.

    Nota
    • Desde a versão 0.10.0 do Hudi, esta propriedade vem configurada como true por padrão.

    • Quando index.type=bucket, definir o parâmetro index.global.enabled como true não surte efeito, pois índices Bucket não suportam alterações entre partições. Consequentemente, mesmo com o índice global ativado, a deduplicação não funcionará em múltiplas partições.

  • Em cenários de atualização de janela longa, como a modificação de dados de um mês atrás, aumente o valor do parâmetro index.state.ttl, medido em dias.

    O índice é a estrutura de dados central do Hudi para identificar duplicatas. O parâmetro index.state.ttl controla por quanto tempo o estado do índice persiste. O valor padrão anterior era de 1,5 dias. Um valor igual ou inferior a 0 indica retenção permanente do estado do índice.

    Nota

    A partir da versão 0.10.0 do Hudi, o valor padrão desta propriedade é 0.

Merge On Read apresenta apenas arquivos de log

  • Causa: o Hudi cria arquivos Parquet somente após a compactação; caso contrário, gera apenas arquivos de log. Por padrão, tabelas Merge On Read utilizam compactação assíncrona, acionando um job de compactação a cada cinco commits.

  • Solução: ajuste o parâmetro compaction.delta_commits para acionar jobs de compactação com maior frequência.

Erro: "multi-statement be found"

  • Problema

    Um job do Flink que grava dados em uma instância do AnalyticDB for MySQL (ADB) falha e reinicia. Os logs exibem um erro semelhante ao seguinte: Caused by: java.sql.SQLSyntaxErrorException: [13000, 2024101216171419216823505703151806929] multi-statement be found.

    at java.util.TimerThread.run(Timer.java:505)
    Caused by: java.sql.BatchUpdateException: [13000, 2024101216400819216823505703151079281] multi-statement be found.
    	at sun.reflect.GeneratedConstructorAccessor115.newInstance(Unknown Source)
    	at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
    	at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
    	at com.mysql.cj.util.Util.handleNewInstance(Util.java:192)
    	at com.mysql.cj.util.Util.getInstance(Util.java:167)
    	at com.mysql.cj.util.Util.getInstance(Util.java:174)
    	at com.mysql.cj.jdbc.exceptions.SQLError.createBatchUpdateException(SQLError.java:224)
    	at com.mysql.cj.jdbc.ClientPreparedStatement.executePreparedBatchAsMultiStatement(ClientPreparedStatement.java:584)
    	at com.mysql.cj.jdbc.ClientPreparedStatement.executeBatchInternal(ClientPreparedStatement.java:431)
    	at com.mysql.cj.jdbc.StatementImpl.executeBatch(StatementImpl.java:795)
    	at com.ververica.cdc.connectors.shaded.com.zaxxer.hikari.pool.ProxyStatement.executeBatch(ProxyStatement.java:127)
    	at com.ververica.cdc.connectors.shaded.com.zaxxer.hikari.pool.HikariProxyPreparedStatement.executeBatch(HikariProxyPreparedStatement...)
    	at com.ververica.connectors.mysql.table.sink.MySqlOutputFormat.executeSql(MySqlOutputFormat.java:567)
    	... 6 more
  • Causa

    Esse erro indica um problema de compatibilidade entre a versão 8.x do driver JDBC do MySQL e um banco de dados AnalyticDB for MySQL (ADB) com ALLOW_MULTI_QUERIES=true ativado.

  • Solução

    1. Entre em contato com o suporte técnico para obter um conector personalizado do ADB 3.0 que utilize a versão 5.1.46 do driver JDBC do MySQL. Aplique este conector à sua tarefa do Flink. Consulte Gerenciar conectores personalizados para obter instruções sobre o uso de conectores personalizados.

    2. Configure o parâmetro allowMultiQueries=true na URI da tabela ADB, por exemplo: jdbc:mysql://xxxxx.ads.aliyuncs.com:3306/xxx?allowMultiQueries=true.

Erro: No suitable driver found

  • Causa

    O conector personalizado não consegue localizar o driver necessário.

  • Solução

Perda ou sobrescrita de dados ao gravar no Elasticsearch com Flink

  • **Causa 1: Conflito entre doc_as_upsert e um Ingest Pipeline do Elasticsearch**

    Quando as configurações abaixo são aplicadas simultaneamente, o Elasticsearch pode processar atualizações parciais e transformações de pipeline em uma ordem incompatível, resultando em sobrescrita inesperada ou perda de documentos:

    • sink.bulk-flush.update.doc_as_upsert = 'true' na cláusula WITH do DDL do Flink

    • Um Ingest Pipeline do Elasticsearch atribuído ao índice de destino

    Nessa configuração, o Elasticsearch aplica o Ingest Pipeline antes de processar a atualização parcial. Dependendo da lógica do pipeline e da versão do Elasticsearch, essa interação pode fazer com que campos do documento sejam sobrescritos com valores incorretos ou que documentos sejam descartados silenciosamente.

  • Causa 2: Múltiplas tarefas sink do Flink gravando no mesmo índice

    Se dois ou mais jobs do Flink — ou várias instâncias sink paralelas com faixas de chaves não sobrepostas — gravarem no mesmo índice do Elasticsearch sem coordenação no tratamento de chaves, os documentos poderão ser sobrescritos mutuamente. Isso causa perda de dados para o sink que gravar por último em um determinado ID de documento.

  • Soluções

    • **Para o conflito entre doc_as_upsert e Ingest Pipeline:** Remova o parâmetro sink.bulk-flush.update.doc_as_upsert = 'true' do DDL do Flink e remova ou reatribua o Ingest Pipeline do Elasticsearch do índice de destino. Transfira qualquer lógica de transformação de dados que era tratada pelo pipeline para o próprio job do Flink — por exemplo, usando uma ProcessFunction ou colunas computadas antes do operador sink.

    • Para múltiplas tarefas sink sobrescrevendo dados: Garanta que cada job do Flink grave em um índice separado do Elasticsearch ou consolide as gravações em um único job do Flink que gerencie as chaves dos documentos de forma determinística.