Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions docs/reference/online-stores/milvus.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,8 @@ online_store:
| `nlist` | `128` | `nlist` index parameter, used when `index_params` is unset. |
| `index_params` | unset | Index build parameters passed to Milvus, e.g. `{M: 16, efConstruction: 200}` for HNSW. Defaults to `{nlist: <nlist>}`, or no parameters for `AUTOINDEX`. |
| `search_params` | unset | Search parameters passed to Milvus, e.g. `{ef: 64}` for HNSW or `{level: 2}` for `AUTOINDEX`. Defaults to `{nprobe: 10}`, or no parameters for `AUTOINDEX`. |
| `consistency_level` | unset | `Strong`, `Bounded`, `Session` or `Eventually`. Sent with every read and search. When unset, Milvus uses the collection's level. |
| `collection_consistency_level` | unset | `Strong`, `Bounded`, `Session` or `Eventually`. Set when Feast creates a collection. When unset, Milvus uses its default (`Bounded`). |
| `vector_enabled` | `true` | Enables vector search. |
| `varchar_max_length` | `65535` | Default `max_length` of VARCHAR fields. Override per field with the `max_length` tag. |
| `enable_openai_compatible_store` | `false` | Store numeric features as native Milvus numeric types. |
Expand Down Expand Up @@ -130,6 +132,18 @@ online_store:
Index parameters only apply when Feast creates a collection. To change them for an existing
collection, run `feast teardown` and `feast apply`, then materialize again.

## Consistency level

By default Milvus uses `Bounded` consistency, so a read issued straight after materialization may
not see the newest writes for a short time. Set `consistency_level: Strong` if reads must always see
the latest writes, at the cost of higher read latency. See the
[Milvus consistency documentation](https://milvus.io/docs/consistency.md).

`consistency_level` applies to Feast's reads and searches and takes effect immediately.
`collection_consistency_level` sets the collection's own default, which Milvus uses for requests
that don't specify a level, such as those from other clients. It only applies when Feast creates a
collection; to change it, run `feast teardown` and `feast apply`, then materialize again.

## Collection loading

Feast creates collections together with their indexes, which makes Milvus load them straight away.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,16 @@ class MilvusOnlineStoreConfig(FeastConfigBaseModel, VectorStoreConfig):
# Search params, e.g. {"ef": 64}, or {"level": 2} for AUTOINDEX.
# Defaults to {"nprobe": 10}, or {} for AUTOINDEX.
search_params: Optional[Dict[str, Any]] = None
# Sent with every read and search. When unset, Milvus uses the
# collection's level.
consistency_level: Optional[
Literal["Strong", "Bounded", "Session", "Eventually"]
] = None
# Set when Feast creates a collection. When unset, Milvus uses its
# default (Bounded).
collection_consistency_level: Optional[
Literal["Strong", "Bounded", "Session", "Eventually"]
] = None
username: Optional[StrictStr] = ""
password: Optional[StrictStr] = ""
enable_openai_compatible_store: Optional[bool] = False
Expand Down Expand Up @@ -421,6 +431,7 @@ def _get_or_create_collection(
dimension=config.online_store.embedding_dim,
schema=schema,
index_params=index_params,
**_collection_consistency_kwargs(config.online_store),
)
else:
self._ensure_loaded(collection_name)
Expand Down Expand Up @@ -584,6 +595,7 @@ def online_read(
collection_name=collection_name,
filter=query_filter_for_entities,
output_fields=output_fields,
**_consistency_kwargs(config.online_store),
)
# Group hits by composite key.
grouped_hits: Dict[str, Any] = {}
Expand Down Expand Up @@ -863,6 +875,7 @@ def _combine_exprs(*parts: Optional[str]) -> Optional[str]:
limit=top_k,
output_fields=output_fields,
filter=combined_filter,
**_consistency_kwargs(config.online_store),
)

elif embedding is not None and config.online_store.vector_enabled:
Expand All @@ -880,6 +893,7 @@ def _combine_exprs(*parts: Optional[str]) -> Optional[str]:
limit=top_k,
output_fields=output_fields,
filter=metadata_filter_expr,
**_consistency_kwargs(config.online_store),
)

elif query_string is not None:
Expand Down Expand Up @@ -915,6 +929,7 @@ def _combine_exprs(*parts: Optional[str]) -> Optional[str]:
filter=combined_filter or text_filter,
output_fields=output_fields,
limit=top_k,
**_consistency_kwargs(config.online_store),
)

