Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Upsert Kafka

Última atualização: Jun 27, 2026

O conector Upsert Kafka lê e grava em tópicos do Kafka usando semântica de upsert. Como source, converte registros do Kafka em um stream de changelog. Como sink, grava dados de changelog de volta no Kafka e utiliza mensagens com valor nulo (tombstone) para representar exclusões.

Categoria

Detalhes

Tipos suportados

Tabela source, tabela sink e sink de ingestão de dados

Modo de execução

Modo streaming

Formatos de dados

avro, avro-confluent, csv, json e raw

Tipos de API

SQL e job YAML de ingestão de dados

Atualizar ou excluir dados em uma tabela sink

Sim

Pré-requisitos

Antes de começar, verifique se você tem:

Funcionamento

O stream de changelog upsert produzido pela source contém apenas entradas +I (inserção), +U (atualização posterior) e -D (exclusão). Diferentemente de um stream de retração, nunca contém mensagens -U (atualização anterior).

Tipo de entrada de changelog

Nome longo

Semântica

+I

Inserção

Novo registro cuja chave não existia anteriormente

+U

Atualização posterior

Registro cuja chave já existe — sobrescreve o valor anterior

-D

Exclusão

Registro com valor nulo — exclui a entrada correspondente àquela chave

Como tabela source, o conector lê registros do Kafka e os converte em um stream de changelog:

  • Registro com chave existente é interpretado como UPDATE (+U)

  • Registro com chave inexistente é interpretado como INSERT (+I)

  • Registro com valor nulo é interpretado como DELETE (-D)

Como tabela sink ou sink de ingestão de dados, o conector consome um stream de changelog upstream:

  • Registros INSERT e UPDATE_AFTER são gravados como mensagens normais do Kafka

  • Registros DELETE são gravados como mensagens tombstone (valor nulo para a chave correspondente)

O Flink particiona a saída do sink por chave primária. Assim, todas as mensagens da mesma chave chegam à mesma partição na ordem correta.

Exemplo

O exemplo completo a seguir lê eventos de visualização de página de um site a partir de uma tabela source Kafka, agrega visualizações de página (PV) e visitantes únicos (UV) por região e grava os resultados em um sink Upsert Kafka.

Tabela source — lê eventos brutos de visualização de página:

CREATE TABLE pageviews (
  user_id    BIGINT,
  page_id    BIGINT,
  viewtime   TIMESTAMP,
  user_region STRING,
  WATERMARK FOR viewtime AS viewtime - INTERVAL '2' SECOND
) WITH (
  'connector'                  = 'kafka',
  'topic'                      = '<yourTopicName>',
  'properties.bootstrap.servers' = '<host:port>',
  'format'                     = 'json'
);

Tabela sink — armazena resultados agregados com semântica de upsert:

CREATE TABLE pageviews_per_region (
  user_region STRING,
  pv          BIGINT,
  uv          BIGINT,
  PRIMARY KEY (user_region) NOT ENFORCED
) WITH (
  'connector'                  = 'upsert-kafka',
  'topic'                      = '<yourTopicName>',
  'properties.bootstrap.servers' = '<host:port>',
  'key.format'                 = 'avro',
  'value.format'               = 'avro'
);
Toda tabela Upsert Kafka deve definir uma chave primária. As colunas da chave primária determinam como o Flink particiona a saída e como o conector identifica operações de upsert.

Consulta — agrega e grava resultados:

INSERT INTO pageviews_per_region
SELECT
  user_region,
  COUNT(*),
  COUNT(DISTINCT user_id)
FROM pageviews
GROUP BY user_region;

SQL

Sintaxe

CREATE TABLE upsert_kafka_sink (
  user_region STRING,
  pv          BIGINT,
  uv          BIGINT,
  PRIMARY KEY (user_region) NOT ENFORCED
) WITH (
  'connector'                  = 'upsert-kafka',
  'topic'                      = '<yourTopicName>',
  'properties.bootstrap.servers' = '...',
  'key.format'                 = 'avro',
  'value.format'               = 'avro'
);

Parâmetros WITH

Parâmetros gerais

Parâmetro

Obrigatório

Padrão

Descrição

