Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Gravar e consumir dados em tabelas Paimon

Última atualização: Jun 27, 2026

Este tópico descreve como inserir, atualizar, sobrescrever ou excluir dados em tabelas Paimon pelo console do Realtime Compute for Apache Flink. Também explica como consumir dados dessas tabelas e especificar offsets de consumo.

Pré-requisitos

Crie previamente um catálogo Paimon e uma tabela Paimon. Para mais informações, consulte Gerenciar catálogos Paimon.

Limitações

As tabelas Paimon são suportadas apenas no Ververica Runtime (VVR) 8.0.5 e versões posteriores.

Gravar dados em uma tabela Paimon

Sincronizar dados e esquema com CTAS/CDAS

Para mais informações, consulte Gerenciar catálogos Paimon.

Inserir ou atualizar dados com INSERT INTO

Use a instrução INSERT INTO para inserir ou atualizar dados diretamente em uma tabela Paimon.

Sobrescrever dados com INSERT OVERWRITE

A sobrescrita limpa os dados existentes e grava novos dados. Use a instrução INSERT OVERWRITE para sobrescrever uma tabela Paimon inteira ou uma partição específica. Os exemplos a seguir ilustram esse uso.

Nota
  • Sobrescreva toda a tabela não particionada my_table.

    INSERT OVERWRITE my_table SELECT ...;
  • Sobrescreva a partição dt=20240108,hh=06 na tabela my_table.

    INSERT OVERWRITE my_table PARTITION (`dt` = '20240108', `hh` = '06') SELECT ...;
  • Sobrescreva partições dinamicamente na tabela my_table. As partições presentes no resultado da instrução SELECT são sobrescritas; as demais permanecem inalteradas.

    INSERT OVERWRITE my_table SELECT ...;
  • Sobrescreva toda a tabela particionada my_table.

    INSERT OVERWRITE my_table /*+ OPTIONS('dynamic-partition-overwrite' = 'false') */ SELECT ...;

Excluir dados usando DELETE

Use a instrução DELETE para remover dados de uma tabela de chave primária Paimon. Execute instruções DELETE somente na Exploração de Dados.

-- Delete all data where currency = 'UNKNOWN' from the my_table table.
DELETE FROM my_table WHERE currency = 'UNKNOWN';

Filtrar mensagens de exclusão

Em tabelas de chave primária Paimon, as mensagens DELETE removem, por padrão, os dados associados à chave primária correspondente. Para impedir que a tabela processe essas mensagens, use uma SQL hint e defina o parâmetro abaixo como true. Isso filtra as mensagens de exclusão.

Parâmetro

Descrição

Tipo

Padrão

ignore-delete

Define se as mensagens de exclusão devem ser filtradas.

Booleano

false

Ajustar o paralelismo do sink

Ajuste o paralelismo do operador de sink definindo o parâmetro a seguir por meio de uma SQL hint.

Parâmetro

Descrição

Tipo

Padrão

sink.parallelism

Define o paralelismo do operador de sink Paimon.

Inteiro

Nenhum

Por exemplo, a instrução SQL abaixo define o paralelismo do operador de sink Paimon como 10.

INSERT INTO t /*+ OPTIONS('sink.parallelism' = '10') */ SELECT * FROM s;

Consumir dados de uma tabela Paimon

Jobs de streaming

Nota

Para tabelas de chave primária Paimon consumidas por jobs de streaming, configure um produtor de changelog.

Por padrão, um operador de source Paimon em um job de streaming produz primeiro todos os dados da tabela no início do job e, em seguida, os dados incrementais a partir desse momento.

Consumir dados a partir de um offset específico

