Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Tabelas de chave primária e somente de anexação

Última atualização: Jun 27, 2026

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 deduplicate

Deduplicação em que apenas o primeiro registro é relevante

Tabela de chave primária, mecanismo first-row

Agregação em tempo real (totais acumulados, rastreamento de máx/mín)

Tabela de chave primária, mecanismo aggregation

Montagem de tabela ampla a partir de múltiplas fontes de streaming

Tabela de chave primária, mecanismo partial-update

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 bucket ou defina 'bucket' = '-1'

Não suporta gravações simultâneas de múltiplos deployments do Flink; suporta atualizações entre partições

Modo de bucket fixo

Defina 'bucket' = '<num>' (inteiro > 0)

<num> representa a quantidade de buckets para toda a tabela (não particionada) ou por partição (particionada); permite alterar a contagem de buckets

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.

  • aggregation ou partial-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

dynamic-bucket.target-row-num

Máximo de registros por bucket

2000000

dynamic-bucket.initial-buckets

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:

  1. Pause todos os deployments que leem ou gravam na tabela.

  2. 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>');
  3. 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' e hh = '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';
  4. 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 .

deduplicate (padrão)

O Paimon mantém apenas o registro mais recente para cada conjunto de chaves primárias e descarta os anteriores. Se o registro mais recente for uma exclusão (DELETE), o sistema remove todos os registros com essas chaves primárias.

CREATE TABLE T (
  k INT,
  v1 DOUBLE,
  v2 STRING,
  PRIMARY KEY (k) NOT ENFORCED
) WITH (
  'merge-engine' = 'deduplicate' -- Optional: deduplicate is the default
);

Resultados do exemplo:

  • Grave +I(1, 2.0, 'apple'), +I(1, 4.0, 'banana'), +I(1, 8.0, 'cherry') — a consulta retorna (1, 8.0, 'cherry').

  • Grave +I(1, 2.0, 'apple'), +I(1, 4.0, 'banana'), -D(1, 4.0, 'banana') — a consulta não retorna linhas.

Indicado para: Sincronização CDC, pipelines de upsert padrão.

first-row

O Paimon mantém apenas o primeiro registro para cada conjunto de chaves primárias. Como os consumidores downstream veem apenas eventos de changelog do tipo INSERT, a produção de changelog é mais eficiente do que com deduplicate.

CREATE TABLE T (
  k INT,
  v1 DOUBLE,
  v2 STRING,
  PRIMARY KEY (k) NOT ENFORCED
) WITH (
  'merge-engine' = 'first-row'
);

Resultado do exemplo:

Grave +I(1, 2.0, 'apple'), +I(1, 4.0, 'banana'), +I(1, 8.0, 'cherry') — a consulta retorna (1, 2.0, 'apple').

Notas de uso:

  • Defina changelog-producer como lookup para habilitar o consumo em streaming.

  • O mecanismo first-row não processa alterações DELETE e UPDATE_BEFORE. Para ignorá-las em vez de gerar erro, defina 'first-row.ignore-delete' = 'true'.

  • Campos de sequência não são suportados.

Indicado para: Pipelines de deduplicação onde apenas o primeiro registro recebido é relevante.

aggregation

O Paimon aplica uma função de agregação a cada coluna que não seja chave primária nos registros que compartilham as mesmas chaves primárias. Especifique a função por coluna usando fields.<field-name>.aggregate-function. Colunas sem uma função explícita usam last_non_null_value como padrão.

CREATE TABLE T (
  product_id BIGINT,
  price DOUBLE,
  sales BIGINT,
  PRIMARY KEY (product_id) NOT ENFORCED
) WITH (
  'merge-engine' = 'aggregation',
  'fields.price.aggregate-function' = 'max',
  'fields.sales.aggregate-function' = 'sum'
);

Resultado do exemplo:

Grave +I(1, 23.0, 15) e +I(1, 30.2, 20) — a consulta retorna (1, 30.2, 35).

Funções de agregação suportadas:

Função

Tipos suportados

sum

DECIMAL, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT, DOUBLE