connector

Sim

Defina como upsert-kafka.

properties.bootstrap.servers

Sim

Lista separada por vírgulas de endereços de brokers Kafka no formato host:port.

topic

Sim

Tópico Kafka para leitura ou gravação.

key.format

Sim

Formato de serialização para a chave da mensagem. Valores válidos: csv, json, avro, debezium-json, canal-json, maxwell-json, avro-confluent, raw.

value.format

Sim

Formato de serialização para o valor da mensagem. Equivalente a format — configure apenas um dos dois.

value.fields-include

Sim

ALL

Controla quais colunas são incluídas no valor da mensagem. ALL inclui todas as colunas; EXCEPT_KEY exclui os campos de chave definidos por key.fields.

properties.*

Não

Parâmetros adicionais do cliente Kafka. O prefixo properties. é removido antes de passar a configuração para o cliente. O sufixo deve ser uma chave de configuração válida de producer ou consumer do Kafka. Não use este mecanismo para definir key.deserializer ou value.deserializer — o conector gerencia esses parâmetros internamente.

key.fields-prefix

Não

Prefixo personalizado aplicado a todos os nomes de campos de chave no schema da tabela, usado para evitar colisões de nomes com campos de valor. O prefixo é removido durante a análise e geração da chave. Se definido, value.fields-include deve ser configurado como EXCEPT_KEY.

Parâmetros específicos do sink

Parâmetro

Obrigatório

Padrão

Descrição

sink.parallelism

Não

Herdado do operador upstream

Paralelismo do operador sink do Kafka.

sink.buffer-flush.max-rows

Não

0 (desativado)

Número máximo de registros a serem armazenados em buffer antes do flush. Quando várias atualizações chegam para a mesma chave, apenas o registro mais recente é mantido no buffer. Isso reduz o volume de gravação e evita mensagens tombstone desnecessárias. Tanto sink.buffer-flush.max-rows quanto sink.buffer-flush.interval devem ser maiores que zero para ativar o buffer.

sink.buffer-flush.interval

Não

0 (desativado)

Frequência de liberação do buffer. Unidades suportadas: ms, s, min, h — por exemplo, '1 s'. Aplica-se o mesmo comportamento de deduplicação de chave de sink.buffer-flush.max-rows. Ambos os parâmetros devem ser maiores que zero para ativar o buffer.

Ingestão de dados

O conector Upsert Kafka pode servir como sink em um job YAML de ingestão de dados. Os dados são gravados no formato JSON e os campos de chave primária são incluídos no corpo da mensagem.

Sintaxe

sink:
  type: upsert-kafka
  name: upsert-kafka Sink
  properties.bootstrap.servers: localhost:9092
  # ApsaraMQ for Kafka
  aliyun.kafka.accessKeyId: ${secret_values.kafka-ak}
  aliyun.kafka.accessKeySecret: ${secret_values.kafka-sk}
  aliyun.kafka.instanceId: ${instancd-id}
  aliyun.kafka.endpoint: ${endpoint}
  aliyun.kafka.regionId: ${region-id}

Parâmetros

Parâmetro

Obrigatório

Padrão

Descrição

type

Sim

Defina como upsert-kafka.

name

Não

Nome de exibição para o sink.

properties.bootstrap.servers

Sim

Lista separada por vírgulas de endereços de brokers Kafka no formato host:port.

properties.*

Não

Parâmetros adicionais do producer Kafka. O prefixo properties. é removido antes de passar a configuração para o cliente. O sufixo deve ser uma chave de configuração válida de producer do Kafka.

sink.delivery-guarantee

Não

at-least-once

Semântica de entrega. Valores válidos: none (sem garantia; dados podem ser perdidos ou duplicados), at-least-once (padrão; sem perda de dados, mas duplicações são possíveis), exactly-once (usa transações do Kafka para eliminar duplicatas).

sink.add-tableId-to-header-enabled

Não

false

Quando ativado, grava namespace, schemaName e tableName no cabeçalho da mensagem Kafka.

aliyun.kafka.accessKeyId

Não

AccessKey ID da sua conta Alibaba Cloud. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Criar um par de AccessKey.

aliyun.kafka.accessKeySecret

