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:
Uma instância e uma tabela do AnalyticDB for PostgreSQL. Consulte Criar uma instância e CREATE TABLE
Uma lista de permissões de endereços IP configurada para a instância. Consulte Configurar uma lista de permissões de endereços IP
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 |
|
|
Sim |
— |
STRING |
Defina como |
|
|
Sim |
— |
STRING |
URL JDBC (Java Database Connectivity) no formato |
|
|
Sim |
— |
STRING |
Nome da tabela no banco de dados. |
|
|
Sim |
— |
STRING |
Nome de usuário do banco de dados AnalyticDB for PostgreSQL. |
|
|
Sim |
— |
STRING |
Senha do banco de dados AnalyticDB for PostgreSQL. |
|
|
Não |
|
INTEGER |
Número máximo de tentativas de nova execução após falha na gravação. |
|
|
Não |
|
STRING |
Nome do schema do banco de dados. |
|
|
Não |
|
STRING |
Define se os identificadores diferenciam maiúsculas de minúsculas. Valores válidos: |
|
|
Não |
|
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 |
|
|
Sim |
— |
STRING |
Nome do schema. Suporta expressões regulares para assinar vários schemas simultaneamente. |
|
|
Sim |
|
INTEGER |
Porta da instância do AnalyticDB for PostgreSQL. |
|
|
Sim |
|
STRING |
Plug-in de decodificação lógica do PostgreSQL. Defina como |
|
|
Sim |
— |
STRING |
Nome do slot de decodificação lógica. Consulte Orientações sobre nomes de slots. |
|
|
Sim |
— |
STRING |
Propriedades de configuração do cliente Debezium. Por exemplo, defina |
|
|
Não |
|
BOOLEAN |
Define se o snapshot incremental deve ser ativado. Valores válidos: |
|
|
Não |
|
STRING |
Modo de inicialização para consumo de dados. Consulte Modos de inicialização. |
|
|
Não |
|
STRING |
Forma como os eventos de alteração são codificados no fluxo de mudanças. Valores válidos: |
|
|
Não |
30 segundos |
DURATION |
Intervalo para envio de pacotes de heartbeat, em milissegundos. Consulte Heartbeat e retenção de WAL. |
|
|
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 |
|
|
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). |
|
|
Ignora a varredura de dados históricos e lê apenas a partir do final atual do WAL. |
|
|
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 |
|
|
|
|
|
|
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.namepara todas as tabelas source do job.Entre jobs diferentes do Flink: Atribua um
slot.nameexclusivo para cada job. Reutilizar o nome de um slot entre jobs causa o erroPSQLException: 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 |
|
|
Não |
|
STRING |
Estratégia de gravação para a primeira tentativa. Consulte Modos de gravação e resolução de conflitos. |
|
|
Não |
|
STRING |
Define como tratar conflitos de chave primária ou índice. Consulte Modos de gravação e resolução de conflitos. |
|
|
Não |
|
INTEGER |
Quantidade de registros gravados por lote. |
|
|
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 |
|
|
Não |
|
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.
|
** |
Descrição |
Observações |
|
|
Insere registros diretamente. O tratamento de conflitos segue a definição de |
Adequado para a maioria dos casos de uso. |
|
|
Atualiza automaticamente os registros existentes em caso de conflito. |
Exige chave primária. |
|
|
Insere registros usando o comando |
Requer VVR 11.1 ou posterior. |
|
** |
Comportamento em caso de conflito |
Observações |
|
|
Reporta um erro. |
|
|
|
Descarta silenciosamente o registro conflitante. |
|
|
|
Atualiza o registro conflitante. |
Recomendado para tabelas sem chave primária. Reduz o throughput de gravação. |
|
|
Atualiza o registro conflitante. |
Exige chave primária. Mais eficiente que |
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 |
|
|
Não |
|
STRING |
Política de cache para consultas à tabela de dimensão. Consulte Políticas de cache. |
|
|
Não |
|
LONG |
Número máximo de linhas mantidas no cache. Tem efeito apenas quando |
|
|
Não |
|
LONG |
Tempo de vida das entradas de cache em milissegundos. O comportamento depende da configuração de |
|
|
Não |
|
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 |
Comportamento |
Quando usar |
|
|
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 ( |
Tabelas de dimensão pequenas ou médias, nas quais a atualização dos dados tolera recargas periódicas. |
|
|
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 |
Tabelas de dimensão grandes, onde armazenar tudo em cache é inviável, ou quando o cache parcial de baixa latência é aceitável. |
|
|
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 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
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 |
|
|
Total de registros gravados no sink. |
|
|
Registros gravados por segundo. |
|
|
Total de bytes gravados no sink. |
|
|
Bytes gravados por segundo. |
|
|
Tempo gasto na operação de gravação mais recente. |