FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

fix: Make Milvus placeholder vectors valid on Milvus servers by simonhearne · Pull Request #6881 · feast-dev/feast · GitHub

Repository navigation

Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension .baseline  (1) .md  (1) .py  (3) All 3 file types selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
4 changes: 2 additions & 2 deletions .secrets.baseline

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

45 changes: 41 additions & 4 deletions docs/reference/online-stores/milvus.md
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# Redis online store
# Milvus online store

## Description

Expand All @@ -19,11 +19,11 @@ Feast supports both milvus-lite 2.x and 3.x. However, if you upgrade from milvus
See the [milvus-lite GitHub page](https://github.com/milvus-io/milvus-lite) for more details.
{% endhint %}

You can get started by using any of the other templates (e.g. `feast init -t gcp` or `feast init -t snowflake` or `feast init -t aws`), and then swapping in Redis as the online store as seen below in the examples.
You can get started by using any of the other templates (e.g. `feast init -t gcp` or `feast init -t snowflake` or `feast init -t aws`), and then swapping in Milvus as the online store as seen below in the examples.

## Examples

Connecting to a local MilvusDB instance:
Using Milvus Lite, which stores data in a local file:

{% code title="feature_store.yaml" %}
```yaml
Expand All @@ -33,18 +33,55 @@ provider: local
online_store:
type: milvus
path: "data/online_store.db"
connection_string: "localhost:6379"
embedding_dim: 128
index_type: "FLAT"
metric_type: "COSINE"
```
{% endcode %}

Connecting to a self-hosted Milvus server:

{% code title="feature_store.yaml" %}
```yaml
project: my_feature_repo
registry: data/registry.db
provider: local
online_store:
type: milvus
host: "http://localhost"
port: 19530
username: "username"
password: "password"
embedding_dim: 128
index_type: "IVF_FLAT"
metric_type: "COSINE"
```
{% endcode %}

## Configuration options

| Option | Default | Description |
|:-------|:--------|:------------|
| `path` | `""` | Path to a Milvus Lite database file. Used when `provider: local` and `path` is set. |
| `host` | `http://localhost` | Milvus server host, including the scheme. |
| `port` | `19530` | Milvus server port. |
| `username` / `password` | `""` | Credentials, sent as the token `username:password`. |
| `embedding_dim` | `128` | Dimension of vector fields. |
| `index_type` | `FLAT` | Index type for vector fields with `vector_index=True`. |
| `metric_type` | `COSINE` | Default metric when a field does not set `vector_search_metric`. |
| `nlist` | `128` | `nlist` index parameter. |
| `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. |

The full set of configuration options is available in [MilvusOnlineStoreConfig](https://rtd.feast.dev/en/latest/#feast.infra.online_stores.milvus.MilvusOnlineStoreConfig).

## Feature views without vectors

Milvus requires every collection to have a vector field. For feature views that have no vector
feature, Feast adds a 2-dimensional `_placeholder_vector` field with a FLAT index and fills it with zeros.
It is never returned or searched.

## Functionality Matrix

The set of functionality supported by online stores is described in detail [here](overview.md#functionality).
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,13 @@
DataType.BOOL,
}

# Milvus requires every collection to have a vector field, so feature views
# without one get a small placeholder vector. Milvus servers reject vectors
# with fewer than 2 dimensions and non-finite values, and a collection can only
# be loaded once every vector field is indexed.
PLACEHOLDER_VECTOR_FIELD = "_placeholder_vector"
PLACEHOLDER_VECTOR_DIM = 2


def _milvus_escape_string(s: str) -> str:
"""Escape a string for safe use inside a Milvus single-quoted literal.
Expand Down Expand Up @@ -348,9 +355,9 @@ def _get_or_create_collection(
if not has_vector_field:
fields.append(
FieldSchema(
name="_placeholder_vector",
name=PLACEHOLDER_VECTOR_FIELD,
dtype=DataType.FLOAT_VECTOR,
dim=1,
dim=PLACEHOLDER_VECTOR_DIM,
)
)
schema = CollectionSchema(
Expand All @@ -368,14 +375,12 @@ def _get_or_create_collection(
index_params = self.client.prepare_index_params()
indices_added = False
for vector_field in schema.fields:
if (
vector_field.dtype
in [
DataType.FLOAT_VECTOR,
DataType.BINARY_VECTOR,
]
and vector_field.name in vector_field_dict
):
if vector_field.dtype not in [
DataType.FLOAT_VECTOR,
DataType.BINARY_VECTOR,
]:
continue
if vector_field.name in vector_field_dict:
metric = vector_field_dict[
vector_field.name
].vector_search_metric
Expand All @@ -387,7 +392,20 @@ def _get_or_create_collection(
index_name=f"vector_index_{vector_field.name}",
params={"nlist": config.online_store.nlist},
)
indices_added = True
else:
# Vector fields that aren't searched (the placeholder,
# or arrays without vector_index) still need an index,
# otherwise Milvus servers refuse to load the collection.
index_params.add_index(
collection_name=collection_name,
field_name=vector_field.name,
metric_type="L2"
if vector_field.name == PLACEHOLDER_VECTOR_FIELD
else config.online_store.metric_type,
index_type="FLAT",
index_name=f"vector_index_{vector_field.name}",
)
indices_added = True
if indices_added:
self.client.create_index(
collection_name=collection_name,
Expand Down Expand Up @@ -423,6 +441,15 @@ def online_write_batch(
collection_field_types = {
field["name"]: field["type"] for field in collection["fields"]
}
# Collections created by older Feast versions may use a 1-dim placeholder.
placeholder_dim = next(
(
int(field.get("params", {}).get("dim", PLACEHOLDER_VECTOR_DIM))
for field in collection["fields"]
if field["name"] == PLACEHOLDER_VECTOR_FIELD
),
PLACEHOLDER_VECTOR_DIM,
)
schema_internal_fields = {"event_ts", "created_ts"}
collection_has_native_numerics = any(
collection_field_types.get(name) in MILVUS_NATIVE_NUMERIC_TYPES
Expand Down Expand Up @@ -472,8 +499,8 @@ def online_write_batch(
for field in required_fields:
if field not in single_entity_record:
field_type = collection_field_types.get(field, DataType.VARCHAR)
if field == "_placeholder_vector":
single_entity_record[field] = [float("nan")]
if field == PLACEHOLDER_VECTOR_FIELD:
single_entity_record[field] = [0.0] * placeholder_dim
else:
single_entity_record[field] = _default_for_milvus_type(
field_type
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
"""Integration tests for the Milvus online store against a Milvus server or Zilliz Cloud.

These tests only run when ``ZILLIZ_URI`` and ``ZILLIZ_TOKEN`` are set, e.g.::

ZILLIZ_URI=https://<cluster>.zillizcloud.com ZILLIZ_TOKEN=<user:password> \
pytest --integration sdk/python/tests/integration/online_store/test_milvus_remote.py

They also run against a self-hosted Milvus server, e.g.
``ZILLIZ_URI=http://localhost:19530 ZILLIZ_TOKEN=root:Milvus``.
"""

import os
import time
import uuid
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Callable, Dict, Iterator, List, Optional, TypeVar
from urllib.parse import urlparse

import pytest

from feast import Entity, FeatureView
from feast.field import Field
from feast.infra.online_stores.milvus_online_store.milvus import MilvusOnlineStore
from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto
from feast.protos.feast.types.Value_pb2 import Value as ValueProto
from feast.repo_config import RepoConfig
from feast.types import Float32, Int64, String
from feast.value_type import ValueType

T = TypeVar("T")

ZILLIZ_URI = os.environ.get("ZILLIZ_URI")
ZILLIZ_TOKEN = os.environ.get("ZILLIZ_TOKEN")

pytestmark = [
pytest.mark.integration,
pytest.mark.skipif(
not (ZILLIZ_URI and ZILLIZ_TOKEN),
reason="ZILLIZ_URI and ZILLIZ_TOKEN must be set to run Milvus server tests",
),
]


def _connection_config() -> Dict[str, Any]:
assert ZILLIZ_URI and ZILLIZ_TOKEN
parsed = urlparse(ZILLIZ_URI)
default_port = 443 if parsed.scheme == "https" else 19530
username, _, password = ZILLIZ_TOKEN.partition(":")
return {
"host": f"{parsed.scheme}://{parsed.hostname}",
"port": parsed.port or default_port,
"username": username,
"password": password,
}


def _repo_config(tmp_path: Path, project: str, **online_store: Any) -> RepoConfig:
return RepoConfig(
project=project,
provider="local",
registry=str(tmp_path / "registry.db"),
online_store={
"type": "milvus",
"embedding_dim": 2,
**_connection_config(),
**online_store,
},
entity_key_serialization_version=3,
repo_path=tmp_path,
)


@pytest.fixture
def project() -> str:
return f"feast_it_{uuid.uuid4().hex[:8]}"


@pytest.fixture
def store() -> Iterator[MilvusOnlineStore]:
store = MilvusOnlineStore()
yield store
if store.client is not None:
for name in list(store._collections):
store.client.drop_collection(name)


def _eventually(fn: Callable[[], T], check: Callable[[T], bool]) -> T:
"""Retry ``fn`` until ``check`` passes, tolerating Bounded consistency."""
deadline = time.monotonic() + 15
while True:
result = fn()
if check(result) or time.monotonic() > deadline:
return result
time.sleep(0.5)


def _entity_key(driver_id: int) -> EntityKeyProto:
return EntityKeyProto(
join_keys=["driver_id"], entity_values=[ValueProto(int64_val=driver_id)]
)


def _write_rows(
store: MilvusOnlineStore,
config: RepoConfig,
fv: FeatureView,
rows: Dict[int, Dict[str, ValueProto]],
) -> None:
now = datetime.now(timezone.utc)
store.online_write_batch(
config,
fv,
[(_entity_key(k), dict(v), now, now) for k, v in rows.items()],
progress=None,
)


def _read(
store: MilvusOnlineStore,
config: RepoConfig,
fv: FeatureView,
driver_ids: List[int],
features: List[str],
) -> List[Optional[Dict[str, ValueProto]]]:
results = store.online_read(
config, fv, [_entity_key(d) for d in driver_ids], features
)
return [values for _, values in results]


def _scalar_feature_view() -> FeatureView:
return FeatureView(
name="driver_stats",
entities=[
Entity(
name="driver_id", join_keys=["driver_id"], value_type=ValueType.INT64
)
],
ttl=timedelta(days=1),
schema=[
Field(name="driver_id", dtype=Int64),
Field(name="trips_today", dtype=Float32),
Field(name="city", dtype=String),
],
)


def test_scalar_feature_view_round_trip(
tmp_path: Path, project: str, store: MilvusOnlineStore
) -> None:
"""Feature views without vectors rely on the placeholder vector."""
config = _repo_config(tmp_path, project)
fv = _scalar_feature_view()
store.update(config, [], [fv], [], [], partial=False)

_write_rows(
store,
config,
fv,
{
1: {
"trips_today": ValueProto(float_val=3.0),
"city": ValueProto(string_val="Paris"),
}
},
)
rows = _eventually(
lambda: _read(store, config, fv, [1], ["trips_today", "city"]),
lambda rows: rows[0] is not None,
)

assert rows[0] is not None and rows[0]["city"].string_val == "Paris"
Loading
Loading

Back | FazBrowse Home | New Git URL