O conector Postgres CDC lê um snapshot completo de um banco de dados PostgreSQL e captura dados de alteração com semântica de processamento exactly-once.
Visão geral
O conector Postgres CDC oferece os seguintes recursos:
|
Categoria |
Detalhes |
|
Tipos suportados |
Fonte SQL, fonte Flink CDC Nota
Utilize o conector {{XREF_0}} para tabelas sink e tabelas de lookup (dimensão). |
|
Modo de execução |
Streaming |
|
Formato de dados |
Não aplicável |
|
Métricas |
|
|
Tipos de API |
SQL e Flink CDC |
|
Atualização/exclusão no sink |
Não aplicável |
Recursos
A partir da versão VVR 8.0.6, o conector Postgres CDC integra-se ao framework de snapshot incremental. Ele lê todos os dados históricos e passa automaticamente a ler logs de alteração do WAL, garantindo semântica exactly-once.
Principais recursos:
Processamento unificado de stream e batch. Lê dados completos e incrementais em um único job.
Leitura simultânea de snapshots. Permite escalabilidade horizontal para maior desempenho.
Transição contínua entre fases completa e incremental. Reduz automaticamente a escala para economizar recursos.
Leitura retomável. Retoma a partir de pontos de interrupção durante a fase de snapshot, aumentando a estabilidade.
Leitura sem locks. Não exige bloqueios, evitando impacto nas operações online.
Pré-requisitos
O conector Postgres CDC lê streams CDC por meio da replicação lógica do PostgreSQL. Há suporte para ApsaraDB RDS for PostgreSQL, Amazon RDS for PostgreSQL e PostgreSQL auto-gerenciado.
A configuração varia conforme o tipo de implantação. {{XREF_2}}.
Após configurar, verifique os seguintes itens:
wal_level está definido como
logicalpara habilitar a decodificação lógica.-
REPLICA IDENTITY de cada tabela assinada está definida como
FULL, garantindo que eventos deINSERTeUPDATEincluam valores anteriores das colunas para manter a consistência dos dados.NotaREPLICA IDENTITYé uma configuração no nível da tabela no PostgreSQL que controla se eventos deINSERTeUPDATEincluem valores anteriores das colunas. REPLICA IDENTITY. max_wal_sendersemax_replication_slotssuperam a soma dos slots em uso com os slots necessários para o job do Flink.A conta possui privilégios de
SUPERUSERou combina as permissõesLOGINeREPLICATION, além deSELECTnas tabelas assinadas.
Se sua tabela Postgres contiver colunas geradas, defina o parâmetro publish_generated_columns como
storedao criar o slot. Caso contrário, os schemas das fases de snapshot e incremental podem divergir.
Notas de uso
O recurso de snapshot incremental requer VVR 8.0.6 ou posterior.
Slots de replicação
Jobs Flink PostgreSQL CDC utilizam slots de replicação para evitar a purga prematura do WAL e garantir a consistência dos dados. Slots mal gerenciados podem causar uso excessivo de disco ou atrasos na leitura. Siga estas práticas recomendadas:
-
Limpe slots não utilizados prontamente
O Flink não exclui automaticamente os slots de replicação após a parada de um job ou reinício sem estado, para evitar perda de dados do WAL.
-
Caso o job não seja reiniciado, exclua manualmente seu slot de replicação para liberar espaço em disco.
NotaGerenciamento de ciclo de vida: Trate os slots de replicação como recursos no nível do job e gerencie-os juntamente com as iniciações e paradas do job.
-
Evite reutilizar slots antigos
Sempre utilize um novo nome de slot. Reutilizar um slot antigo força o job a ler dados históricos acumulados do WAL na inicialização, atrasando o processamento de novos dados.
-
O PostgreSQL exige um slot por conexão. Cada job deve usar um nome de slot exclusivo.
NotaConvenção de nomenclatura: Ao personalizar
slot.name, evite nomes com sufixos numéricos comomy_slot_1, para prevenir conflitos com slots temporários.
-
Comportamento do slot com snapshot incremental ativado
Pré-requisitos: Os checkpoints devem estar habilitados e a tabela de origem precisa ter uma chave primária definida.
-
Regras de criação de slots:
Snapshot incremental desativado: Apenas paralelismo 1 é suportado. Um slot global é utilizado.
-
Snapshot incremental ativado:
Fase de snapshot: Cada subtarefa de origem concorrente cria um slot temporário. O formato de nomenclatura é
${slot.name}_${task_id}.Fase incremental: Todos os slots temporários são recuperados automaticamente. Apenas um slot global é mantido.
Número máximo de slots: Paralelismo da origem + 1 (durante a fase de snapshot)
-
Recursos e desempenho
Se os slots disponíveis ou o espaço em disco forem limitados, reduza o paralelismo do snapshot para utilizar menos slots temporários. Isso diminui a velocidade de leitura do snapshot.
-
Caso o sink downstream suporte gravações idempotentes, defina
scan.incremental.snapshot.backfill.skip = truepara ignorar o backfill do WAL durante a fase de snapshot e acelerar a inicialização.Isso fornece apenas semântica at-least-once e não é adequado para computações com estado (agregações ou lookup joins), pois alterações históricas necessárias podem ser perdidas.
-
Quando os snapshots incrementais estão desativados, não há suporte para checkpoints durante a fase de snapshot.
Reutilizar assinatura do Postgres
O conector Postgres CDC depende de uma publicação para determinar quais alterações de tabelas são enviadas para um slot. Se múltiplos jobs compartilharem a mesma publicação, suas configurações serão sobrescritas.
Causa
O modo padrão publication.autocreate.mode é filtered, que inclui apenas as tabelas presentes na configuração do conector. Isso modifica a publicação na inicialização do job, podendo afetar outros jobs.
Solução
-
Crie uma publicação no PostgreSQL que inclua todas as tabelas monitoradas ou crie uma publicação separada por job.
-- Create a publication named my_flink_pub that includes all tables (or specified tables, creating one publication per job) CREATE PUBLICATION my_flink_pub FOR TABLE table_a, table_b; -- Or more simply, include all tables in the database CREATE PUBLICATION my_flink_pub FOR ALL TABLES;NotaAssinar todas as tabelas não é recomendado para bancos de dados grandes devido ao uso excessivo de largura de banda e CPU no cluster Flink.
-
Adicione as seguintes configurações do Flink:
debezium.publication.name = 'my_flink_pub'(Especifica o nome da publicação)debezium.publication.autocreate.mode = 'disabled'(Impede que o Flink tente criar ou modificar a publicação na inicialização)
Isso garante isolamento completo e evita que novos jobs afetem os existentes.
SQL
Sintaxe
CREATE TABLE postgrescdc_source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>',
'decoding.plugin.name'= 'pgoutput',
'scan.incremental.snapshot.enabled' = 'true',
-- Skipping backfill can speed up reads and reduce resource usage, but it may cause data duplication. Enable this if the downstream sink is idempotent.
'scan.incremental.snapshot.backfill.skip' = 'false',
-- In a production environment, set this to 'filtered' or 'disabled' and manage the publication manually instead of through Flink.
'debezium-publication.autocreate.mode' = 'disabled'
-- If you have multiple sources, configure a different publication for each source.
--'debezium.publication.name' = 'my_flink_pub'
);
Opções do conector
|
Opção |
Descrição |
Tipo de dados |
Obrigatório |
Padrão |
Observações |
|
connector |
Nome do conector. |
STRING |
Sim |
– |
O valor deve ser |
|
hostname |
Endereço IP ou nome do host do banco de dados PostgreSQL. |
STRING |
Sim |
– |
– |
|
username |
Nome de usuário do serviço de banco de dados PostgreSQL. |
STRING |
Sim |
– |
– |
|
password |
Senha do serviço de banco de dados PostgreSQL. |
STRING |
Sim |
– |
– |
|
database-name |
Nome do banco de dados PostgreSQL. |
STRING |
Sim |
– |
Nome do banco de dados. |
|
schema-name |
Nome do schema PostgreSQL. Suporta regex. |
STRING |
Sim |
– |
O nome do schema aceita expressões regulares para leitura de dados de múltiplos schemas. |
|
table-name |
Nome da tabela PostgreSQL. Suporta regex. |
STRING |
Sim |
– |
O nome da tabela aceita expressões regulares para leitura de dados de múltiplas tabelas. |
|
port |
Número da porta. |
INTEGER |
Não |
5432 |
– |
|
decoding.plugin.name |
Nome do plugin de decodificação lógica do PostgreSQL. |
STRING |
Não |
decoderbufs |
Determinado pelo plugin instalado no serviço PostgreSQL. Plugins suportados:
|
|
slot.name |
Nome do slot de decodificação lógica. |
STRING |
Obrigatório para VVR 8.0.1 e posteriores. Opcional para versões anteriores. |
|
Defina um Sem valor padrão para VVR 8.0.1+. |
|
debezium.* |
Propriedades e parâmetros do Debezium |
STRING |
Não |
– |
Oferece controle mais granular sobre o comportamento do cliente Debezium. Por exemplo, |
|
scan.incremental.snapshot.enabled |
Define se os snapshots incrementais devem ser ativados. |
BOOLEAN |
Não |
false |
Nota
|
|
scan.startup.mode |
Modo de inicialização para consumo de dados. |
STRING |
Não |
initial |
Valores válidos:
|
|
changelog-mode |
Modo de changelog para codificar alterações de stream. |
String |
Não |
all |
Modos de changelog suportados:
|
|
heartbeat.interval.ms |
Intervalo para envio de pacotes de heartbeat. |
Duration |
Não |
30s |
A unidade é milissegundos. O conector Postgres CDC envia heartbeats ativamente ao banco de dados para avançar o offset do slot. Quando as alterações na tabela são pouco frequentes, definir este valor garante a recuperação oportuna dos logs WAL. |
|
scan.incremental.snapshot.chunk.key-column |
Especifica uma coluna para ser usada como chave de chunk na divisão de shards durante a fase de snapshot. |
STRING |
Não |
– |
Por padrão, a primeira coluna da chave primária é selecionada. |
|
scan.incremental.close-idle-reader.enabled |
Define se leitores ociosos devem ser fechados após a conclusão do snapshot. |
Boolean |
Não |
false |
Para ativar esta configuração, defina |
|
scan.incremental.snapshot.backfill.skip |
Define se a leitura de logs durante a fase de snapshot deve ser ignorada. |
Boolean |
Não |
false |
Valores válidos:
|
Mapeamentos de tipos
Mapeamentos de tipos PostgreSQL para Flink:
|
PostgreSQL CDC |
Flink |
|
SMALLINT |
SMALLINT |
|
INT2 |
|
|
SMALLSERIAL |
|
|
SERIAL2 |
|
|
INTEGER |
INT |
|
SERIAL |
|
|
BIGINT |
BIGINT |
|
BIGSERIAL |
|
|
REAL |
FLOAT |
|
FLOAT4 |
|
|
FLOAT8 |
DOUBLE |
|
DOUBLE PRECISION |
|
|
NUMERIC(p, s) |
DECIMAL(p, s) |
|
DECIMAL(p, s) |
|
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
|
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
|
CHAR(n) |
STRING |
|
CHARACTER(n) |
|
|
VARCHAR(n) |
|
|
CHARACTER VARYING(n) |
|
|
TEXT |
|
|
BYTEA |
BYTES |
Exemplo
CREATE TABLE source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>'
);
SELECT * FROM source;
Flink CDC
O VVR V11.4+ suporta o conector PostgreSQL como fonte Flink CDC.
Sintaxe
source:
type: postgres
name: PostgreSQL Source
hostname: localhost
port: 5432
username: pg_username
password: pg_password
tables: db.scm.tbl
slot.name: test_slot
scan.startup.mode: initial
server-time-zone: UTC
connect.timeout: 120s
decoding.plugin.name: decoderbufs
sink:
type: ...
Opções do conector
|
Opção |
Descrição |
Obrigatório |
Tipo de dados |
Padrão |
Observações |
|
type |
Nome do conector. |
Sim |
STRING |
– |
Deve ser |
|
name |
Nome da fonte de dados. |
Não |
STRING |
– |
– |
|
hostname |
Nome de domínio ou endereço IP do servidor de banco de dados PostgreSQL. |
Sim |
STRING |
– |
– |
|
port |
Porta do banco de dados PostgreSQL. |
Não |
INTEGER |
5432 |
– |
|
username |
Nome de usuário do PostgreSQL. |
Sim |
STRING |
– |
– |
|
password |
Senha do PostgreSQL. |
Sim |
STRING |
– |
– |
|
tables |
Nomes das tabelas a serem capturadas. Suporta regex. |
Sim |
STRING |
– |
Importante
Atualmente, apenas tabelas no mesmo banco de dados podem ser capturadas. Um ponto (.) é tratado como separador para um nome totalmente qualificado. Para usar um ponto (.) em uma expressão regular correspondendo a qualquer caractere, escape-o com uma barra invertida. Exemplo: |
|
slot.name |
Nome do slot de replicação do PostgreSQL. |
Sim |
STRING |
– |
O nome deve obedecer às regras de nomenclatura de slots de replicação do PostgreSQL e pode conter letras minúsculas, números e sublinhados. |
|
decoding.plugin.name |
Nome do plugin de decodificação lógica do PostgreSQL instalado no servidor. |
Não |
STRING |
|
Valores válidos: |
|
tables.exclude |
Tabelas a excluir. Esta opção entra em vigor após a opção |
Não |
STRING |
– |
Consulte a opção |
|
server-time-zone |
Fuso horário da sessão do servidor de banco de dados, como "Asia/Shanghai". |
Não |
STRING |
– |
Se não definido, o fuso horário padrão do sistema ( |
|
scan.incremental.snapshot.chunk.size |
Tamanho (número de linhas) de cada chunk no framework de snapshot incremental. |
Não |
INTEGER |
8096 |
Quando o snapshot incremental está ativado, a tabela é dividida em vários chunks para leitura. Os dados de um chunk são armazenados em cache na memória antes de serem totalmente consumidos. Um chunk menor resulta em um número total maior de chunks para a tabela. Embora isso reduza a granularidade da recuperação de falhas, pode levar a erros de falta de memória (OOM) e menor throughput geral. Portanto, é necessário encontrar um equilíbrio e definir um tamanho de chunk razoável. |
|
scan.snapshot.fetch.size |
Número máximo de registros a buscar por vez ao ler todos os dados de uma tabela. |
Não |
INTEGER |
1024 |
– |
|
scan.startup.mode |
Modo de inicialização para consumo de dados. |
Não |
STRING |
initial |
Valores válidos:
|
|
scan.incremental.close-idle-reader.enabled |
Define se leitores ociosos devem ser fechados após a conclusão do snapshot. |
Não |
BOOLEAN |
false |
Para ativar esta configuração, defina |
|
scan.lsn-commit.checkpoints-num-delay |
Número de checkpoints a adiar antes de começar a confirmar offsets LSN. |
Não |
INTEGER |
3 |
Offsets LSN de checkpoint são confirmados de forma rotativa para evitar a impossibilidade de recuperação de estado. |
|
connect.timeout |
Tempo máximo que o conector aguarda para se conectar ao servidor de banco de dados PostgreSQL antes de atingir o timeout. |
Não |
DURATION |
30s |
Este valor não pode ser inferior a 250 milissegundos. |
|
connect.max-retries |
Número máximo de tentativas de nova conexão pelo conector. |
Não |
INTEGER |
3 |
– |
|
connection.pool.size |
Tamanho do pool de conexões. |
Não |
INTEGER |
20 |
– |
|
jdbc.properties.* |
Permite passar propriedades personalizadas de URL JDBC. |
Não |
STRING |
20 |
É possível passar propriedades personalizadas, como |
|
heartbeat.interval |
Intervalo para envio de eventos de heartbeat visando rastrear o offset de log WAL mais recente disponível. |
Não |
DURATION |
30s |
– |
|
debezium.* |
Passa propriedades do Debezium para o Debezium Embedded Engine, usado para capturar alterações de dados do servidor PostgreSQL. |
Não |
STRING |
– |
Propriedades do conector Debezium PostgreSQL: Documentação do Debezium. |
|
chunk-meta.group.size |
Tamanho dos metadados do chunk. |
Não |
STRING |
1000 |
Se os metadados forem maiores que este valor, eles são passados em várias partes. |
|
metadata.list |
Lista de metadados legíveis passados para o downstream, utilizáveis no módulo de transformação. |
Não |
STRING |
false |
Use vírgulas (,) como separadores. Atualmente, os metadados disponíveis são: |
|
scan.incremental.snapshot.unbounded-chunk-first.enabled |
Despacha o chunk ilimitado primeiro durante a fase de leitura do snapshot. |
Não |
STRING |
false |
Recurso experimental. Ativá-lo pode reduzir o risco de erros OOM quando o TaskManager sincroniza o último chunk durante a fase de snapshot. Recomendamos adicionar isso antes da primeira inicialização do job. |