VECTOR_SEARCH is a function that finds the most semantically similar items based on a specified high-dimensional numerical vector, querying a Milvus vector table from a Flink SQL job and returning the top-K most similar entries.
Limitations
-
Version support: Ververica Runtime (VVR) 11.3 and later support stream mode. VVR 11.4 and later support batch mode.
-
The
OUTPUT_COLUMNSparameter is supported from VVR 11.9 of Realtime Compute for Apache Flink. -
Vector table: Milvus and DLF Paimon (VVR 11.8 and later) are supported as vector tables.
-
Stream type: Only insert-only streams are supported (streams containing only
INSERTmessages). -
Execution mode:
VECTOR_SEARCHruns in stream mode only; batch mode is not supported.
Syntax
VECTOR_SEARCH(
TABLE <SEARCH_TABLE>,
DESCRIPTOR(<COLUMN_TO_SEARCH>),
<COLUMN_TO_QUERY>,
<TOP_K>[,
<CONFIG>]
)
Parameters
|
Parameter |
Data type |
Description |
|
async |
Boolean |
Specifies whether to enable async mode. If the connector of the vector table does not support the specified mode, the engine returns an error. |
|
max-concurrent-operations |
Integer |
The maximum number of concurrent requests in asynchronous mode. |
|
output-mode |
Enum |
The vector feature column in the input, such as the embedding vector of an image or text uploaded by the user. The column is matched against the indexed vector column by similarity. |
|
TOP_K |
INT |
The maximum number of similar records to return. |
|
CONFIG |
MAP<STRING,STRING> |
The Runtime parameters that you can configure. |
|
OUTPUT_COLUMNS |
DESC |
Optional. Specifies the columns returned in the vector search result. You can select top-level columns of the vector table and the similarity column |
Return value
VECTOR_SEARCH returns a table. If OUTPUT_COLUMNS is not specified, each row contains all columns of the vector table and an additional score column, which is consistent with earlier versions. If OUTPUT_COLUMNS is specified, each row contains only the specified columns, in the order declared in DESCRIPTOR.
The score column is of the DOUBLE type and describes the similarity between the input data and each search result. You can also select or drop this column by using OUTPUT_COLUMNS. If a column named score already exists in the vector table, the engine selects a unique name for the generated similarity column, such as score0. In this case, use the actual column name in OUTPUT_COLUMNS.
OUTPUT_COLUMNS cannot be empty, does not support nested fields, and cannot contain duplicate columns.
Runtime parameters
|
Parameter |
Data type |
Default value |
Description |
|
async |
Boolean |
(none) |
Specifies whether to enable async mode. If the connector of the vector table does not support the specified mode, the engine returns an error. By default, the engine selects the execution mode based on the modes supported by the connector. If the connector supports both asynchronous and synchronous modes, the engine prefers the asynchronous mode to improve overall throughput. |
|
max-concurrent-operations |
Integer |
10 |
The maximum number of concurrent requests in asynchronous mode. |
|
output-mode |
Enum |
|
The output mode of asynchronous operations. Valid values:
For more information about the two values, see Async I/O. |
|
timeout |
Duration |
3min |
The timeout period from the first invocation until the asynchronous operation is completed. The period may cover multiple retries and is reset on failover. |
Example
Test data
The vector table vector_table contains the following data:
|
id |
topic |
vector_index |
|
1 |
'Spark' |
[5, 12, 13] |
|
2 |
'Flink' |
[-5, -12, -13] |
|
3 |
'Batch' |
[5, 12, 13] |
The query table query_table contains the following data:
|
id |
user_keyword |
embedding |
|
1 |
'Spark' |
[5, 12, 13] |
|
2 |
'Flink' |
[-5, -12, -13] |
Query
For each row of query_table, the following SQL statement searches vector_table and returns the two most similar records.
SELECT user_keyword, topic
FROM
query_table,
LATERAL TABLE (VECTOR_SEARCH(
SEARCH_TABLE => TABLE vector_table,
COLUMN_TO_SEARCH => DESCRIPTOR(vector_index),
COLUMN_TO_QUERY => query_table.embedding,
TOP_K => 2,
MAP['async', 'false'] -- Enable synchronous mode
))
Results
|
user_keyword |
topic |
score |
|
'Spark' |
'Batch' |
1.0 |
|
'Spark' |
'BigData' |
0.6538461538461539 |
|
'Flink' |
'Streaming' |
1.0 |
|
'Flink' |
'BigData' |
-0.6538461538461539 |