Consuma dados a partir de um offset definido das seguintes formas:

  • Caso não seja necessário consumir os dados completos da tabela Paimon na inicialização do job, mas apenas os dados incrementais subsequentes, defina 'scan.mode' = 'latest' usando uma SQL hint.

    SELECT * FROM t /*+ OPTIONS('scan.mode' = 'latest') */;
  • Para ignorar os dados completos e consumir apenas dados incrementais a partir de um momento específico, use uma SQL hint para definir o parâmetro scan.timestamp-millis. O valor desse parâmetro representa o número de milissegundos desde a Unix Epoch (1970-01-01 00:00:00 UTC) até o momento desejado.

    SELECT * FROM t /*+ OPTIONS('scan.timestamp-millis' = '1678883047356') */;
  • Para consumir dados gravados após um horário específico e continuar consumindo os dados incrementais seguintes, adote um dos métodos abaixo.

    Nota

    Esse método de consumo lê arquivos de dados modificados após o horário especificado. Devido à compactação, esses arquivos podem conter uma pequena quantidade de dados gravados antes do horário definido. Adicione uma cláusula WHERE ao seu job SQL para filtrar os dados conforme necessário.

    • Não configure nenhuma SQL hint. Ao iniciar o job, selecione Specify source's start time. Na caixa de diálogo Job Start, escolha Stateless start-up, ative a opção Specify source's start time e defina o horário alvo.

    • Defina o parâmetro scan.file-creation-time-millis por meio de uma SQL hint.

      SELECT * FROM t /*+ OPTIONS('scan.file-creation-time-millis' = '1678883047356') */;
  • Quando não for necessário consumir os dados completos e você quiser processar apenas dados incrementais a partir de um arquivo de snapshot específico, use uma SQL hint para definir o parâmetro scan.snapshot-id. O valor corresponde ao ID do arquivo de snapshot desejado.

    SELECT * FROM t /*+ OPTIONS('scan.snapshot-id' = '3') */;
  • Se o objetivo for consumir todos os dados de um arquivo de snapshot específico e depois continuar com os dados incrementais, utilize uma SQL hint para definir os parâmetros 'scan.mode' = 'from-snapshot-full' e scan.snapshot-id. O valor de scan.snapshot-id deve ser o ID do arquivo de snapshot escolhido.

    SELECT * FROM t /*+ OPTIONS('scan.mode' = 'from-snapshot-full', 'scan.snapshot-id' = '1') */;

Especificar um consumer ID

Um consumer ID armazena o progresso de consumo de uma tabela Paimon. Seu uso é indicado principalmente nos cenários a seguir:

  • Ao definir um consumer ID, o progresso de consumo correspondente é salvo no arquivo de metadados da tabela Paimon. Isso permite que o job retome o consumo do ponto de interrupção, mesmo que seja reiniciado posteriormente em modo stateless.

  • Após a definição de um consumer ID, snapshots ainda não consumidos não são excluídos por expiração. Isso evita erros causados quando a velocidade de consumo não acompanha a velocidade de expiração dos snapshots.

Defina o parâmetro consumer-id para atribuir um Consumer ID a um operador de source Paimon em um job de streaming. O valor do Consumer ID pode ser qualquer string. Na primeira criação de um Consumer ID, o offset inicial segue as regras descritas em Consumir uma tabela Paimon a partir de um offset especificado. Posteriormente, basta reutilizar o mesmo Consumer ID para retomar o consumo da tabela Paimon.

Por exemplo, a instrução SQL abaixo define um consumer ID chamado test-id para um operador de source Paimon. Para redefinir o offset de um consumer ID específico, também é possível definir 'consumer.ignore-progress' = 'true'.

SELECT * FROM t /*+ OPTIONS('consumer-id' = 'test-id') */;
Nota

Arquivos de snapshot não consumidos por um consumer ID não são excluídos ao expirar. Se consumer IDs obsoletos não forem limpos, os arquivos de snapshot e os arquivos de dados históricos correspondentes nunca serão removidos, ocupando espaço de armazenamento indefinidamente. Defina o parâmetro de tabela consumer.expiration-time para limpar consumer IDs inativos por um período determinado. Por exemplo, 'consumer.expiration-time' = '3d' indica que consumer IDs sem uso por 3 dias serão removidos.