Não

AccessKey secret da sua conta Alibaba Cloud. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Criar um par de AccessKey.

aliyun.kafka.instanceId

Não

ID da instância ApsaraMQ for Kafka. Visualize os detalhes da instância no console Kafka da Alibaba Cloud. Obrigatório ao gravar no ApsaraMQ for Kafka.

aliyun.kafka.endpoint

Não

Endpoint da API para ApsaraMQ for Kafka. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Endpoints.

aliyun.kafka.regionId

Não

ID da região da instância ApsaraMQ for Kafka. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Endpoints.

Alterações de tipo suportadas

O conector de ingestão de dados Upsert Kafka suporta todos os tipos de operação de changelog. Para consumir os dados gravados downstream, use o conector SQL Flink Upsert Kafka com um schema fixo.

Exemplo de ingestão de dados

O exemplo a seguir sincroniza uma tabela MySQL com o ApsaraMQ for Kafka usando um job YAML de ingestão de dados:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: ${mysql.source.table}
  server-id: 8601-8604

sink:
  type: upsert-kafka
  name: Upsert Kafka Sink
  properties.bootstrap.servers: ${upsert.kafka.bootstraps.server}
  aliyun.kafka.accessKeyId: ${upsert.kafka.aliyun.ak}
  aliyun.kafka.accessKeySecret: ${upsert.kafka.aliyun.sk}
  aliyun.kafka.instanceId: ${upsert.kafka.aliyun.instanceid}
  aliyun.kafka.endpoint: ${upsert.kafka.aliyun.endpoint}
  aliyun.kafka.regionId: ${upsert.kafka.aliyun.regionid}

route:
  - source-table: ${mysql.source.table}
    sink-table: ${upsert.kafka.topic}

Semântica de entrega

Por padrão, o sink Upsert Kafka grava com garantias at-least-once. O Flink pode gravar registros duplicados com a mesma chave no tópico. Como o conector opera em modo upsert, o último registro de qualquer chave específica prevalece ao ler o tópico novamente como source. Gravações duplicadas são, portanto, idempotentes — elas não corrompem o estado final.

Para ativar a semântica exactly-once, o cluster Kafka deve suportar transações Kafka (Apache Kafka 0.11 ou posterior) e o recurso de transação deve estar habilitado. Defina sink.delivery-guarantee como exactly-once na configuração YAML de ingestão de dados.

Limites

  • Requer Flink executando no Ververica Runtime (VVR) 2.0.0 ou posterior.

  • Suporta Apache Kafka 0.10 ou posterior tanto para leitura quanto para gravação.

  • Suporta apenas os parâmetros de cliente documentados para Apache Kafka 2.8. Consulte as referências de configuração de consumer e producer.

  • Tabelas source Upsert Kafka sempre iniciam a partir de earliest-offset — isso é fixo e não configurável. O conector precisa ler todos os dados históricos de alterações para reconstruir um changelog completo. Usar qualquer outro modo de inicialização (como latest-offset ou por timestamp) resultaria em um changelog incompleto e causaria problemas de integridade de dados nos cálculos downstream.

  • Para semântica exactly-once no sink, o cluster Kafka de destino deve executar Apache Kafka 0.11 ou posterior com transações habilitadas.

Métricas de monitoramento

Métricas da tabela source

Métrica

Descrição

numRecordsIn

Número total de registros lidos

numRecordsInPerSecond

Registros lidos por segundo

numBytesIn

Total de bytes lidos

numBytesInPerSecond

Bytes lidos por segundo

currentEmitEventTimeLag

Latência entre o tempo do evento e o tempo de emissão

currentFetchEventTimeLag

Latência entre o tempo do evento e o tempo de busca

sourceIdleTime

Tempo em que a source permaneceu ociosa

pendingRecords

Número de registros aguardando processamento

Métricas da tabela sink

Métrica

Descrição

numRecordsOut

Número total de registros gravados

numRecordsOutPerSecond

Registros gravados por segundo

numBytesOut

Total de bytes gravados

numBytesOutPerSecond

Bytes gravados por segundo

currentSendTime

Tempo gasto para enviar o último lote

Próximos passos