AnalyticDB for PostgreSQL oferece um conector de Change Data Capture (CDC) desenvolvido internamente que assina dados completos e incrementais com base no recurso de replicação lógica do PostgreSQL. Esse conector integra-se perfeitamente ao Flink e captura alterações de dados em tempo real nas tabelas de origem para sincronização e processamento de fluxo contínuo. Isso permite que as empresas respondam rapidamente a requisitos dinâmicos de dados. Este tópico descreve como usar o Realtime Compute for Apache Flink CDC para assinar dados completos e incrementais do AnalyticDB for PostgreSQL em tempo real.
Limites
-
Este recurso está disponível apenas para instâncias do AnalyticDB for PostgreSQL V7.0 que executam a versão secundária do mecanismo 7.2.1.4 ou posterior.
NotaVisualize a minor engine version na página Basic Information da instância no console. Caso a versão não atenda aos requisitos acima, upgrade the minor engine version.
O modo serverless do AnalyticDB for PostgreSQL não tem suporte.
Pré-requisitos
A instância do AnalyticDB for PostgreSQL e o fully managed Flink workspace devem estar na mesma VPC.
-
Ajuste as parameter settings da instância do AnalyticDB for PostgreSQL:
Ative a replicação lógica definindo o parâmetro
wal_levelcomo logical.Para instâncias da High-availability Edition do AnalyticDB for PostgreSQL, defina os parâmetros
hot_standby,hot_standby_feedbackesync_replication_slotscomo on. Essa configuração garante que a assinatura lógica não seja interrompida por um failover entre primário e secundário.
Utilize uma initial account ou um privileged user with the RDS_SUPERUSER permission na instância do AnalyticDB for PostgreSQL. O usuário deve ter o privilégio REPLICATION:
ALTER USER <username> WITH REPLICATION;.Adicione o bloco CIDR do workspace do Flink a uma whitelist da instância do AnalyticDB for PostgreSQL.
Baixe o arquivo flink-sql-connector-adbpg-cdc-3.3.jar e upload the CDC connector para o seu workspace do Flink.
Procedimento
Etapa 1: Preparar uma tabela de teste e dados de teste
Faça logon no console do AnalyticDB for PostgreSQL e localize a instância desejada. Clique em ID da instância.
No canto inferior direito da página Basic Information, clique em Log On to Database.
-
Crie um banco de dados de teste e uma tabela de origem chamada adbpg_source_table. Em seguida, insira 50 linhas de dados nessa tabela.
-- Create a test database. CREATE DATABASE testdb; -- Switch to the testdb database and create a schema. CREATE SCHEMA testschema; -- Create a source table named adbpg_source_table. CREATE TABLE testschema.adbpg_source_table( id int, username text, PRIMARY KEY(id) ); -- Insert 50 rows of data into the adbpg_source_table table. INSERT INTO testschema.adbpg_source_table(id, username) SELECT i, 'username'||i::text FROM generate_series(1, 50) AS t(i); -
Crie uma tabela de destino chamada adbpg_sink_table para que o Flink grave os dados resultantes.
CREATE TABLE testschema.adbpg_sink_table( id int, username text, score int );
Etapa 2: Criar um job do Flink
Faça logon no console do Realtime Compute. Na aba Fully Managed Flink, localize o workspace desejado e clique em Console na coluna Actions.
No painel de navegação à esquerda, escolha .
-
Na barra de menu superior, clique em
. Selecione New Blank Stream Draft e configure os parâmetros abaixo.Parâmetro do job
Descrição
Exemplo
Name
Nome do job.
NotaO nome do job deve ser exclusivo dentro do projeto atual.
adbpg-test
Location
Pasta onde o arquivo de código do job será armazenado.
Também é possível clicar em ícone
ao lado de uma pasta existente para criar uma subpasta.Job Drafts
Engine Version
Versão do mecanismo Flink utilizada pelo job. Para mais informações sobre números de versão, mapeamentos de versões e marcos do ciclo de vida, consulte Flink engine versions.
vvr-6.0.7-flink-1.15
Clique em Create.
Etapa 3: Escrever o código do job e implantá-lo
-
Crie uma origem chamada datagen_source para gerar dados simulados e outra origem chamada source_adbpg para capturar alterações de dados em tempo real do banco de dados AnalyticDB for PostgreSQL. Em seguida, faça o join das duas origens e grave os resultados em uma tabela de destino chamada sink_adbpg. Os dados processados serão gravados no AnalyticDB for PostgreSQL.
Copie o seguinte código do job para o editor.
---Create a Datagen source table to generate streaming data using the Datagen connector. CREATE TEMPORARY TABLE datagen_source ( id INT, score INT ) WITH ( 'connector' = 'datagen', 'fields.id.kind'='sequence', 'fields.id.start'='1', 'fields.id.end'='100', 'fields.score.kind'='random', 'fields.score.min'='70', 'fields.score.max'='100' ); --Create an adbpg source table to capture data changes of the adbpg_source_table table based on slot.name and pgoutput using the adbpg-cdc connector. CREATE TEMPORARY TABLE source_adbpg( id int, username varchar, PRIMARY KEY(id) NOT ENFORCED ) WITH( 'connector' = 'adbpg-cdc', 'hostname' = 'gp-bp16v8cgx46ns****-master.gpdb.rds.aliyuncs.com', 'port' = '5432', 'username' = 'account****', 'password' = 'password****', 'database-name' = 'testdb', 'schema-name' = 'testschema', 'table-name' = 'adbpg_source_table', 'slot.name' = 'flink', 'decoding.plugin.name' = 'pgoutput' ); --Create an adbpg sink table to write the processed results to the destination table adbpg_sink_table in the database. CREATE TEMPORARY TABLE sink_adbpg ( id int, username varchar, score int ) WITH ( 'connector' = 'adbpg', 'url' = 'jdbc:postgresql://gp-bp16v8cgx46ns****-master.gpdb.rds.aliyuncs.com:5432/testdb', 'tablename' = 'testschema.adbpg_sink_table', 'username' = 'account****', 'password' = 'password****', 'maxRetryTimes' = '2', 'batchsize' = '5000', 'conflictMode' = 'ignore', 'writeMode' = 'insert', 'retryWaitTime' = '200' ); -- Write the join results of the datagen_source and source_adbpg tables to the adbpg sink table. INSERT INTO sink_adbpg SELECT ts.id,ts.username,ds.score FROM datagen_source AS ds JOIN source_adbpg AS ts ON ds.id = ts.id;Parâmetros
Parâmetro
Obrigatório
Tipo de dado
Descrição
connector
Sim
STRING
Tipo do conector. Defina o valor como
adbpg-cdcpara a tabela de origem eadbpgpara a tabela de destino.hostname
Sim
STRING
Endpoint interno da instância do AnalyticDB for PostgreSQL. Obtenha o internal endpoint na página Basic Information da instância.
username
Sim
STRING
Conta e senha do banco de dados da instância do AnalyticDB for PostgreSQL.
password
Sim
STRING
database-name
Sim
STRING
Nome do banco de dados.
schema-name
Sim
STRING
Nome do schema. Este parâmetro aceita expressões regulares, permitindo a assinatura de múltiplos schemas simultaneamente.
table-name
Sim
STRING
Nome da tabela. Suporta expressões regulares para assinar várias tabelas ao mesmo tempo.
port
Sim
INTEGER
Porta do AnalyticDB for PostgreSQL. O valor é fixo em 5432.
decoding.plugin.name
Sim
STRING
Nome do plug-in de decodificação lógica do PostgreSQL. Valor fixo: pgoutput.
slot.name
Sim
STRING
Nome do slot de decodificação lógica.
-
Para tabelas de origem no mesmo job do Flink, utilize o mesmo valor para
slot.name. -
Se jobs diferentes do Flink envolverem a mesma tabela, defina um
slot.nameexclusivo para cada job. Isso evita o erro:PSQLException: ERROR: replication slot "debezium" is active for PID 974.
debezium.*
Não
STRING
Controla o comportamento do cliente Debezium com maior granularidade. Por exemplo, definir
'debezium.snapshot.mode' = 'never'desativa o recurso de snapshot. Para mais detalhes, consulte as propriedades de configuração.scan.incremental.snapshot.enabled
Não
BOOLEAN
Define se snapshots incrementais devem ser ativados. Valores válidos:
-
false (padrão): Snapshots incrementais desativados.
-
true: Snapshots incrementais ativados.
scan.startup.mode
Não
STRING
Modo de inicialização para consumo de dados. Valores válidos:
-
initial (padrão): Na primeira execução do job, todos os dados históricos são verificados antes da leitura dos dados mais recentes do WAL (write-ahead logging). Isso proporciona uma transição contínua entre dados completos e incrementais.
-
latest-offset: O job ignora dados históricos na primeira execução e começa a ler a partir do final do WAL (posição mais recente do log). Captura apenas alterações ocorridas após o início do conector.
-
snapshot: Verifica todos os dados históricos e lê novas entradas WAL geradas durante essa varredura completa. O job é encerrado ao concluir a varredura total.
changelog-mode
Não
STRING
Modo de changelog para codificar alterações no fluxo. Valores válidos:
-
ALL (padrão): Suporta todos os tipos de operação, incluindo
INSERT,DELETE,UPDATE_BEFOREeUPDATE_AFTER. -
UPSERT: Suporta apenas operações
UPSERT, abrangendoINSERT,DELETEeUPDATE_AFTER.
heartbeat.interval.ms
Não
DURATION
Intervalo para envio de pacotes de heartbeat. O valor padrão é 30 segundos, especificado em milissegundos.
O conector CDC do AnalyticDB for PostgreSQL envia heartbeats ao banco de dados para garantir o avanço contínuo do offset do slot. Se os dados da tabela não mudarem com frequência, configure este parâmetro adequadamente para limpar logs WAL prontamente e evitar desperdício de espaço em disco.
scan.incremental.snapshot.chunk.key-column
Não
STRING
Especifique uma coluna para particionamento (chunking) durante a fase de snapshot. Por padrão, a primeira coluna da chave primária é selecionada.
url
Sim
STRING
Formato:
jdbc:postgresql://<Address>:<PortId>/<DatabaseName>. -
No topo da página de desenvolvimento do job, clique em Deep Check para validar a sintaxe.
Clique em Deploy e, em seguida, clique em Confirm.
No canto superior direito, clique em Operations. Na página Deployments, clique em Start.
Etapa 4: Visualizar os dados gravados pelo Flink
-
Execute as instruções abaixo no banco de dados de teste para visualizar os dados gravados pelo Flink.
SELECT * FROM testschema.adbpg_sink_table; SELECT COUNT(*) FROM testschema.adbpg_sink_table; -
Insira mais 50 linhas de dados na tabela de origem. Depois, verifique o total de linhas incrementais que o Flink gravou na tabela de destino.
-- Insert 50 rows of incremental data into the source table. INSERT INTO testschema.adbpg_source_table(id, username) SELECT i, 'username'||i::text FROM generate_series(51, 100) AS t(i); -- Check the new data in the destination table. SELECT COUNT(*) FROM testschema.adbpg_sink_table where id > 50;O resultado é exibido abaixo.
count ------- 50 (1 row)
Observações de uso
-
Gerencie os Replication Slots prontamente para evitar desperdício de espaço em disco.
Para prevenir perda de dados causada pela limpeza de logs WAL correspondentes a um checkpoint durante o reinício de um job do Flink, o sistema não exclui Replication Slots automaticamente. Portanto, caso confirme que um job do Flink não precisará mais ser reiniciado, exclua manualmente o Replication Slot correspondente para liberar os recursos ocupados. Além disso, se a posição confirmada de um Replication Slot não avançar por um longo período, o AnalyticDB for PostgreSQL não conseguirá limpar as entradas WAL posteriores a essa posição. Isso pode causar acúmulo de dados WAL não utilizados, consumindo grande quantidade de espaço em disco.
Durante a operação normal de uma instância do AnalyticDB for PostgreSQL, a semântica de processamento de dados exactly-once é garantida. No entanto, em cenários de falha, apenas a semântica at-least-once tem suporte.
-
O conector CDC altera o parâmetro REPLICA IDENTITY da tabela assinada para
FULLa fim de garantir a consistência da sincronização de dados. Essa alteração gera os seguintes efeitos:Maior uso de espaço em disco. Em cenários com operações frequentes de atualização ou exclusão, essa configuração aumenta o tamanho dos logs WAL, resultando em maior consumo de disco.
Redução no desempenho de gravação. O desempenho pode ser significativamente afetado em cenários de escrita com alta concorrência.
Pressão elevada nos checkpoints. Logs WAL maiores exigem que os checkpoints processem mais dados, o que pode prolongar o tempo necessário para sua conclusão.
Melhores práticas
O Flink CDC suporta o desenvolvimento de jobs usando a API Flink SQL ou a API DataStream. É possível utilizar o Flink CDC para implementar sincronização integrada de dados completos e incrementais para tabelas únicas ou múltiplas em um banco de dados de origem. Também é viável realizar computações, como joins de tabelas, em fontes de dados distintas. O framework Flink garante a semântica de processamento de eventos exactly-once durante todo o procedimento de processamento de dados. Contudo, o Flink CDC não é adequado para sincronizar um banco de dados inteiro compatível com PostgreSQL, pois não suporta sincronização de DDL e exige a definição da estrutura de cada tabela no Flink SQL, tornando a manutenção complexa.
Esta seção utiliza um exemplo de sincronização de dados do AnalyticDB for PostgreSQL para o Kafka para descrever as melhores práticas no desenvolvimento de jobs Flink CDC SQL. Antes de desenvolver o job Flink CDC, certifique-se de ter preparado e configurado os recursos conforme descrito na seção Prerequisites.
Etapa 1: Preparar tabelas de teste
Na instância do AnalyticDB for PostgreSQL, crie duas tabelas de origem.
CREATE TABLE products (
product_id SERIAL PRIMARY KEY,
product_name VARCHAR(200) NOT NULL,
sku CHAR(12) NOT NULL,
description TEXT,
price NUMERIC(10,2) NOT NULL,
discount_price DECIMAL(10,2),
stock_quantity INTEGER DEFAULT 0,
weight REAL,
volume DOUBLE PRECISION,
dimensions BOX,
release_date DATE,
is_featured BOOLEAN DEFAULT FALSE,
rating FLOAT,
warranty_period INTERVAL,
metadata JSON,
tags TEXT[]
);
CREATE TABLE documents (
document_id UUID PRIMARY KEY,
title VARCHAR(200) NOT NULL,
content TEXT,
summary TEXT,
publication_date TIMESTAMP WITHOUT TIME ZONE,
last_updated TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
author_id BIGINT,
file_data BYTEA,
xml_content XML,
json_metadata JSON,
reading_time INTERVAL,
is_public BOOLEAN DEFAULT TRUE,
views_count INTEGER DEFAULT 0,
category VARCHAR(50),
tags TEXT[]
);
Etapa 2: Preparar recursos do Kafka
Adicione o bloco CIDR do workspace do Flink à lista de permissões da instância do Kafka.
Etapa 3: Criar um job do Flink
Faça logon no console do Realtime Compute. Na aba Fully Managed Flink, localize o workspace desejado e clique em Console na coluna Actions.
No painel de navegação à esquerda, escolha .
Na barra de menu superior, clique em
. Selecione New Blank Stream Draft e configure os parâmetros do job.Clique em Create.
Etapa 4: Escrever o código do job e implantá-lo
-
Escreva um job SQL no workspace do Flink. Copie o código do job abaixo para o editor e substitua as configurações pelos seus valores reais.
-- Use one source to capture data from multiple tables CREATE TEMPORARY TABLE ADBPGSource( table_name STRING METADATA FROM 'table_name' VIRTUAL, row_kind STRING METADATA FROM 'row_kind' VIRTUAL, product_id BIGINT, product_name STRING, sku STRING, description STRING, price STRING, discount_price STRING, stock_quantity INT, weight STRING, volume STRING, dimensions STRING, release_date STRING, is_featured BOOLEAN, rating FLOAT, warranty_period STRING, metadata STRING, tags STRING, document_id STRING, title STRING, content STRING, summary STRING, publication_date STRING, last_updated STRING, author_id BIGINT, file_data STRING, xml_content STRING, json_metadata STRING, reading_time STRING, is_public BOOLEAN, views_count INT, category STRING ) WITH ( 'connector' = 'adbpg-cdc', 'hostname' = 'gp-2zev887z58390***-master.gpdb.rds.aliyuncs.com', 'port' = '5432', 'username' = 'account****', 'password' = 'password****', 'database-name' = 'testdb', 'schema-name' = 'public', 'table-name' = '(products|documents)', 'slot.name' = 'flink', 'decoding.plugin.name' = 'pgoutput', 'debezium.snapshot.mode' = 'never' ); CREATE TEMPORARY TABLE KafkaProducts ( product_id BIGINT, product_name STRING, sku STRING, description STRING, price STRING, discount_price STRING, stock_quantity INT, weight STRING, volume STRING, dimensions STRING, release_date STRING, is_featured BOOLEAN, rating FLOAT, warranty_period STRING, metadata STRING, tags STRING, PRIMARY KEY(product_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = '****', 'properties.bootstrap.servers' = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092', 'key.format'='avro', 'value.format'='avro' ); CREATE TEMPORARY TABLE KafkaDocuments ( document_id STRING, title STRING, content STRING, summary STRING, publication_date STRING, last_updated STRING, author_id BIGINT, file_data STRING, xml_content STRING, json_metadata STRING, reading_time STRING, is_public BOOLEAN, views_count INT, category STRING, tags STRING, PRIMARY KEY(document_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = '****', 'properties.bootstrap.servers' = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092', 'key.format'='avro', 'value.format'='avro' ); -- Use a STATEMENT SET to wrap multiple statements BEGIN STATEMENT SET; -- Use the table_name METADATA to route data to the destination table INSERT INTO KafkaProducts SELECT product_id,product_name,sku,description,price,discount_price,stock_quantity,weight,volume,dimensions,release_date,is_featured,rating,warranty_period,metadata,tags FROM ADBPGSource WHERE table_name = 'products'; INSERT INTO KafkaDocuments SELECT document_id,title,content,summary,publication_date,last_updated,author_id,file_data,xml_content,json_metadata,reading_time,is_public,views_count,category,tags FROM ADBPGSource WHERE table_name = 'documents'; END;Observe os seguintes pontos sobre este job SQL:
Para tarefas de sincronização de múltiplas tabelas, recomenda-se usar uma única tabela de origem para capturar dados de várias tabelas, conforme demonstrado neste exemplo SQL. Todas as colunas de todas as tabelas de origem devem ser definidas nesta tabela única. Se houver nomes de colunas duplicados, mantenha apenas um. Ao gravar na tabela de destino, utilize o
METADATAtable_name para rotear os dados para a tabela especificada. Essa abordagem exige a criação de apenas um Replication Slot no AnalyticDB for PostgreSQL, reduzindo o uso de recursos do banco de dados de origem, melhorando o desempenho da sincronização e simplificando a manutenção futura.Utilize o parâmetro
table-namepara especificar múltiplas tabelas de origem. Coloque os nomes das tabelas entre parênteses e separe-os por barras verticais (|), por exemplo:(table1|table2|table3).Definir
debezium.snapshot.modecomoneversignifica que apenas dados incrementais da tabela de origem serão sincronizados. Para sincronizar dados completos e incrementais, altere a configuração parainitial.
No topo da página de desenvolvimento do job, clique em Deep Check para validar a sintaxe.
Clique em Deploy e, em seguida, clique em OK.
No canto superior direito, clique em Operations. Na página Deployments, clique em Start.
Etapa 5: Inserir dados de teste
Na instância do AnalyticDB for PostgreSQL, atualize os dados nas duas tabelas de origem e observe the message changes no tópico do Kafka.
Utilize a instrução SQL abaixo para inserir dados de teste:
INSERT INTO products (
product_name, sku, description, price, discount_price, stock_quantity, weight, volume, dimensions, release_date, is_featured, rating, warranty_period, metadata, tags
) VALUES (
'Test Product', 'Test-2025', 'A piece of test product data', 299.99, 279.99, 150, 50.5, 120.75, '(10,20),(30,40)', '2023-05-01', TRUE, 4.8, INTERVAL '1 year', '{"brand": "TechCo", "model": "X1"}', '{"Test1", "Test2"}'
);
Referências
Para mais informações sobre como assinar os dados completos do AnalyticDB for PostgreSQL, consulte Read and write full data in real time using Flink.