Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Milvus

Última atualização: Jun 27, 2026

Este tópico descreve como usar o conector do Milvus.

Visão geral

O Milvus é um banco de dados vetorial altamente escalável, projetado para processar dados não estruturados em grande escala, como imagens, texto e áudio. Ele oferece suporte a buscas de similaridade eficientes, sendo ideal para casos de uso como sistemas de recomendação, recuperação de imagens e busca semântica. O conector do Milvus oferece os seguintes recursos:

Categoria

Detalhes

Tipos suportados

Tabela sink, tabela vetorial

Modo de execução

Streaming

Formato de dados

Nenhum

Métricas de monitoramento específicas

Nenhuma

Tipo de API

SQL

Suporte a atualize/exclua

Sim

Recursos

O conector do Milvus integra o Apache Flink ao banco de dados vetorial Milvus para criar um pipeline de dados confiável e de alto desempenho, voltado para cenários de busca vetorial em tempo real. Os principais recursos incluem:

  • Gravações com alta concorrência: Permite configurar o paralelismo do sink.

  • Nova tentativa automática: Repete operações com falha para aumentar a estabilidade.

  • Buffer em lotes: Melhora o desempenho de gravação ao agrupar registros antes do flush.

  • Semântica at-least-once: Garante consistência eventual por meio de atualizações idempotentes baseadas na chave primária.

  • Busca vetorial: Possibilita pesquisas de similaridade vetorial em tempo real diretamente no Flink SQL.

Pré-requisitos

Limitações

  • A gravação em uma tabela sink exige o Ververica Runtime (VVR) 11.1 ou superior.

  • Consultas em tabelas vetoriais requerem o VVR 11.3 ou superior.

  • Somente o Milvus 2.4.x é suportado.

  • O conector do Milvus suporta apenas a semântica at-least-once.

Sintaxe

CREATE TEMPORARY TABLE milvus_sink (
  id BIGINT,
  f1 STRING,
  f2 BOOLEAN,
  f3 TINYINT,
  f4 SMALLINT,
  f5 INTEGER,
  f6 DATE,
  f7 TIME(3),
  f8 TIMESTAMP_LTZ(3),
  f9 TIMESTAMP(3),
  f10 FLOAT,
  f11 DOUBLE,
  f12 DECIMAL(10, 2),
  f13 ARRAY<FLOAT>,
  f14 ARRAY<DOUBLE>,
  f15 ARRAY<INTEGER>,
  f16 ARRAY<BIGINT>,
  PRIMARY KEY (id) NOT ENFORCED  -- Required. Milvus supports only BIGINT or STRING as the primary key.
) WITH (
  'connector'='milvus',
  'endpoint'='<yourEndpoint>',
  'port'='<yourPort>',
  'userName'='<yourUserName>',
  'password'='<yourPassword>',
  'databaseName'='<yourDatabaseName>',
  'collectionName'='<yourCollectionName>'
);

Opções do conector

Gerais

Opção

Descrição

Tipo de dados

Obrigatório

Valor padrão

Observações

connector

Nome do conector.

String

Sim

Defina como milvus.

endpoint

Endpoint (endereço IP ou nome de domínio) do banco de dados Milvus.

String

Sim

Consulte Configurações de acesso à rede e segurança.

port

Número da porta do banco de dados Milvus.

INTEGER

Não

19530

username

Nome de usuário para acesso ao banco de dados Milvus.

STRING

Sim

Nenhuma.

password

Senha para acesso ao banco de dados Milvus.

STRING

Sim

databaseName

Nome do banco de dados Milvus.

STRING

Sim

collectionName

Nome da coleção do Milvus.

STRING

Sim

partitionName

Nome da partição de destino para gravação.

STRING

Não

_default

partitionKey.enabled

Indica se a coleção utiliza um campo escalar como chave de partição.

BOOLEAN

Não

false

maxRetries

Quantidade de novas tentativas para operações com falha.

INTEGER

Não

3

Nenhuma.

Específicas do sink

Opção

Descrição

Tipo de dados

Obrigatório

Valor padrão

Observações

sink.parallelism

Paralelismo do operador sink.

INTEGER

Não

Se não for definido, o paralelismo será herdado do operador anterior.

sink.maxRetries

Número máximo de tentativas após falha na gravação.

INTEGER

Não

3

No VVR 11.3 e versões posteriores, esta opção foi descontinuada. Utilize maxRetries em seu lugar.

sink.buffer-flush.max-rows

