Tous les produits
Search
Centre de documentation

AnalyticDB:Intégrer des données vectorielles avec Realtime Compute for Flink

Dernière mise à jour :Sep 02, 2026

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.

  • 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 fastann sur 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

  1. Créer des index structurés et vectoriels.

  2. Écrire les exemples de données vectorielles dans un topic Kafka.

  3. Créer une table de mappage et importer les données.

Créer des index structurés et vectoriels

  1. 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.

  2. 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
  3. 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);
  4. 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

  1. 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>
  2. 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

  1. Créez une tâche Flink.

    1. 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.

    2. Dans le volet de navigation de gauche, cliquez sur SQL Development. Cliquez sur New, sélectionnez Blank stream draft, puis cliquez sur Next.

    3. Dans la boîte de dialogue New draft, configurez les paramètres du brouillon.

      Paramètre

      Description

      Exemple

      File Name

      Le nom du brouillon.

      Remarque

      Le 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

  2. 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.

  3. 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.

  4. Exécutez l'instruction suivante pour créer une tâche d'importation :

    INSERT INTO vector_ingest SELECT * FROM vector_kafka;