results = [
Expand Down Expand Up @@ -1053,6 +1068,22 @@ def _search_params(online_config: MilvusOnlineStoreConfig) -> Dict[str, Any]:
return {"nprobe": 10}


def _consistency_kwargs(online_config: MilvusOnlineStoreConfig) -> Dict[str, Any]:
"""Read and search kwargs; only pass a level when configured."""
if online_config.consistency_level:
return {"consistency_level": online_config.consistency_level}
return {}


def _collection_consistency_kwargs(
online_config: MilvusOnlineStoreConfig,
) -> Dict[str, Any]:
"""create_collection kwargs; only pass a level when configured."""
if online_config.collection_consistency_level:
return {"consistency_level": online_config.collection_consistency_level}
return {}


def _milvus_token(online_config: MilvusOnlineStoreConfig) -> str:
"""Return the token to authenticate with: ``token``, else ``username:password``."""
if online_config.token:
Expand Down
48 changes: 48 additions & 0 deletions sdk/python/tests/integration/online_store/test_milvus_remote.py
Original file line number Diff line number Diff line change
Expand Up @@ -312,3 +312,51 @@ def test_autoindex_with_search_level(
lambda hits: len(hits) == 1,
)
assert hits[0]["city"].string_val == "Paris"


@pytest.mark.parametrize(
"collection_consistency_level, consistency_level",
[(None, None), ("Strong", None), (None, "Strong")],
)
def test_consistency_level(
tmp_path: Path,
project: str,
store: MilvusOnlineStore,
collection_consistency_level: Optional[str],
consistency_level: Optional[str],
) -> None:
online_store = {
key: value
for key, value in {
"collection_consistency_level": collection_consistency_level,
"consistency_level": consistency_level,
}.items()
if value
}
config = _repo_config(tmp_path, project, **online_store)
fv = _scalar_feature_view()
store.update(config, [], [fv], [], [], partial=False)

assert store.client is not None
description = store.client.describe_collection(f"{project}_{fv.name}")
# Only collection_consistency_level sets the collection's level; Milvus
# defaults to Bounded.
assert description["consistency_level_name"] == (
collection_consistency_level or "Bounded"
)

if consistency_level == "Strong":
# A Strong read sees the previous write, even on a Bounded collection.
_write_rows(
store,
config,
fv,
{
1: {
"trips_today": ValueProto(float_val=1.0),
"city": ValueProto(string_val="Oslo"),
}
},
)
rows = _read(store, config, fv, [1], ["city"])
assert rows[0] is not None and rows[0]["city"].string_val == "Oslo"
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@

from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional
from typing import Any, Dict, List, Optional, Tuple
from unittest.mock import MagicMock, patch

import pytest
from pydantic import ValidationError
from pymilvus import DataType, MilvusClient
from pymilvus.client.types import LoadState

Expand Down Expand Up @@ -496,3 +497,85 @@ def test_autoindex_search_with_level(tmp_path: Path) -> None:

assert len(results) == 1
assert results[0][2] is not None and results[0][2]["city"].string_val == "Paris"


def _consistency_calls(
**online_store: Any,
) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]:
"""Return the kwargs of create_collection calls and of query/search calls."""
with patch(f"{MILVUS_MODULE}.MilvusClient") as mock_client_cls:
mock_client = _mock_client(mock_client_cls, has_collection=False)
mock_client.describe_collection.return_value = {
"collection_name": "test_milvus_driver_embeddings",
"fields": [
{"name": "driver_id_pk", "type": DataType.VARCHAR, "params": {}},
{"name": "event_ts", "type": DataType.INT64, "params": {}},
{"name": "created_ts", "type": DataType.INT64, "params": {}},
{
"name": "embedding",
"type": DataType.FLOAT_VECTOR,
"params": {"dim": 2},
},
{"name": "city", "type": DataType.VARCHAR, "params": {}},
],
}
mock_client.search.return_value = [[]]
mock_client.query.return_value = []

store = MilvusOnlineStore()
config = _mock_config(embedding_dim=2, **online_store)
fv = _vector_feature_view()
store.online_read(config, fv, [_entity_key(1)], ["city"])
store.retrieve_online_documents_v2(
config, fv, ["embedding", "city"], embedding=[1.0, 0.0], top_k=1
)
store.retrieve_online_documents_v2(
config, fv, ["city"], embedding=None, top_k=1, query_string="Paris"
)
create_calls = [
call.kwargs for call in mock_client.create_collection.call_args_list
]
read_calls = [
call.kwargs
for call in mock_client.query.call_args_list
+ mock_client.search.call_args_list
]
return create_calls, read_calls


def test_consistency_levels_not_sent_when_unset() -> None:
create_calls, read_calls = _consistency_calls()

assert len(create_calls) == 1 and len(read_calls) == 3
assert all("consistency_level" not in kwargs for kwargs in create_calls)
assert all("consistency_level" not in kwargs for kwargs in read_calls)


def test_consistency_level_applied_to_reads_and_searches_only() -> None:
create_calls, read_calls = _consistency_calls(consistency_level="Strong")

assert "consistency_level" not in create_calls[0]
assert len(read_calls) == 3
assert all(kwargs["consistency_level"] == "Strong" for kwargs in read_calls)


def test_collection_consistency_level_applied_on_create_only() -> None:
create_calls, read_calls = _consistency_calls(collection_consistency_level="Strong")

assert create_calls[0]["consistency_level"] == "Strong"
assert all("consistency_level" not in kwargs for kwargs in read_calls)


def test_collection_and_read_consistency_levels_independent() -> None:
create_calls, read_calls = _consistency_calls(
collection_consistency_level="Session", consistency_level="Strong"
)

assert create_calls[0]["consistency_level"] == "Session"
assert all(kwargs["consistency_level"] == "Strong" for kwargs in read_calls)


@pytest.mark.parametrize("field", ["consistency_level", "collection_consistency_level"])
def test_invalid_consistency_level_rejected(field: str) -> None:
with pytest.raises(ValidationError):
MilvusOnlineStoreConfig(**{field: "Immediate"})
Loading