Transmitir e consumir resultados de INSERT OVERWRITE

Por padrão, operações INSERT OVERWRITE não geram dados de changelog, impedindo que jobs de streaming downstream consumam os dados excluídos e importados. Caso precise consumir esses dados, defina 'streaming-read-overwrite' = 'true' em um job de consumo de streaming usando uma SQL hint.

SELECT * FROM t /*+ OPTIONS('streaming-read-overwrite' = 'true') */;

Jobs em lote

Em jobs em lote, o operador de source Paimon lê, por padrão, o snapshot mais recente e retorna o estado atualizado da tabela Paimon.

Time travel em lote

Consulte o estado de uma tabela Paimon em um momento específico definindo o parâmetro scan.timestamp-millis em uma SQL hint. O valor representa o número de milissegundos decorridos desde a Unix Epoch (1970-01-01 00:00:00 UTC) até o instante desejado.

SELECT * FROM t /*+ OPTIONS('scan.timestamp-millis' = '1678883047356') */;

Também é possível consultar o estado da tabela Paimon no momento da criação de um snapshot. Para isso, use uma SQL hint e defina o parâmetro scan.snapshot-id com o ID do snapshot desejado.

SELECT * FROM t /*+ OPTIONS('scan.snapshot-id' = '3') */;

Consultar alterações entre snapshots

Para consultar alterações de dados em uma tabela Paimon entre dois snapshots, defina o parâmetro incremental-between usando uma SQL hint. Por exemplo, para visualizar todas as alterações ocorridas entre o snapshot 12 e o snapshot 20, utilize a seguinte instrução SQL.

SELECT * FROM t /*+ OPTIONS('incremental-between' = '12.20') */;
Nota

Como jobs em lote não suportam o consumo de mensagens DELETE, essas mensagens são descartadas por padrão. Para consumir mensagens DELETE em um job em lote, consulte a tabela de sistema Audit Log. Exemplo: SELECT * FROM .

Ajustar o paralelismo do source

Por padrão, o Paimon infere automaticamente o paralelismo do operador de source com base em informações como o número de partições e buckets. Utilize uma SQL hint para definir os parâmetros abaixo e ajustar o paralelismo manualmente.

Parâmetro

Tipo

Padrão

Descrição

scan.parallelism

Inteiro

Nenhum

Define o paralelismo do operador de source Paimon.

scan.infer-parallelism

Booleano

true

Indica se o paralelismo do operador de source Paimon deve ser inferido automaticamente.

scan.infer-parallelism.max

Inteiro

1024

Limite superior para o paralelismo inferido automaticamente pelo operador de source Paimon.

A instrução SQL a seguir exemplifica como definir o paralelismo do operador de source Paimon como 10.

SELECT * FROM t /*+ OPTIONS('scan.parallelism' = '10') */;

Usar tabelas Paimon como tabelas de dimensão

Tabelas Paimon também funcionam como tabelas de dimensão. Para detalhes sobre a sintaxe de JOINs com tabelas de dimensão, consulte Instrução JOIN de tabela de dimensão.

Gravar e consumir o tipo VARIANT

No Ververica Runtime (VVR) 11,1 e versões posteriores, as tabelas Paimon suportam o tipo de dados semiestruturado VARIANT. Esse tipo permite converter strings JSON VARCHAR para o tipo VARIANT usando PARSE_JSON ou TRY_PARSE_JSON. Gravar e consumir diretamente o tipo VARIANT melhora significativamente o desempenho de consultas e processamento de JSON.

O código abaixo apresenta um exemplo:

CREATE TABLE `my-catalog`.`my_db`.`my_tbl` (
  k BIGINT,
  info VARIANT
);
INSERT INTO `my-catalog`.`my_db`.`my_tbl` 
SELECT k, PARSE_JSON(jsonStr) FROM T;

Documentação relacionada