Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector do AnalyticDB for PostgreSQL

Última atualização: Jun 27, 2026

O conector do AnalyticDB for PostgreSQL permite usar o AnalyticDB for PostgreSQL como tabela source (beta), de dimensão ou sink em jobs SQL do Realtime Compute for Apache Flink. O AnalyticDB for PostgreSQL é um data warehouse de processamento massivamente paralelo (MPP) para análises online em grande escala.

Suporte: Source (beta) · Dimensão · Sink | Streaming · Batch | API SQL | O sink suporta atualizações e exclusões

A leitura de uma tabela source exige um conector personalizado configurado via Flink CDC. Para instruções de configuração, consulte Usar o Flink CDC para assinar dados completos e incrementais em tempo real .

Pré-requisitos

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

Limitações

  • O AnalyticDB for PostgreSQL V7.0 requer o VVR 8.0.1 ou posterior.

  • Não há suporte para bancos de dados PostgreSQL autogerenciados.

Sintaxe

Use connector='adbpg' para tabelas de dimensão e sink, e connector='adbpg-cdc' para tabelas source.

CREATE TEMPORARY TABLE adbpg_table (
  id      INT,
  len     INT,
  content VARCHAR,
  PRIMARY KEY (id)
) WITH (
  'connector'  = 'adbpg',
  'url'        = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName'  = '<table>',
  'userName'   = '<username>',
  'password'   = '<password>'
);

Opções do conector

Gerais

Estas opções aplicam-se a todos os tipos de tabela.

Opção

Obrigatória

Padrão

Tipo

Descrição

connector

Sim

STRING

Defina como adbpg-cdc para tabelas source ou adbpg para tabelas de dimensão e sink.

url

Sim

STRING

URL JDBC (Java Database Connectivity) no formato jdbc:postgresql://<host>:<port>/<database>.

tableName

Sim

STRING

Nome da tabela no banco de dados.

userName

Sim

STRING

Nome de usuário do banco de dados AnalyticDB for PostgreSQL.

password

Sim

STRING

Senha do banco de dados AnalyticDB for PostgreSQL.

maxRetryTimes

Não

3

INTEGER

Número máximo de tentativas de nova execução após falha na gravação.

targetSchema

Não

public

STRING

Nome do schema do banco de dados.

caseSensitive

Não

false

STRING

Define se os identificadores diferenciam maiúsculas de minúsculas. Valores válidos: true, false.

connectionMaxActive

Não

5

INTEGER

Quantidade máxima de conexões no pool. O conector libera conexões ociosas automaticamente. Definir esse valor muito alto pode causar contagens anormais de conexões no servidor.

Específicas para source (beta)

Estas opções aplicam-se apenas a tabelas source (connector='adbpg-cdc').

Opção

Obrigatória

Padrão

Tipo

Descrição

schema-name

Sim

STRING

Nome do schema. Suporta expressões regulares para assinar vários schemas simultaneamente.

port

Sim

5432

INTEGER

Porta da instância do AnalyticDB for PostgreSQL.

decoding.plugin.name

Sim

pgoutput

STRING

Plug-in de decodificação lógica do PostgreSQL. Defina como pgoutput.

slot.name

Sim

STRING

Nome do slot de decodificação lógica. Consulte Orientações sobre nomes de slots.

debezium.*

Sim

STRING

Propriedades de configuração do cliente Debezium. Por exemplo, defina 'debezium.snapshot.mode' = 'never' para ignorar o snapshot inicial. Consulte Propriedades do conector.

scan.incremental.snapshot.enabled

Não

false

BOOLEAN

Define se o snapshot incremental deve ser ativado. Valores válidos: true, false.

scan.startup.mode

Não

initial

STRING

Modo de inicialização para consumo de dados. Consulte Modos de inicialização.

changelog-mode

Não

ALL

STRING

Forma como os eventos de alteração são codificados no fluxo de mudanças. Valores válidos: ALL, UPSERT.

heartbeat.interval.ms

Não

30 segundos

DURATION

Intervalo para envio de pacotes de heartbeat, em milissegundos. Consulte Heartbeat e retenção de WAL.