Limite máximo de registros no buffer (incluindo inserções, upserts e exclusões). O flush é acionado ao atingir essa quantidade.

INTEGER

Não

10000

Configure como 0 para desativar esse gatilho.

sink.buffer-flush.interval

Intervalo de tempo em milissegundos (ms) para realizar o flush dos registros em buffer. O flush ocorre quando esse intervalo é excedido.

INTEGER

Não

1000

Configure como 0 para desativar o flush baseado em tempo.

sink.ignoreDelete

Define se as operações de exclusão devem ser ignoradas.

BOOLEAN

Não

false

Valores válidos:

  • true: Ignora operações de exclusão.

  • false: Processa operações de exclusão.

Específicas da tabela vetorial

Parâmetro

Descrição

Tipo de dados

Obrigatório

Valor padrão

Observações

search.metric

Métrica utilizada para medir a similaridade entre vetores.

String

Não

L2

Para detalhes sobre as métricas de similaridade suportadas, consulte a documentação do Milvus. A versão 2.4 do Milvus suporta atualmente as seguintes métricas:

  • L2: Distância euclidiana

  • IP: Produto interno

  • COSINE: Similaridade de cosseno

Mapeamento de tipos

Tipo Flink SQL

Tipo Milvus

STRING

VarChar(n)

BOOLEAN

Bool

TINYINT

Int8

SMALLINT

Int16

INTEGER

Int32

BIGINT

Int64

DATE

VarChar(n)

TIME(3)

VarChar(n)

TIMESTAMP_LTZ(3)

Int64

Nota

Armazenado como epoch time em milissegundos.

TIMESTAMP(3)

VarChar(n)

FLOAT

Float

DOUBLE

Double

DECIMAL(10, 2)

VarChar(n)

ARRAY<FLOAT>

FloatVector

Nota

Após criar uma coleção no Milvus, crie um índice para o campo vetorial.

ARRAY<DOUBLE>

Array<Double>[m]

ARRAY<INTEGER>

Array<Int32>[m]

ARRAY<BIGINT>

Array<Int64>[m]

Exemplos

  • Gravar dados de streaming no Milvus

-- Create a mock data source that generates 100 rows per second. simulated streaming data.
CREATE TEMPORARY TABLE mock_source (
    id STRING,
    vector ARRAY<FLOAT>,       -- Vector passed as a FLOAT array.
    event_time AS PROCTIME() -- Processing time attribute.
) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '100',  -- Generates 100 rows per second.
    'fields.id.kind' = 'sequence',
    'fields.id.start' = '1',
    'fields.id.end' = '1000'
);

CREATE TEMPORARY TABLE milvus_sink (
  id STRING,               -- Unique identifier, such as a device ID.
  vector ARRAY<FLOAT>,     -- Vector data. The array length must be consistent with the source.
  timestamp BIGINT         -- Timestamp for stream processing.
  PRIMARY KEY (id) NOT ENFORCED  -- Required. Milvus supports only BIGINT or STRING as the primary key.
) WITH (
  'connector'='milvus',
  'endpoint'='xxx',
  'port'='19530',
  'userName'='xxx',
  'password'='xxx',
  'databaseName'='xxxx',
  'collectionName'='xxxx'
);

-- Transform data and write to Milvus.
INSERT INTO milvus_sink
SELECT 
    id,
    vector,
    UNIX_TIMESTAMP() * 1000 AS timestamp  -- Current timestamp in milliseconds.
FROM mock_source;
  • Executar uma busca vetorial

CREATE TEMPORARY TABLE milvus_table (
  id STRING,               -- Unique identifier.
  vector ARRAY<FLOAT>,     -- Vector data. The array length must be consistent with the source.
  PRIMARY KEY (id) NOT ENFORCED  -- Required. Milvus only supports BIGINT or STRING as the primary key.
) WITH (
  'connector'='milvus',
  'endpoint'='xxx',
  'port'='19530',
  'userName'='xxx',
  'password'='xxx',
  'databaseName'='xxxx',
  'collectionName'='xxxx'
);
-- Find the top 2 most similar items for the vector [1.1, 2.2, 3.3].
SELECT * FROM 
LATERAL TABLE(
  VECTOR_SEARCH(
    TABLE milvus_table, 
    DESCRIPTOR(vector), 
    ARRAY[1.1, 2.2, 3.3], 
    2));

Nota: A coleção do Milvus deve estar carregada na memória antes da execução da consulta. Para mais informações, consulte a documentação do Milvus sobre Carregar uma coleção.