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.
Uma tabela de chave primária Paimon aceita mensagens de todos os tipos, incluindo INSERT, UPDATE_BEFORE, UPDATE_AFTER e DELETE. Durante a gravação, ela mescla dados com a mesma chave primária conforme o mecanismo de mesclagem de dados.
Já uma tabela append-only Paimon, também chamada de tabela sem chave primária, aceita exclusivamente mensagens do tipo INSERT.
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.
A instrução
INSERT OVERWRITEé suportada apenas em jobs em lote.Por padrão, uma operação
INSERT OVERWRITEnão gera dados de changelog. Jobs de streaming downstream não conseguem consumir os dados excluídos e importados. Caso precise consumir esse tipo de dado, consulte Transmitir e consumir os resultados de uma instrução INSERT OVERWRITE.
-
Sobrescreva toda a tabela não particionada
my_table.INSERT OVERWRITE my_table SELECT ...; -
Sobrescreva a partição
dt=20240108,hh=06na 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çãoSELECTsã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
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.
NotaEsse 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
WHEREao 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-millispor 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'escan.snapshot-id. O valor descan.snapshot-iddeve 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') */;
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') */;
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
Durante a gravação e o consumo de dados em tabelas Paimon, utilize SQL hints para modificar temporariamente os parâmetros da tabela. Para mais informações, consulte Gerenciar tabelas Paimon.
Para conhecer os recursos e funções básicos das tabelas de chave primária Paimon e das tabelas append-only Paimon, consulte Tabelas de chave primária e append-only Paimon.
Saiba mais sobre otimizações comuns para tabelas de chave primária Paimon e tabelas Append Scalable em diferentes cenários em Otimização de desempenho do Paimon.
O consumo de dados de tabelas Paimon depende de arquivos de snapshot. Se o tempo de expiração do snapshot for muito curto ou o job de consumo apresentar baixa eficiência, o arquivo de snapshot em uso pode ser excluído após a expiração, gerando o erro
File xxx not found, Possible causesno job. Para resolver esse problema, consulte Resolver o erro "File xxx not found, Possible causes" em jobs de leitura Paimon.