scan.incremental.snapshot.chunk.key-column

Não

Primeira coluna da chave primária

STRING

Coluna usada como chave de fragmento durante as leituras de snapshot.

Modos de inicialização

Valor

Comportamento

initial (padrão)

Executa uma varredura completa dos dados históricos na primeira inicialização e, em seguida, lê a partir da posição mais recente do Write-Ahead Logging (WAL).

latest-offset

Ignora a varredura de dados históricos e lê apenas a partir do final atual do WAL.

snapshot

Realiza uma varredura completa dos dados históricos e captura os dados do WAL gerados durante essa varredura. Interrompe a operação após a conclusão da varredura.

Modos de changelog

Valor

Eventos capturados

ALL (padrão)

INSERT, DELETE, UPDATE_BEFORE, UPDATE_AFTER

UPSERT

INSERT, DELETE, UPDATE_AFTER

Orientações sobre nomes de slots

Gerencie os slots de replicação com cuidado para evitar conflitos:

  • Em um único job do Flink: Use o mesmo slot.name para todas as tabelas source do job.

  • Entre jobs diferentes do Flink: Atribua um slot.name exclusivo para cada job. Reutilizar o nome de um slot entre jobs causa o erro PSQLException: ERROR: replication slot "debezium" is active for PID 974.

Heartbeat e retenção de WAL

Quando uma tabela source recebe poucas atualizações, o conector pode não avançar o offset do slot de replicação. Isso impede que o AnalyticDB for PostgreSQL recupere espaço em disco do WAL, pois o banco de dados retém os arquivos WAL até que o offset do slot avance além deles. O envio periódico de heartbeats mantém o offset do slot em movimento, permitindo a liberação de arquivos WAL antigos. Se as alterações na tabela forem pouco frequentes, defina heartbeat.interval.ms com um valor adequado à sua frequência de atualização.

Específicas para sink

Estas opções aplicam-se apenas a tabelas sink (connector='adbpg').

Opção

Obrigatória

Padrão

Tipo

Descrição

writeMode

Não

insert

STRING

Estratégia de gravação para a primeira tentativa. Consulte Modos de gravação e resolução de conflitos.

conflictMode

Não

strict

STRING

Define como tratar conflitos de chave primária ou índice. Consulte Modos de gravação e resolução de conflitos.

batchSize

Não

500

INTEGER

Quantidade de registros gravados por lote.

flushIntervalMs

Não

INTEGER

Tempo máximo, em milissegundos, que o conector aguarda antes de liberar os registros armazenados em cache. Caso o buffer não atinja o batchSize dentro desse período, todos os registros no buffer serão gravados imediatamente.

retryWaitTime

Não

100

INTEGER

Intervalo entre as tentativas de nova execução, em milissegundos.

Modos de gravação e resolução de conflitos

As opções writeMode e conflictMode funcionam em conjunto para controlar a gravação de registros e o tratamento de conflitos.

**writeMode**

Descrição

Observações

insert (padrão)

Insere registros diretamente. O tratamento de conflitos segue a definição de conflictMode.

Adequado para a maioria dos casos de uso.

upsert

Atualiza automaticamente os registros existentes em caso de conflito.

Exige chave primária.

copy

Insere registros usando o comando COPY do PostgreSQL.

Requer VVR 11.1 ou posterior.

**conflictMode**

Comportamento em caso de conflito

Observações

strict (padrão)

Reporta um erro.

ignore

Descarta silenciosamente o registro conflitante.

update

Atualiza o registro conflitante.

Recomendado para tabelas sem chave primária. Reduz o throughput de gravação.

upsert

Atualiza o registro conflitante.

Exige chave primária. Mais eficiente que update para tabelas com chave.

Específicas para tabela de dimensão

Estas opções aplicam-se apenas a tabelas de dimensão (uso de connector='adbpg' em um JOIN com FOR SYSTEM_TIME AS OF).

Opção

Obrigatória

Padrão

Tipo

Descrição

cache

Não

ALL

STRING

Política de cache para consultas à tabela de dimensão. Consulte Políticas de cache.

cacheSize

Não

100000

LONG

