AnalyticDB for PostgreSQL est un entrepôt de données cloud-native qui prend en charge l'intégration de données vectorielles via le connecteur flink-adbpg-connector. Cette rubrique utilise un exemple d'importation de données depuis Message Queue for Apache Kafka pour illustrer le chargement de données vectorielles dans AnalyticDB for PostgreSQL.
Prérequis
Vous avez créé une instance AnalyticDB for PostgreSQL. Pour plus d'informations, consultez la section Créer une instance.
-
Vous avez créé un espace de travail Flink entièrement géré et celui-ci se trouve dans le même VPC que l'instance AnalyticDB for PostgreSQL. Pour plus d'informations, consultez la section Activer Realtime Compute for Apache Flink.
Si vous utilisez un cluster Flink open source auto-géré, assurez-vous que le flink-adbpg-connector est installé dans le répertoire
$FLINK_HOME/lib.Si vous utilisez le service entièrement géré, aucune action n'est requise.
-
L'extension de récupération vectorielle FastANN est installée dans votre base de données AnalyticDB for PostgreSQL.
Exécutez la commande
\dx fastannsur un client psql pour vérifier si l'extension est installée.Si les informations relatives à l'extension FastANN s'affichent, l'extension est installée.
Si aucune information ne s'affiche, soumettez un ticket afin de contacter le support technique pour l'installation.
Vous avez acheté et déployé une instance Message Queue for Apache Kafka dans le même VPC que l'instance AnalyticDB for PostgreSQL. Pour plus d'informations, consultez la section Acheter et déployer une instance.
Vous avez ajouté les plages CIDR de l'espace de travail Flink et de l'instance Kafka à la liste d'autorisation d'adresses IP de l'instance AnalyticDB for PostgreSQL. Pour plus d'informations, consultez la section Configurer une liste d'autorisation d'adresses IP.
Exemple de données
AnalyticDB for PostgreSQL fournit des exemples de données à des fins de test. Pour télécharger les données, cliquez sur vector_sample_data.csv.
Le tableau suivant décrit le schéma des exemples de données.
|
Champ |
Type |
Description |
|
id |
bigint |
L'identifiant de la voiture. |
|
market_time |
timestamp |
L'heure de lancement de la voiture. |
|
color |
varchar(10) |
La couleur de la voiture. |
|
price |
int |
Le prix de la voiture. |
|
feature |
float4[] |
Le vecteur de caractéristiques de l'image de la voiture. |
Procédure
Créer des index structurés et vectoriels
Connectez-vous à votre base de données AnalyticDB for PostgreSQL. Les étapes suivantes utilisent un client psql. Pour plus d'informations, consultez la section Se connecter à une base de données à l'aide de psql.
-
Exécutez les instructions suivantes pour créer une base de données de test et basculer vers celle-ci :
CREATE DATABASE adbpg_test; \c adbpg_test -
Exécutez les instructions suivantes pour créer une table de destination :
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); -
Exécutez les instructions suivantes pour créer des index structurés et vectoriels :
-- 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');
Écrire des exemples de données vectorielles dans Kafka
-
Exécutez la commande suivante pour créer un topic Kafka :
bin/kafka-topics.sh --create --topic vector_ingest --partitions 1 \ --bootstrap-server <your_broker_list> -
Exécutez la commande suivante pour écrire les exemples de données vectorielles dans le topic Kafka :
bin/kafka-console-producer.sh \ --bootstrap-server <your_broker_list> \ --topic vector_ingest < ../vector_sample_data.csv
<your_broker_list> : endpoint de l'instance. Vous pouvez obtenir l'endpoint depuis la section Access Point Information de la page Instance Details dans la console Message Queue for Apache Kafka.
Créer une table de mappage et importer les données
-
Créez une tâche Flink.
Connectez-vous à la console Realtime Compute for Apache Flink. Dans l'onglet Fully Managed Flink, localisez l'espace de travail cible et cliquez sur Console dans la colonne Actions.
Dans le volet de navigation de gauche, cliquez sur SQL Development. Cliquez sur New, sélectionnez Blank stream draft, puis cliquez sur Next.
-
Dans la boîte de dialogue New draft, configurez les paramètres du brouillon.
Paramètre
Description
Exemple
File Name
Le nom du brouillon.
RemarqueLe nom du brouillon doit être unique dans le projet actuel.
adbpg-test
Storage Location
Le dossier dans lequel le fichier de code du brouillon est stocké.
Vous pouvez également cliquer sur l'icône
à côté d'un dossier existant pour créer un sous-dossier.Drafts
Engine Version
La version du moteur Flink pour la tâche actuelle. Pour plus d'informations sur les versions du moteur, les correspondances de versions et les jalons du cycle de vie, consultez Versions du moteur.
vvr-6.0.6-flink-1.15
-
Exécutez l'instruction suivante pour créer une table de mappage pour 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' );Pour obtenir la description des paramètres, consultez la section Écrire des données dans AnalyticDB for PostgreSQL.
-
Exécutez l'instruction suivante pour créer une table de mappage pour 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' );Le tableau suivant décrit les paramètres.
Paramètre
Obligatoire
Description
connector
Oui
Le nom du connecteur. La valeur doit être
kafka.properties.bootstrap.servers
Oui
L'endpoint de l'instance Message Queue for Apache Kafka. Vous pouvez obtenir l'endpoint depuis la section Endpoint information de la page Instance details dans la console Message Queue for Apache Kafka.
topic
Oui
Le nom du topic Kafka.
format
Oui
Le format des valeurs des messages Kafka. Les formats suivants sont pris en charge :
-
csv
-
json
-
avro
-
debezium-json
-
canal-json
-
maxwell-json
-
avro-confluent
-
raw
csv.field-delimiter
Oui
Le délimiteur de champ pour le format CSV.
scan.startup.mode
Oui
Spécifie le décalage à partir duquel le consommateur Kafka commence à lire les données. Valeurs valides :
-
earliest-offset: Commence la lecture à partir du décalage disponible le plus ancien. -
latest-offset: Commence la lecture à partir du dernier décalage.
-
-
Exécutez l'instruction suivante pour créer une tâche d'importation :
INSERT INTO vector_ingest SELECT * FROM vector_kafka;