product

DECIMAL, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT, DOUBLE

count

INTEGER, BIGINT

max, min

CHAR, VARCHAR, DECIMAL, TINYINT, SMALLINT, INTEGER, BIGINT, FLOAT, DOUBLE, DATE, TIME, TIMESTAMP, TIMESTAMP_LTZ

first_value, last_value

Todos os tipos, incluindo nulo

first_not_null_value, last_non_null_value

Todos os tipos

listagg

STRING

bool_and, bool_or

BOOLEAN

Apenas sum , product e count suportam alterações de retração (UPDATE_BEFORE e DELETE). Para ignorar alterações de retração em uma coluna específica, defina 'fields.<field-name>.ignore-retract' = 'true' .

Defina changelog-producer como lookup ou full-compaction para habilitar o consumo em streaming.

Indicado para: Métricas de agregação em tempo real, totais acumulados, rastreamento de máx/mín.

partial-update

O Paimon atualiza colunas individuais de um registro existente usando o valor não nulo mais recente dos registros recebidos que compartilham as mesmas chaves primárias. Valores nulos não sobrescrevem valores existentes.

CREATE TABLE T (
  k INT,
  v1 DOUBLE,
  v2 BIGINT,
  v3 STRING,
  PRIMARY KEY (k) NOT ENFORCED
) WITH (
  'merge-engine' = 'partial-update'
);

Resultado do exemplo:

Grave +I(1, 23.0, 10, NULL), +I(1, NULL, NULL, 'This is a book'), +I(1, 25.2, NULL, NULL) — a consulta retorna (1, 25.2, 10, 'This is a book').

Notas de uso:

  • Defina changelog-producer como lookup ou full-compaction para habilitar o consumo em streaming.

  • O mecanismo partial-update não processa alterações DELETE e UPDATE_BEFORE. Para ignorá-las, defina 'partial-update.ignore-delete' = 'true'.

Grupos de sequência

Ao montar uma tabela ampla a partir de múltiplas tabelas upstream, use grupos de sequência para controlar independentemente a ordem de atualização de diferentes conjuntos de colunas e lidar com dados fora de ordem entre fontes.

No exemplo a seguir, as colunas a e b são atualizadas em ordem crescente de g_1, e as colunas c e d são atualizadas em ordem crescente de g_2:

CREATE TABLE T (
  k INT,
  a STRING,
  b STRING,
  g_1 INT,
  c STRING,
  d STRING,
  g_2 INT,
  PRIMARY KEY (k) NOT ENFORCED
) WITH (
  'merge-engine' = 'partial-update',
  'fields.g_1.sequence-group' = 'a,b',
  'fields.g_2.sequence-group' = 'c,d'
);

Combine grupos de sequência com funções de agregação adicionando fields.<field-name>.aggregate-function para qualquer coluna no grupo. O exemplo a seguir rastreia o valor não nulo mais recente de a, o máximo de b, o valor não nulo mais recente de c e a soma de d:

CREATE TABLE T (
  k INT,
  a STRING,
  b INT,
  g_1 INT,
  c STRING,
  d INT,
  g_2 INT,
  PRIMARY KEY (k) NOT ENFORCED
) WITH (
  'merge-engine' = 'partial-update',
  'fields.g_1.sequence-group' = 'a,b',
  'fields.b.aggregate-function' = 'max',
  'fields.g_2.sequence-group' = 'c,d',
  'fields.d.aggregate-function' = 'sum'
);

Para mais informações, consulte Mecanismo de Mesclagem.

Indicado para: Montagem de tabela ampla unindo colunas de múltiplas fontes de streaming.

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 metadados op_t como 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

none

Nenhum changelog é gerado.

Mais baixo

Quando não há necessidade de consumo em streaming

input

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.

lookup (recomendado)

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

full-compaction

Gera um changelog após cada compactação completa de arquivos pequenos.

Até várias horas

Menor que lookup

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 a lookup e full-compaction e 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

'bucket' = '-1'

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

'bucket' = '<num>' (inteiro > 0)

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