Número máximo de linhas mantidas no cache. Tem efeito apenas quando cache=LRU.

cacheTTLMs

Não

Long.MAX_VALUE

LONG

Tempo de vida das entradas de cache em milissegundos. O comportamento depende da configuração de cache. Consulte Políticas de cache.

maxJoinRows

Não

1024

INTEGER

Quantidade máxima de linhas a serem unidas por registro de entrada.

Políticas de cache

A escolha da política de cache correta equilibra o throughput de consulta e a atualização dos dados.

**Valor de cache**

Comportamento

Quando usar

ALL (padrão)

Carrega toda a tabela de dimensão na memória antes do início do job. Todas as consultas acessam o cache. Ao expirar o cache (cacheTTLMs), recarrega a tabela inteira. Se uma chave de junção não for encontrada no cache, ela não existe.

Tabelas de dimensão pequenas ou médias, nas quais a atualização dos dados tolera recargas periódicas.

LRU

Armazena em cache um subconjunto de linhas. Em caso de falta no cache, busca os dados no banco e atualiza o cache. Remove as entradas menos usadas recentemente quando o cache atinge o cacheSize. O parâmetro cacheTTLMs controla a expiração individual das entradas.

Tabelas de dimensão grandes, onde armazenar tudo em cache é inviável, ou quando o cache parcial de baixa latência é aceitável.

None

Sem cache. Toda consulta vai diretamente ao banco de dados.

Quando os dados precisam sempre ser lidos atualizados do banco de dados.

Mapeamentos de tipos de dados

Tipo do AnalyticDB for PostgreSQL

Tipo do Flink SQL

BOOLEAN

BOOLEAN

SMALLINT

INT

INT

INT

BIGINT

BIGINT

FLOAT

DOUBLE

VARCHAR

VARCHAR

TEXT

VARCHAR

TIMESTAMP

TIMESTAMP

DATE

DATE

Exemplos

Tabela sink

Este exemplo lê de uma source datagen e grava em uma tabela sink do AnalyticDB for PostgreSQL.

CREATE TEMPORARY TABLE datagen_source (
  `name` VARCHAR,
  `age`  INT
)
COMMENT 'datagen source table'
WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adbpg_sink (
  name VARCHAR,
  age  INT
) WITH (
  'connector' = 'adbpg',
  'url'       = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName' = '<table>',
  'userName'  = '<username>',
  'password'  = '<password>'
);

INSERT INTO adbpg_sink
SELECT * FROM datagen_source;

Tabela de dimensão

Este exemplo une um fluxo datagen a uma tabela de dimensão do AnalyticDB for PostgreSQL usando uma junção temporal.

CREATE TEMPORARY TABLE datagen_source (
  a         INT,
  b         BIGINT,
  c         STRING,
  `proctime` AS PROCTIME()
)
COMMENT 'datagen source table'
WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE adbpg_dim (
  a INT,
  b VARCHAR,
  c VARCHAR
) WITH (
  'connector' = 'adbpg',
  'url'       = 'jdbc:postgresql://<host>:<port>/<database>',
  'tableName' = '<table>',
  'userName'  = '<username>',
  'password'  = '<password>'
);

CREATE TEMPORARY TABLE blackhole_sink (
  a INT,
  b STRING
)
COMMENT 'blackhole sink table'
WITH (
  'connector' = 'blackhole'
);

INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN adbpg_dim FOR SYSTEM_TIME AS OF T.proctime AS H
  ON T.a = H.a;

Tabela source (beta)

Para configuração e exemplos de tabela source, consulte Usar o Flink CDC para assinar dados completos e incrementais em tempo real.

Métricas

As métricas da tabela sink estão disponíveis no console do Realtime Compute for Apache Flink. Tabelas de dimensão não expõem métricas. Para definições das métricas, consulte Métricas.

Métrica

Descrição

numRecordsOut

Total de registros gravados no sink.

numRecordsOutPerSecond

Registros gravados por segundo.

numBytesOut

Total de bytes gravados no sink.

numBytesOutPerSecond

Bytes gravados por segundo.

currentSendTime

Tempo gasto na operação de gravação mais recente.

Próximos passos