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
Crie um cluster do Milvus. Para mais informações, consulte Início rápido: Criar uma instância do Milvus.
Crie uma coleção no Milvus. Se você planeja gravar em uma partição específica, verifique se ela já existe.
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 |
|
|
endpoint |
Endpoint (endereço IP ou nome de domínio) do banco de dados Milvus. |
String |
Sim |
||
|
port |
Número da porta do banco de dados Milvus. |
INTEGER |
Não |
|
|
|
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 |
|
|
|
partitionKey.enabled |
Indica se a coleção utiliza um campo escalar como chave de partição. |
BOOLEAN |
Não |
|
|
|
maxRetries |
Quantidade de novas tentativas para operações com falha. |
INTEGER |
Não |
|
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 |
|
No VVR 11.3 e versões posteriores, esta opção foi descontinuada. Utilize |
|
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 |
|
Configure como |
|
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 |
|
Configure como |
|
sink.ignoreDelete |
Define se as operações de exclusão devem ser ignoradas. |
BOOLEAN |
Não |
|
Valores válidos:
|
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 |
|
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:
|
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.