AnalyticDB for PostgreSQL, um data warehouse nativo da cloud, permite a integração de dados vetoriais por meio do flink-adbpg-connector. Este tópico utiliza um exemplo de importação de dados do Message Queue for Apache Kafka para demonstrar como carregar dados vetoriais no AnalyticDB for PostgreSQL.
Pré-requisitos
Uma instância do AnalyticDB for PostgreSQL deve estar criada. Para mais informações, consulte Crie uma instância.
-
Um workspace do Flink totalmente gerenciado deve estar criado na mesma VPC da instância do AnalyticDB for PostgreSQL. Para mais informações, consulte Ative o Realtime Compute for Apache Flink.
Caso utilize um cluster Flink open source autogerenciado, verifique se o flink-adbpg-connector está instalado no diretório
$FLINK_HOME/lib.Se você utilizar o service totalmente gerenciado, nenhuma ação é necessária.
-
A extensão de recuperação vetorial FastANN deve estar instalada no seu banco de dados AnalyticDB for PostgreSQL.
Execute o comando
\dx fastannem um cliente psql para verificar se a extensão está instalada.Se o sistema retornar informações sobre a extensão FastANN, ela já está instalada.
Caso nenhuma informação seja retornada, Envie um ticket para solicitar a instalação ao suporte técnico.
Uma instância do Message Queue for Apache Kafka deve estar adquirida e implantada na mesma VPC da instância do AnalyticDB for PostgreSQL. Para mais informações, consulte Adquira e implante uma instância.
Os blocos CIDR do workspace do Flink e da instância do Kafka devem estar adicionados à lista de permissões de endereços IP da instância do AnalyticDB for PostgreSQL. Para mais informações, consulte Configure uma lista de permissões de endereços IP.
Dados de exemplo
O AnalyticDB for PostgreSQL disponibiliza dados de exemplo para testes. Para baixar os dados, clique em vector_sample_data.csv.
A tabela a seguir descreve o esquema dos dados de exemplo.
|
Campo |
Tipo |
Descrição |
|
id |
bigint |
O ID do carro. |
|
market_time |
timestamp |
A data de lançamento do carro no mercado. |
|
color |
varchar(10) |
A cor do carro. |
|
price |
int |
O preço do carro. |
|
feature |
float4[] |
O vetor de características da imagem do carro. |
Procedimento
Criar índices estruturados e vetoriais
Conecte-se ao seu banco de dados AnalyticDB for PostgreSQL. Os passos a seguir utilizam um cliente psql. Para mais informações, consulte Conecte-se a um banco de dados usando psql.
-
Execute as instruções abaixo para criar e acessar um banco de dados de teste:
CREATE DATABASE adbpg_test; \c adbpg_test -
Execute as instruções a seguir para criar uma tabela de destino:
CREATE SCHEMA IF NOT EXISTS vector_test; CREATE TABLE IF NOT EXISTS vector_test.car_info ( id bigint NOT NULL, market_time timestamp, color varchar(10), price int, feature float4[], PRIMARY KEY(id) ) DISTRIBUTED BY(id); -
Execute as instruções abaixo para criar os índices estruturados e vetoriais:
-- Change the storage format of the vector column to PLAIN. ALTER TABLE vector_test.car_info ALTER COLUMN feature SET STORAGE PLAIN; -- Create structured indexes. CREATE INDEX ON vector_test.car_info(market_time); CREATE INDEX ON vector_test.car_info(color); CREATE INDEX ON vector_test.car_info(price); -- Create a vector index. CREATE INDEX ON vector_test.car_info USING ann(feature) WITH (dim='10', pq_enable='0');
Gravar dados vetoriais de exemplo no Kafka
-
Execute o comando a seguir para criar um tópico do Kafka:
bin/kafka-topics.sh --create --topic vector_ingest --partitions 1 \ --bootstrap-server <your_broker_list> -
Execute o comando abaixo para gravar os dados vetoriais de exemplo no tópico do Kafka:
bin/kafka-console-producer.sh \ --bootstrap-server <your_broker_list> \ --topic vector_ingest < ../vector_sample_data.csv
<your_broker_list>: O endpoint da instância. Obtenha o endpoint na seção Access Point Information da página Instance Details no console do Message Queue for Apache Kafka.
Criar uma tabela de mapeamento e importar dados
-
Crie um job do Flink.
Faça login no console do Realtime Compute for Apache Flink. Na aba Fully Managed Flink, localize o workspace desejado e clique em Console na coluna Actions.
No painel de navegação à esquerda, clique em SQL Development. Clique em New, selecione Blank stream draft e, em seguida, clique em Next.
-
Na caixa de diálogo New draft, defina os parâmetros do rascunho.
Parâmetro
Descrição
Exemplo
File Name
Nome do rascunho.
NotaO nome do rascunho deve ser único no projeto atual.
adbpg-test
Storage Location
Pasta onde o arquivo de código do rascunho será armazenado.
Também é possível clicar no ícone
ao lado de uma pasta existente para criar uma subpasta.Drafts
Engine Version
Versão do mecanismo Flink para o job atual. Para mais detalhes sobre versões do mecanismo, mapeamentos de versão e marcos do ciclo de vida, consulte Engine versions.
vvr-6.0.6-flink-1.15
-
Execute a instrução a seguir para criar uma tabela de mapeamento para o AnalyticDB for PostgreSQL:
CREATE TABLE vector_ingest ( id INT, market_time TIMESTAMP, color VARCHAR(10), price int, feature VARCHAR )WITH ( 'connector' = 'adbpg-nightly-1.13', 'url' = 'jdbc:postgresql://<your_instance_url>:5432/adbpg_test', 'tablename' = 'car_info', 'username' = '<your_username>', 'password' = '<your_password>', 'targetschema' = 'vector_test', 'maxretrytimes' = '2', 'batchsize' = '3000', 'batchwritetimeoutms' = '10000', 'connectionmaxactive' = '20', 'conflictmode' = 'ignore', 'exceptionmode' = 'ignore', 'casesensitive' = '0', 'writemode' = '1', 'retrywaittime' = '200' );Para descrições dos parâmetros, consulte Grave dados no AnalyticDB for PostgreSQL.
-
Execute a instrução abaixo para criar uma tabela de mapeamento para o Kafka:
CREATE TABLE vector_kafka ( id INT, market_time TIMESTAMP, color VARCHAR(10), price int, feature string ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = '<your_broker_list>', 'topic' = 'vector_ingest', 'format' = 'csv', 'csv.field-delimiter' = '\t', 'scan.startup.mode' = 'earliest-offset' );A tabela a seguir descreve os parâmetros.
Parâmetro
Obrigatório
Descrição
connector
Sim
Nome do conector. O valor deve ser
kafka.properties.bootstrap.servers
Sim
Endpoint da instância do Message Queue for Apache Kafka. Você pode obter o endpoint na seção de informações de endpoint da página de detalhes da instância no console do Message Queue for Apache Kafka.
topic
Sim
Nome do tópico do Kafka.
format
Sim
Formato dos valores das mensagens do Kafka. Os seguintes formatos são suportados:
-
csv
-
json
-
avro
-
debezium-json
-
canal-json
-
maxwell-json
-
avro-confluent
-
raw
csv.field-delimiter
Sim
Delimitador de campos para o formato CSV.
scan.startup.mode
Sim
Define o offset a partir do qual o consumidor do Kafka começa a ler os dados. Valores válidos:
-
earliest-offset: Inicia a leitura a partir do offset mais antigo disponível. -
latest-offset: Inicia a leitura a partir do offset mais recente.
-
-
Execute a instrução a seguir para criar um job de importação:
INSERT INTO vector_ingest SELECT * FROM vector_kafka;