O Apache Paimon oferece suporte a dois tipos de tabela: tabelas de chave primária e tabelas somente de anexação. Cada tipo possui semânticas distintas para gravações, atualizações e consumo em streaming. Escolha o tipo que melhor atende aos requisitos do seu pipeline de dados.
Escolha um tipo de tabela
|
Cenário |
Configuração recomendada |
|
Sincronização CDC, upserts padrão |
Tabela de chave primária, mecanismo |
|
Deduplicação em que apenas o primeiro registro é relevante |
Tabela de chave primária, mecanismo |
|
Agregação em tempo real (totais acumulados, rastreamento de máx/mín) |
Tabela de chave primária, mecanismo |
|
Montagem de tabela ampla a partir de múltiplas fontes de streaming |
Tabela de chave primária, mecanismo |
|
Ingestão de logs, streaming de eventos |
Tabela somente de anexação (escalável) |
|
Substituição de fila de mensagens com preservação de ordem |
Tabela de fila de anexação |
Tabelas de chave primária
Uma tabela de chave primária exige uma ou mais chaves primárias, definidas na criação da tabela. O sistema mescla registros com as mesmas chaves primárias conforme o mecanismo de mesclagem configurado.
CREATE TABLE T (
dt STRING,
shop_id BIGINT,
user_id BIGINT,
num_orders INT,
total_amount INT,
PRIMARY KEY (dt, shop_id, user_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'bucket' = '4'
);
Neste exemplo, a chave de partição é dt, as chaves primárias são dt, shop_id e user_id, e a tabela utiliza 4 buckets fixos.
Modos de bucket
O bucket é a menor unidade de leitura e gravação em uma tabela Paimon. As partições (ou a tabela inteira, no caso de tabelas não particionadas) se subdividem em buckets, permitindo leituras e gravações paralelas.
|
Modo |
Configuração |
Observações |
|
Modo de bucket dinâmico (padrão) |
Omita |
Não suporta gravações simultâneas de múltiplos deployments do Flink; suporta atualizações entre partições |
|
Modo de bucket fixo |
Defina |
|
Dimensionamento de buckets: Recomendamos definir o tamanho total dos dados em cada bucket como 2 GB e evitar valores superiores a 5 GB. Poucos buckets limitam o paralelismo de gravação; muitos buckets geram excesso de arquivos pequenos.
Modo de bucket dinâmico: atualização de dados
As atualizações entre partições ocorrem quando as chaves primárias não incluem todas as chaves de partição. O Paimon não consegue determinar a partição de destino apenas pela chave primária, portanto utiliza o RocksDB para manter um mapeamento de chave primária para partição/bucket. Em tabelas grandes, isso pode reduzir significativamente o desempenho em comparação ao modo de bucket fixo, e a inicialização do deployment demora mais enquanto o mapeamento carrega no RocksDB.
O comportamento de atualização depende do mecanismo de mesclagem:
deduplicate: exclui o registro existente e insere um novo registro na partição de destino.aggregationoupartial-update: atualiza o registro existente diretamente na partição atual.first-row: mantém o registro existente e descarta o registro recebido.
As atualizações dentro da partição ocorrem quando as chaves primárias incluem todas as chaves de partição. O Paimon determina a partição a partir da chave primária, mas não o bucket, então mantém um índice em memória que mapeia chaves primárias para buckets. A cada 100 milhões de entradas de mapeamento, o sistema consome aproximadamente 1 GB de memória heap, cobrados apenas das partições em gravação ativa.
Modo de bucket dinâmico: atribuição de buckets
O sistema grava novos registros primeiro nos buckets existentes. Se todos os buckets atingirem a capacidade máxima, ele cria automaticamente um novo bucket.
Configure esse comportamento com os seguintes parâmetros na cláusula WITH:
|
Parâmetro |
Descrição |
Padrão |
|
|
Máximo de registros por bucket |
|
|
|
Contagem inicial de buckets |
Paralelismo do operador de gravação |
Modo de bucket fixo: atribuição de buckets
Por padrão, o Paimon atribui registros aos buckets usando um hash dos valores da chave primária. Para usar um conjunto diferente de colunas, defina bucket-key na cláusula WITH. Especifique várias colunas separadas por vírgulas. As colunas em bucket-key devem ser um subconjunto das chaves primárias.
Por exemplo, 'bucket-key' = 'c1,c2' roteia registros pelo hash de c1 e c2.
No modo de bucket fixo, inclua todas as chaves de partição nas chaves primárias para evitar atualizações entre partições.
Alterar o número de buckets no modo de bucket fixo
A contagem de buckets determina o paralelismo de leitura e gravação. Para alterá-la em uma tabela existente:
Pause todos os deployments que leem ou gravam na tabela.
-
Crie um script e execute a seguinte instrução SQL para configurar o parâmetro de bucket:
ALTER TABLE `<catalog-name>`.`<database-name>`.`<table-name>` SET ('bucket' = '<bucket-num>'); -
Reorganize os dados executando um deployment em lote.
Tabela não particionada: Crie um Blank Batch Draft, cole o SQL abaixo, clique em Deploy e Start: ``
sql INSERT OVERWRITE<catalog-name>.<database-name>.<table-name>SELECT * FROM<catalog-name>.<database-name>.<table-name>;``-
Tabela particionada: Execute um lote por partição que deseja reorganizar:
INSERT OVERWRITE `<catalog-name>`.`<database-name>`.`<table-name>` PARTITION (<partition-spec>) SELECT * FROM `<catalog-name>`.`<database-name>`.`<table-name>` WHERE <partition-condition>;Exemplo — reorganizar a partição onde
dt = '20240312'ehh = '08':INSERT OVERWRITE `<catalog-name>`.`<database-name>`.`<table-name>` PARTITION (dt = '20240312', hh = '08') SELECT * FROM `<catalog-name>`.`<database-name>`.`<table-name>` WHERE dt = '20240312' AND hh = '08';
Após a conclusão bem-sucedida do deployment em lote, retome os deployments pausados.
Mecanismos de mesclagem
Quando vários registros compartilham as mesmas chaves primárias, o Paimon os mescla conforme a configuração merge-engine.
Se seu fluxo de entrada contiver registros fora de ordem, configure o parâmetro sequence.field para controlar a ordem de mesclagem. Consulte Tratamento de dados fora de ordem .
Tratamento de dados fora de ordem
Por padrão, o Paimon mescla registros na ordem de entrada — o último registro recebido é o último a ser mesclado. Para fluxos com registros fora de ordem, defina sequence.field na cláusula WITH. O sistema mesclará registros com as mesmas chaves primárias em ordem crescente do valor da coluna especificada.
Tipos de dados suportados para sequence.field: TINYINT, SMALLINT, INTEGER, BIGINT, TIMESTAMP, TIMESTAMP_LTZ.
Ao usar MySQL como fonte de entrada com a coluna de metadadosop_tcomo campo de sequência, os pares de alteração UPDATE_BEFORE e UPDATE_AFTER compartilham o mesmo valor de sequência. Para garantir que o Paimon processe UPDATE_BEFORE antes de UPDATE_AFTER, defina'sequence.auto-padding' = 'row-kind-flag'.
Produtor de changelog
Uma tabela de chave primária deve gerar um changelog completo (cobrindo operações INSERT, DELETE e UPDATE) para suportar o consumo de streaming downstream, semelhante a um binlog de banco de dados. Configure o método de geração com o parâmetro changelog-producer.
|
Valor |
Descrição |
Latência |
Uso de recursos |
Indicado para |
|
|
Nenhum changelog é gerado. |
— |
Mais baixo |
Quando não há necessidade de consumo em streaming |
|
|
Encaminha os registros de entrada diretamente para os consumidores downstream. |
Mais baixa |
Mais baixo |
Quando a entrada já contém um changelog completo, como um binlog de banco de dados. Opção mais eficiente — sem computação extra. |
|
|
Executa uma busca nos resultados de compactação de arquivos pequenos para produzir um changelog completo. Acionado a cada checkpoint do Flink. |
Nível de minutos |
Mais alto |
Qualquer tipo de entrada; requisitos de latência no nível de minutos |
|
|
Gera um changelog após cada compactação completa de arquivos pequenos. |
Até várias horas |
Menor que |
Qualquer tipo de entrada; latência maior aceitável |
Para full-compaction, defina 'full-compaction.delta-commits' = '<num>' para acionar a compactação completa a cada <num> checkpoints do Flink. A compactação completa consome muitos recursos — defina o intervalo entre 30 minutos e 1 hora.
Por padrão, o Paimon gera um registro de changelog mesmo quando o valor atualizado é idêntico ao anterior. Para suprimir esses registros sem efeito, defina'changelog-producer.row-deduplicate' = 'true'. Esta opção aplica-se apenas alookupefull-compactione requer computação adicional para comparar valores antes/depois. Habilite-a apenas quando for esperado um grande número de registros desnecessários.
Tabelas somente de anexação
Uma tabela somente de anexação não possui chaves primárias e aceita apenas operações INSERT no modo streaming. Utilize este tipo de tabela para cargas de trabalho que não exigem atualizações em streaming, como sincronização de dados de log.
CREATE TABLE T (
dt STRING,
order_id BIGINT,
item_id BIGINT,
amount INT,
address STRING
) PARTITIONED BY (dt) WITH (
'bucket' = '-1'
);
Subtipos
|
Subtipo |
Configuração |
Características |
|
Tabela escalável de anexação |
|
Alto throughput de gravação; sem necessidade de particionamento por hash; suporta ordenação de dados; configuração flexível de paralelismo; suporta conversão direta de tabelas Hive; compactação de arquivos totalmente assíncrona; ordem de consumo difere da ordem de gravação; possível skew de dados quando o paralelismo upstream é igual ao paralelismo de gravação |
|
Tabela de fila de anexação |
|
Preserva ordem por bucket; latência de vários minutos; equivalente a uma partição de tópico Kafka ou ao número de shards em uma instância ApsaraMQ for MQTT |
Atribuição de buckets
Tabela escalável de anexação: O sistema grava os dados diretamente em uma única partição em paralelo. O conceito de bucket é ignorado e não há necessidade de particionamento por hash, o que maximiza o desempenho de gravação.
Tabela de fila de anexação: Por padrão, o sistema atribui registros aos buckets com base nos valores de todas as colunas do registro de dados. Para controlar a atribuição, defina bucket-key na cláusula WITH com uma lista de nomes de colunas separados por vírgulas.
Por exemplo, 'bucket-key' = 'c1,c2' roteia registros pelo hash de c1 e c2.
Especifique bucket-key para reduzir a computação de atribuição de buckets e melhorar a eficiência de gravação.
Ordem de consumo de dados
Tabela escalável de anexação: Os registros podem ser consumidos em uma ordem diferente da ordem de gravação.
Tabela de fila de anexação: Os registros dentro de cada bucket são consumidos na ordem de gravação. Para registros em diferentes partições ou buckets:
Dois registros de partições diferentes: se
'scan.plan-sort-partition' = 'true', o registro com o menor valor de partição é consumido primeiro; caso contrário, o registro na partição criada anteriormente é consumido primeiro.Dois registros da mesma partição, mesmo bucket: o registro gravado anteriormente é consumido primeiro.
Dois registros da mesma partição, buckets diferentes: a ordem de consumo não é garantida porque buckets diferentes podem ser processados simultaneamente.
Próximos passos
Para criar um catálogo e uma tabela Paimon, consulte Gerenciar catálogos Paimon.
Para otimizar o desempenho de tabelas de chave primária, consulte Otimização de desempenho.