Este tutorial demonstra como usar o AnalyticDB for PostgreSQL como tabela de dimensão e tabela de resultados em um job Flink SQL. Esse é um padrão comum em pipelines de enriquecimento em tempo real.
Ao final, você terá um job Flink em execução que lê dados de uma source Datagen, consulta informações de usuários em uma tabela de dimensão do AnalyticDB for PostgreSQL e grava os registros enriquecidos em uma tabela de resultados do AnalyticDB for PostgreSQL.
Limitações
O Realtime Compute for Apache Flink não oferece suporte à leitura de dados do AnalyticDB for PostgreSQL no modo serverless.
O conector do AnalyticDB for PostgreSQL requer o Ververica Runtime (VVR) 6.0.0 ou superior.
O AnalyticDB for PostgreSQL V7.0 requer o VVR 8.0.1 ou posterior.
Para usar um conector personalizado, consulte Gerencie conectores personalizados .
Pré-requisitos
Antes de começar, verifique se você tem:
Um workspace do Flink totalmente gerenciado. Consulte Ative o Flink totalmente gerenciado.
Uma instância do AnalyticDB for PostgreSQL e uma conta privilegiada. Consulte Crie uma instância e Crie uma conta privilegiada.
A instância do AnalyticDB for PostgreSQL e o workspace do Flink totalmente gerenciado na mesma Virtual Private Cloud (VPC).
Se os recursos estiverem em VPCs diferentes, consulte Como o Flink totalmente gerenciado acessa um service entre VPCs?
Etapa 1: Configure a lista de permissões e prepare os dados
Faça logon no console do AnalyticDB for PostgreSQL.
-
Adicione o bloco CIDR do workspace do Flink totalmente gerenciado à lista de permissões da instância do AnalyticDB for PostgreSQL.
Identifique o bloco CIDR do vSwitch usado pelo seu workspace do Flink totalmente gerenciado. Consulte Como configure uma whitelist?
Adicione esse bloco CIDR à lista de permissões da instância do AnalyticDB for PostgreSQL. Consulte Procedimento.
Para acessar a instância pela Internet, adicione o endereço IP público.
Na página de detalhes da instância, clique em Log On to Database no canto superior direito e insira seu nome de usuário e senha. Para obter mais detalhes, consulte Use ferramentas de cliente para conectar-se a uma instância.
-
Crie uma tabela de dimensão chamada
adbpg_dim_tablee insira 50 linhas de dados de exemplo.-- Create the dimension table CREATE TABLE adbpg_dim_table( id int, username text, PRIMARY KEY(id) ); -- Insert 50 rows: id ranges from 1 to 50, username is "username" followed by the row number INSERT INTO adbpg_dim_table(id, username) SELECT i, 'username'||i::text FROM generate_series(1, 50) AS t(i);Execute
SELECT * FROM adbpg_dim_table ORDER BY id;para verificar os dados inseridos. -
Crie uma tabela de resultados chamada
adbpg_sink_tablepara receber a saída do Flink.CREATE TABLE adbpg_sink_table( id int, username text, score int );
Etapa 2: Crie um rascunho de stream
Faça logon no console do Realtime Compute for Apache Flink, localize seu workspace e clique em Console na coluna Actions.
No painel de navegação à esquerda, acesse Development > ETL. No canto superior esquerdo da página do SQL Editor, clique em + e selecione New Blank Stream Draft.
-
Na caixa de diálogo New Draft, configure os seguintes parâmetros.
Parâmetro
Descrição
Exemplo
Name
Nome do rascunho. Deve ser único no projeto.
adbpg-testLocation
Pasta onde o rascunho será salvo. Clique no ícone ao lado de uma pasta existente para criar uma subpasta.
DraftEngine Version
Versão do mecanismo Flink. Para obter detalhes sobre versões e ciclo de vida, consulte Versões do mecanismo.
vvr-8.0.1-flink-1.17 Clique em Create.
Etapa 3: Escrever e implantar o rascunho
-
Copie o SQL abaixo para o editor de código. Este script define três tabelas e um lookup join que enriquece o stream Datagen com dados de usuários do AnalyticDB for PostgreSQL.
-- Source table: Datagen generates sequential IDs (1-50) and random scores (70-100). -- No changes needed in the WITH clause for this example. CREATE TEMPORARY TABLE datagen_source ( id INT, score INT ) WITH ( 'connector' = 'datagen', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '50', 'fields.score.kind' = 'random', 'fields.score.min' = '70', 'fields.score.max' = '100' ); -- Dimension table: backed by AnalyticDB for PostgreSQL. -- Flink queries this table at processing time to look up usernames by ID. -- Replace the WITH clause values with your actual connection details. CREATE TEMPORARY TABLE dim_adbpg( id int, username varchar, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'connector' = 'adbpg', 'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest', 'tablename' = 'adbpg_dim_table', 'username' = 'flinktest', 'password' = '${secret_values.adb_password}', 'maxRetryTimes' = '2', 'cache' = 'lru', 'cacheSize' = '100' ); -- Result table: Flink writes enriched records here. -- Replace the WITH clause values with your actual connection details. CREATE TEMPORARY TABLE sink_adbpg ( id int, username varchar, score int ) WITH ( 'connector' = 'adbpg', 'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest', 'tablename' = 'adbpg_sink_table', 'username' = 'flinktest', 'password' = '${secret_values.adb_password}', 'maxRetryTimes' = '2', 'conflictMode' = 'ignore', 'retryWaitTime' = '200' ); -- Lookup join: for each record from datagen_source, Flink looks up the matching -- row in dim_adbpg at the time the record is processed (PROCTIME()). INSERT INTO sink_adbpg SELECT ts.id, ts.username, ds.score FROM datagen_source AS ds JOIN dim_adbpg FOR SYSTEM_TIME AS OF PROCTIME() AS ts ON ds.id = ts.id;Sobre a sintaxe do lookup join: A expressão
FOR SYSTEM_TIME AS OF PROCTIME()instrui o Flink a consultar a tabela de dimensão no momento exato em que cada registro da source é processado. Isso significa que o enriquecimento usa os dados de dimensão disponíveis no instante do processamento. Os resultados já gravados não são atualizados caso a tabela de dimensão sofra alterações posteriores. -
Atualize os parâmetros de conexão das tabelas de dimensão e de resultados. Substitua os valores de espaço reservado nas cláusulas
WITHpelos detalhes reais de conexão do seu AnalyticDB for PostgreSQL. A tabela source Datagen não requer alterações. Para obter a referência completa de parâmetros e mapeamentos de tipos de dados, consulte Conector do AnalyticDB for PostgreSQL.Parâmetro
Obrigatório
Padrão
Descrição
urlSim
—
URL JDBC no formato
jdbc:postgresql://<Internal endpoint>:<Port>/<Database name>. Encontre esta informação na página Database Connection da instância no console do AnalyticDB for PostgreSQL.tablenameSim
—
Nome da tabela no banco de dados do AnalyticDB for PostgreSQL.
usernameSim
—
Nome de usuário para acesso ao banco de dados.
passwordSim
—
Senha da conta do banco de dados.
targetSchemaNão
publicNome do schema. Especifique este valor apenas se sua tabela não estiver no schema
public.maxRetryTimesNão
—
Número máximo de tentativas após falha de gravação.
cacheNão
—
Política de cache para consultas à tabela de dimensão. Defina como
lrupara manter as entradas acessadas recentemente na memória. O cache LRU reduz o tráfego do banco de dados e melhora o throughput de consulta, mas as entradas em cache podem ficar desatualizadas. Trata-se de um compromisso entre throughput e atualização dos dados. Ajuste ocacheSizee considere sua tolerância a dados obsoletos antes de ativar.cacheSizeNão
—
Quantidade máxima de entradas no cache. Valores maiores reduzem as solicitações ao banco de dados, mas consomem mais memória.
conflictModeNão
—
Ação executada quando uma gravação entra em conflito com uma chave primária ou índice existente. Defina como
ignorepara ignorar linhas conflitantes.retryWaitTimeNão
—
Tempo de espera em milissegundos entre novas tentativas de gravação.
No canto superior direito da página do SQL Editor, clique em Validate para verificar a sintaxe.
Clique em Deploy.
Na página O&M > Deployments, localize sua implantação e clique em Start na coluna Actions.
Etapa 4: Verifique o resultado
Faça logon no console do AnalyticDB for PostgreSQL.
Clique em Log On to Database. Para obter mais detalhes, consulte Conecte-se a uma instância a partir de um cliente.
-
Execute a consulta abaixo para visualizar os registros que o Flink gravou na tabela de resultados.
SELECT * FROM adbpg_sink_table ORDER BY id;O resultado deve conter 50 linhas. Cada linha deve ter um ID de usuário, o nome de usuário correspondente da tabela de dimensão e uma pontuação aleatória entre 70 e 100.
A consulta retorna 7 registros com três colunas: id, username e score. As linhas de 1 a 7 correspondem a username1 até username7, com pontuações de 94, 79, 70, 93, 71, 82 e 87, respectivamente.