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:
Um cluster Kafka. Crie um cluster Kafka DataFlow ou crie recursos no ApsaraMQ for Kafka
Conectividade de rede entre o cluster Flink e o cluster Kafka. Para Kafka no EMR, configure uma VPC e um grupo de segurança. Para ApsaraMQ for Kafka, configure uma lista de permissões
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 |
|
|
Inserção |
Novo registro cuja chave não existia anteriormente |
|
|
Atualização posterior |
Registro cuja chave já existe — sobrescreve o valor anterior |
|
|
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
INSERTeUPDATE_AFTERsão gravados como mensagens normais do KafkaRegistros
DELETEsã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 |
|
|
Sim |
— |
Defina como |
|
|
Sim |
— |
Lista separada por vírgulas de endereços de brokers Kafka no formato |
|
|
Sim |
— |
Tópico Kafka para leitura ou gravação. |
|
|
Sim |
— |
Formato de serialização para a chave da mensagem. Valores válidos: |
|
|
Sim |
— |
Formato de serialização para o valor da mensagem. Equivalente a |
|
|
Sim |
|
Controla quais colunas são incluídas no valor da mensagem. |
|
|
Não |
— |
Parâmetros adicionais do cliente Kafka. O prefixo |
|
|
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, |
Parâmetros específicos do sink
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Não |
Herdado do operador upstream |
Paralelismo do operador sink do Kafka. |
|
|
Não |
|
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 |
|
|
Não |
|
Frequência de liberação do buffer. Unidades suportadas: |
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 |
|
|
Sim |
— |
Defina como |
|
|
Não |
— |
Nome de exibição para o sink. |
|
|
Sim |
— |
Lista separada por vírgulas de endereços de brokers Kafka no formato |
|
|
Não |
— |
Parâmetros adicionais do producer Kafka. O prefixo |
|
|
Não |
|
Semântica de entrega. Valores válidos: |
|
|
Não |
|
Quando ativado, grava |
|
|
Não |
— |
AccessKey ID da sua conta Alibaba Cloud. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Criar um par de AccessKey. |
|
|
Não |
— |
AccessKey secret da sua conta Alibaba Cloud. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Criar um par de AccessKey. |
|
|
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. |
|
|
Não |
— |
Endpoint da API para ApsaraMQ for Kafka. Obrigatório ao gravar no ApsaraMQ for Kafka. Consulte Endpoints. |
|
|
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 (comolatest-offsetou 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 |
|
|
Número total de registros lidos |
|
|
Registros lidos por segundo |
|
|
Total de bytes lidos |
|
|
Bytes lidos por segundo |
|
|
Latência entre o tempo do evento e o tempo de emissão |
|
|
Latência entre o tempo do evento e o tempo de busca |
|
|
Tempo em que a source permaneceu ociosa |
|
|
Número de registros aguardando processamento |
Métricas da tabela sink
|
Métrica |
Descrição |
|
|
Número total de registros gravados |
|
|
Registros gravados por segundo |
|
|
Total de bytes gravados |
|
|
Bytes gravados por segundo |
|
|
Tempo gasto para enviar o último lote |