Este tópico descreve como gravar dados no AnalyticDB for PostgreSQL com o Realtime Compute for Apache Flink.
Limites de uso
Este recurso não oferece suporte ao modo Serverless do AnalyticDB for PostgreSQL.
O conector do AnalyticDB for PostgreSQL é compatível apenas com o Realtime Compute for Apache Flink VVR 6.0.0 ou posterior.
-
O AnalyticDB for PostgreSQL V7.0 requer o Realtime Compute for Apache Flink VVR 8.0.1 ou superior.
NotaSe você utilizar um conector personalizado, consulte Gerencie conectores personalizados para obter mais informações.
Pré-requisitos
Crie um workspace do Fully Managed Flink. Para mais informações, consulte Ative o Realtime Compute for Apache Flink.
Crie uma instância do AnalyticDB for PostgreSQL. Para mais informações, consulte Crie uma instância.
A instância do AnalyticDB for PostgreSQL e o workspace do Fully Managed Flink devem estar na mesma VPC.
Configure uma instância do AnalyticDB for PostgreSQL
Faça login no console do AnalyticDB for PostgreSQL.
Adicione o bloco CIDR do workspace Flink à lista de permissões da instância do AnalyticDB for PostgreSQL. Para saber como configurar a lista de permissões, consulte Configure uma whitelist.
Clique em Log On to Database. Para conhecer outras formas de conexão com o banco de dados, consulte Conecte-se a uma instância usando um cliente.
-
Crie uma tabela na instância do AnalyticDB for PostgreSQL.
Exemplo de instrução SQL para criação de tabela:
CREATE TABLE test_adbpg_table( b1 int, b2 int, b3 text, PRIMARY KEY(b1) );
Configure o Realtime Compute for Flink
Acesse o console do Realtime Compute.
Na aba Fully Managed Flink, clique em Console na coluna Actions do workspace desejado.
No painel de navegação à esquerda, clique em Connectors.
Na página Connectors, clique em Create Custom Connector.
-
Envie o arquivo JAR do conector personalizado.
NotaObtenha o pacote JAR do conector Flink personalizado para AnalyticDB for PostgreSQL. Para mais informações, consulte AnalyticDB PostgreSQL Connector.
A versão do pacote JAR deve corresponder à versão do mecanismo Flink da plataforma Realtime Compute.
-
Após concluir o envio, clique em Next.
O sistema analisa o conector personalizado enviado. Se a análise for bem-sucedida, prossiga para a próxima etapa. Em caso de falha, verifique se o código do conector personalizado está em conformidade com os padrões da comunidade Apache Flink.
-
Clique em Finish.
O conector personalizado criado aparecerá na lista de conectores.
Crie um job Flink
Acesse o console do Realtime Compute. Na aba Fully Managed Flink, clique em Console na coluna Actions do workspace desejado.
No painel de navegação à esquerda, clique em SQL Development. Clique em New, selecione Blank Streaming Job Draft e clique em Next.
-
Na caixa de diálogo New File Draft, configure os parâmetros do job.
Parâmetro do Job
Descrição
Exemplo
Name
Nome do job.
NotaO nome do job deve ser exclusivo no projeto atual.
adbpg-test
Location
Pasta onde o arquivo de código do job será armazenado.
Você também pode clicar no ícone
ao lado de uma pasta existente para criar uma subpasta.Job Drafts
Engine Version
Versão do mecanismo Flink utilizada pelo job atual. Para detalhes sobre numeração de versões, mapeamentos e ciclo de vida, consulte Engine versions.
vvr-6.0.7-flink-1.15
Clique em Create.
Gravar dados no AnalyticDB for PostgreSQL
-
Escreva o código do job.
Crie uma tabela de origem aleatória
datagen_sourcee uma tabela de destinotest_adbpg_tableno AnalyticDB for PostgreSQL. Copie o código abaixo para o editor de texto do job.CREATE TABLE datagen_source ( f_sequence INT, f_random INT, f_random_str STRING ) WITH ( 'connector' = 'datagen', 'rows-per-second'='5', 'fields.f_sequence.kind'='sequence', 'fields.f_sequence.start'='1', 'fields.f_sequence.end'='1000', 'fields.f_random.min'='1', 'fields.f_random.max'='1000', 'fields.f_random_str.length'='10' ); CREATE TABLE test_adbpg_table ( `B1` bigint , `B2` bigint , `B3` VARCHAR , `B4` VARCHAR, PRIMARY KEY(B1) not ENFORCED ) with ( 'connector' = 'adbpg-nightly-1.13', 'password' = 'xxx', 'tablename' = 'test_adbpg_table', 'username' = 'xxxx', 'url' = 'jdbc:postgresql://url:5432/schema', 'maxretrytimes' = '2', 'batchsize' = '50000', 'connectionmaxactive' = '5', 'conflictmode' = 'ignore', 'usecopy' = '0', 'targetschema' = 'public', 'exceptionmode' = 'ignore', 'casesensitive' = '0', 'writemode' = '1', 'retrywaittime' = '200' );Não é necessário alterar os parâmetros da tabela
datagen_source. No entanto, ajuste os parâmetros da tabelatest_adbpg_tableconforme suas necessidades reais de negócio. A tabela a seguir detalha esses parâmetros.Parâmetro
Obrigatório
Descrição
connector
Sim
Nome do conector. Defina este parâmetro como
adbpg-nightly-<número da versão>, por exemplo,adbpg-nightly-1.13.url
Sim
URL JDBC da instância do AnalyticDB for PostgreSQL. Formato:
jdbc:postgresql://<endpoint interno>:<porta>/<nome do banco de dados>. Exemplo:jdbc:postgresql://gp-xxxxxx.gpdb.cn-chengdu.rds.aliyuncs.com:5432/postgres.tablename
Sim
Nome da tabela no AnalyticDB for PostgreSQL.
username
Sim
Conta de banco de dados da instância do AnalyticDB for PostgreSQL.
password
Sim
Senha da conta de banco de dados da instância do AnalyticDB for PostgreSQL.
maxretrytimes
Não
Número máximo de tentativas após falha na execução de SQL. Valor padrão: 3.
batchsize
Não
Quantidade máxima de registros de dados a serem gravados em um único lote. Valor padrão: 50.000.
exceptionmode
Não
Política de tratamento de erros quando ocorre uma exceção durante a gravação de dados. Valores válidos:
-
ignore: ignora os dados que causaram a exceção. Este é o valor padrão.
-
strict: aciona um failover e reporta um erro ao detectar uma exceção na gravação.
conflictmode
Não
Política para lidar com conflitos de chave primária ou índice único. Valores válidos:
-
ignore: ignora conflitos de chave primária e mantém os dados existentes.
-
strict: aciona um failover e reporta um erro em caso de conflito de chave primária.
-
update: atualiza os dados quando há conflito de chave primária.
-
upsert: grava dados usando o método UPSERT quando ocorre conflito de chave primária. Este é o valor padrão.
O AnalyticDB for PostgreSQL implementa o UPSERT utilizando INSERT ON CONFLICT e COPY ON CONFLICT. Se a tabela de destino for particionada, a versão secundária do kernel deve ser V6.3.6.1 ou superior. Para instruções sobre como atualizar a versão secundária do kernel, consulte Atualize a versão do mecanismo.
targetschema
Não
Schema do banco de dados AnalyticDB for PostgreSQL. Valor padrão: public.
writemode
Não
Modo de gravação de dados. Valores válidos:
-
0: grava dados usando BATCH INSERT.
-
1: grava dados usando a API COPY. Este é o valor padrão.
-
2: grava dados usando BATCH UPSERT.
verbose
Não
Define se os logs de runtime do conector devem ser gerados. Valores válidos:
-
0: não gera logs de runtime. Este é o valor padrão.
-
1: gera logs de runtime.
retrywaittime
Não
Intervalo entre tentativas de nova execução quando ocorre uma exceção. Unidade: milissegundos. Valor padrão: 100.
batchwritetimeoutms
Não
Tempo máximo de acumulação para gravação em lote. Ao exceder esse tempo, o lote acumulado é gravado. Unidade: milissegundos. Valor padrão: 50.000.
connectionmaxactive
Não
Parâmetro do pool de conexões. Especifica o número máximo de conexões simultâneas no pool de um único Task Manager. Valor padrão: 5.
casesensitive
Não
Define se nomes de colunas e tabelas diferenciam maiúsculas de minúsculas. Valores válidos:
-
0: não diferencia maiúsculas de minúsculas. Este é o valor padrão.
-
1: diferencia maiúsculas de minúsculas.
NotaHá suporte para mapeamentos de parâmetros e tipos. Para mais informações, consulte a documentação do conector em AnalyticDB for PostgreSQL (ADB PG).
-
-
Inicie o job.
-
Na parte superior da página de desenvolvimento do job, clique em Deploy. Na caixa de diálogo exibida, clique em OK.
NotaClusters de sessão são adequados para desenvolvimento e testes em ambientes de não produção. Eles permitem depurar jobs, melhorar a utilização de recursos do Job Manager (JM) e acelerar a inicialização. Contudo, não recomendamos enviar jobs de produção para clusters de sessão devido a preocupações com estabilidade. Para mais detalhes, veja Depure um job.
Na página Deployments, clique em Resume na coluna Actions do job desejado.
Clique em Resume.
-
Verifique os resultados
Conecte-se ao banco de dados AnalyticDB for PostgreSQL. Consulte Conecte-se a uma instância usando um cliente para instruções.
-
Execute a seguinte instrução para consultar a tabela
test_adbpg_table.SELECT * FROM test_adbpg_table;Os dados foram gravados no AnalyticDB for PostgreSQL conforme esperado. A figura abaixo mostra um exemplo de retorno.
O resultado da consulta contém três colunas:
b1(int4),b2(int4) eb3(text). Várias linhas de dados são retornadas, confirmando que a sincronização para a tabela de destino foi bem-sucedida.