Todos os produtos
Search
Central de documentação

AnalyticDB:Integrar dados vetoriais com o Realtime Compute for Flink

Última atualização: Sep 02, 2026

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.

  • A extensão de recuperação vetorial FastANN deve estar instalada no seu banco de dados AnalyticDB for PostgreSQL.

    Execute o comando \dx fastann em 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

  1. Crie índices estruturados e vetoriais.

  2. Grave os dados vetoriais de exemplo em um tópico do Kafka.

  3. Crie uma tabela de mapeamento e importe dados.

Criar índices estruturados e vetoriais

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

  2. Execute as instruções abaixo para criar e acessar um banco de dados de teste:

    CREATE DATABASE adbpg_test;
    \c adbpg_test
  3. 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);
  4. 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

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

  1. Crie um job do Flink.

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

    2. No painel de navegação à esquerda, clique em SQL Development. Clique em New, selecione Blank stream draft e, em seguida, clique em Next.

    3. Na caixa de diálogo New draft, defina os parâmetros do rascunho.

      Parâmetro

      Descrição

      Exemplo

      File Name

      Nome do rascunho.

      Nota

      O 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

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

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

  4. Execute a instrução a seguir para criar um job de importação:

    INSERT INTO vector_ingest SELECT * FROM vector_kafka;