Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:PostgreSQL CDC

Última atualização: Aug 20, 2026

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

Métricas de monitoramento

  • currentFetchEventTimeLag: Intervalo entre a geração dos dados e o momento em que o operador de origem os recupera.

  • currentEmitEventTimeLag: Intervalo entre a geração dos dados e o momento em que saem do operador de origem.

  • sourceIdleTime: Duração durante a qual a origem não produziu novos dados.

Nota
  • As métricas currentFetchEventTimeLag e currentEmitEventTimeLag são válidas apenas na fase incremental. Na fase de snapshot, o valor é sempre 0.

  • Para mais informações sobre as métricas, consulte {{XREF_1}}.

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.

Importante

A configuração varia conforme o tipo de implantação. {{XREF_2}}.

Após configurar, verifique os seguintes itens:

  • wal_level está definido como logical para habilitar a decodificação lógica.

  • REPLICA IDENTITY de cada tabela assinada está definida como FULL, garantindo que eventos de INSERT e UPDATE incluam valores anteriores das colunas para manter a consistência dos dados.

    Nota

    REPLICA IDENTITY é uma configuração no nível da tabela no PostgreSQL que controla se eventos de INSERT e UPDATE incluem valores anteriores das colunas. REPLICA IDENTITY.

  • max_wal_senders e max_replication_slots superam a soma dos slots em uso com os slots necessários para o job do Flink.

  • A conta possui privilégios de SUPERUSER ou combina as permissões LOGIN e REPLICATION, além de SELECT nas tabelas assinadas.

  • Se sua tabela Postgres contiver colunas geradas, defina o parâmetro publish_generated_columns como stored ao 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.

      Nota

      Gerenciamento 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.

      Nota

      Convenção de nomenclatura: Ao personalizar slot.name, evite nomes com sufixos numéricos como my_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 = true para 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.

    Configuração para evitar timeouts durante a fase de snapshot

    Quando os snapshots incrementais estão desativados, os checkpoints durante a fase de snapshot podem causar failovers por timeout. Configure estes parâmetros em Other Configuration ({{XREF_3}}):

    execution.checkpointing.interval: 10min
    execution.checkpointing.tolerable-failed-checkpoints: 100
    restart-strategy: fixed-delay
    restart-strategy.fixed-delay.attempts: 2147483647

    Parâmetros:

    Parâmetro

    Descrição

    Observações

    execution.checkpointing.interval

    Intervalo entre checkpoints.

    A unidade é um valor de duração, como 10min ou 30s.

    execution.checkpointing.tolerable-failed-checkpoints

    Número de falhas de checkpoint toleráveis antes que o job falhe.

    O produto deste parâmetro pelo intervalo de agendamento do checkpoint resulta no tempo permitido para leitura do snapshot.

    Nota

    Se a tabela for muito grande, defina este parâmetro com um valor maior.

    restart-strategy

    Estratégia de reinício do job.

    Valores válidos:

    • fixed-delay: Estratégia de reinício com atraso fixo.

    • failure-rate: Estratégia de reinício baseada em taxa de falha.

    • exponential-delay: Estratégia de reinício com atraso exponencial.

    Estratégias de Reinício.

    restart-strategy.fixed-delay.attempts

    Número máximo de tentativas de reinício para a estratégia fixed-delay.

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

  1. 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;
    Nota

    Assinar todas as tabelas não é recomendado para bancos de dados grandes devido ao uso excessivo de largura de banda e CPU no cluster Flink.

  2. 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 postgres-cdc.

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:

  • decoderbufs: Suportado no PostgreSQL 9.6 e posteriores. Este plugin deve ser instalado.

  • pgoutput (recomendado): Plugin oficial integrado para PostgreSQL 10 e posteriores.

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.

flink (Anterior à 8.0.1)

Defina um slot.name exclusivo para cada tabela a fim de evitar o erro PSQLException: ERROR: replication slot "debezium" is active for PID 974. {{XREF_4}}.

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, 'debezium.snapshot.mode' = 'never'. Propriedades de configuração.

scan.incremental.snapshot.enabled

Define se os snapshots incrementais devem ser ativados.

BOOLEAN

Não

false

Nota
  • Recurso experimental disponível no VVR 8.0.6+.

  • Benefícios, pré-requisitos e limites: {{XREF_5}}, {{XREF_6}} e {{XREF_7}}.

scan.startup.mode

Modo de inicialização para consumo de dados.

STRING

Não

initial

Valores válidos:

  • initial: Examina todos os dados históricos na primeira inicialização e depois lê os dados mais recentes do WAL.

  • latest-offset: Não examina todos os dados históricos na primeira inicialização. Começa a leitura a partir do final do WAL, lendo apenas as alterações mais recentes feitas após o início do conector.

  • snapshot: Examina todos os dados históricos, lê os novos dados do WAL gerados durante a fase de snapshot e então o job é encerrado.

changelog-mode

Modo de changelog para codificar alterações de stream.

String

Não

all

Modos de changelog suportados:

  • ALL: Suporta todos os tipos, incluindo INSERT, DELETE, UPDATE_BEFORE e UPDATE_AFTER.

  • UPSERT: Suporta apenas o tipo upsert, que inclui INSERT, DELETE e UPDATE_AFTER.

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 execution.checkpointing.checkpoints-after-tasks-finish.enabled como true.

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:

  • true: ignora.

    Na fase incremental, os logs são lidos a partir da marca d'água baixa.

    Se operadores downstream ou armazenamento suportarem idempotência, recomenda-se ignorar a leitura de logs na fase completa. Isso reduz o número de slots WAL, mas garante apenas semântica at-least-once.

  • false: não ignora.

    Ao ler splits na fase completa, os logs entre a marca d'água baixa e a alta são lidos para garantir consistência.

    Se o SQL executar agregações, joins ou operações similares, não recomendamos ignorar a leitura de logs na fase completa.

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 postgres.

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: bdb.schema_\..order_\..

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

pgoutput

Valores válidos: decoderbufs e pgoutput.

tables.exclude

Tabelas a excluir. Esta opção entra em vigor após a opção tables. Suporta regex.

Não

STRING

Consulte a opção tables.

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 (ZoneId.systemDefault()) é usado.

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:

  • initial (padrão): Examina o snapshot na primeira inicialização e depois muda para os dados mais recentes do WAL.

  • latest-offset: Ignora a leitura do snapshot; Começa a leitura a partir do final do WAL, lendo apenas as alterações mais recentes feitas após o início do conector.

  • committed-offset: Ignora a leitura do snapshot; Consome dados do WAL a partir de um offset especificado.

  • snapshot: Consome apenas o snapshot, sem dados incrementais.

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 execution.checkpointing.checkpoints-after-tasks-finish.enabled como true.

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 'jdbc.properties.useSSL' = 'false'.

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: op_ts.

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.

Referências