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

feat: Support Milvus partition keys in the Milvus online store · feast-dev/feast@345d382 · GitHub

Repository navigation

Commit 345d382

Browse files
authored andcommitted
feat: Support Milvus partition keys in the Milvus online store
Multi-tenant feature views, such as a catalogue shared by many brands, benefit from Milvus partition keys: searches filtered on the key only scan the matching partitions. A field can be marked as the partition key with the feature view tag milvus.partition_key, or for every feature view containing the field with the new partition_key store config. The tag takes precedence. The field must be stored as VARCHAR or INT64. Partition keys only apply when a collection is created. Feast logs a warning when an existing collection lacks the configured key. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Signed-off-by: Simon Hearne <simon.hearne@gmail.com>
1 parent 6d15533 commit 345d382

4 files changed

Lines changed: 327 additions & 0 deletions

File tree

‎docs/reference/online-stores/milvus.md‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,7 @@ online_store:
9696
| `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`. |
9797
| `consistency_level` | unset | `Strong`, `Bounded`, `Session` or `Eventually`. Sent with every read and search. When unset, Milvus uses the collection's level. |
9898
| `collection_consistency_level` | unset | `Strong`, `Bounded`, `Session` or `Eventually`. Set when Feast creates a collection. When unset, Milvus uses its default (`Bounded`). |
99+
| `partition_key` | unset | Field to use as the Milvus partition key in feature views that contain it. See [Partition key](#partition-key). |
99100
| `vector_enabled` | `true` | Enables vector search. |
100101
| `varchar_max_length` | `65535` | Default `max_length` of VARCHAR fields. Override per field with the `max_length` tag. |
101102
| `enable_openai_compatible_store` | `false` | Store numeric features as native Milvus numeric types. |
@@ -132,6 +133,51 @@ online_store:
132133
Index parameters only apply when Feast creates a collection. To change them for an existing
133134
collection, run `feast teardown` and `feast apply`, then materialize again.
134135

136+
## Partition key
137+
138+
A [partition key](https://milvus.io/docs/use-partition-key.md) makes Milvus group rows by the key's
139+
value, so searches filtered on it only scan the matching partitions. This suits multi-tenant data,
140+
such as a catalogue shared by many brands.
141+
142+
Set the partition key per feature view with the `milvus.partition_key` tag:
143+
144+
```python
145+
products = FeatureView(
146+
name="products",
147+
entities=[product],
148+
schema=[
149+
Field(name="product_id", dtype=Int64),
150+
Field(name="brand_id", dtype=String),
151+
Field(name="embedding", dtype=Array(Float32), vector_index=True),
152+
Field(name="title", dtype=String),
153+
],
154+
source=products_source,
155+
tags={"milvus.partition_key": "brand_id"},
156+
)
157+
```
158+
159+
or for every feature view that has the field, with `partition_key: brand_id` in the online store
160+
config. The tag takes precedence. Then filter on the key when retrieving:
161+
162+
```python
163+
store.retrieve_online_documents_v2(
164+
features=["products:embedding", "products:title"],
165+
query=query_embedding,
166+
top_k=10,
167+
filters=ComparisonFilter(type="eq", key="brand_id", value="acme"),
168+
)
169+
```
170+
171+
The partition key field must be stored as `VARCHAR` or `INT64`. String fields are always `VARCHAR`,
172+
and integer fields are `VARCHAR` unless native numeric types are enabled, in which case `Int64` works
173+
and `Int32` does not.
174+
175+
{% hint style="warning" %}
176+
The partition key only applies when Feast creates a collection. Existing collections are not
177+
changed; Feast logs a warning for them. To add a partition key to an existing feature view, run
178+
`feast teardown` and `feast apply`, then materialize again.
179+
{% endhint %}
180+
135181
## Consistency level
136182

137183
By default Milvus uses `Bounded` consistency, so a read issued straight after materialization may

‎sdk/python/feast/infra/online_stores/milvus_online_store/milvus.py‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -240,6 +240,9 @@ class MilvusOnlineStoreConfig(FeastConfigBaseModel, VectorStoreConfig):
240240
collection_consistency_level: Optional[
241241
Literal["Strong", "Bounded", "Session", "Eventually"]
242242
] = None
243+
# Field to use as the Milvus partition key in feature views that contain it.
244+
# A feature view's "milvus.partition_key" tag takes precedence.
245+
partition_key: Optional[StrictStr] = None
243246
username: Optional[StrictStr] = ""
244247
password: Optional[StrictStr] = ""
245248
enable_openai_compatible_store: Optional[bool] = False
@@ -386,6 +389,9 @@ def _get_or_create_collection(
386389
dim=PLACEHOLDER_VECTOR_DIM,
387390
)
388391
)
392+
partition_key = _partition_key_name(config.online_store, table)
393+
if partition_key:
394+
_mark_partition_key(fields, partition_key, table.name)
389395
schema = CollectionSchema(
390396
fields=fields, description="Feast feature view data"
391397
)
@@ -440,6 +446,10 @@ def _get_or_create_collection(
440446
self._collections[collection_name] = self.client.describe_collection(
441447
collection_name
442448
)
449+
if collection_exists and partition_key:
450+
_warn_if_partition_key_missing(
451+
self._collections[collection_name], partition_key
452+
)
443453
return self._collections[collection_name]
444454

445455
def _ensure_loaded(self, collection_name: str) -> None:
@@ -1068,6 +1078,63 @@ def _search_params(online_config: MilvusOnlineStoreConfig) -> Dict[str, Any]:
10681078
return {"nprobe": 10}
10691079

10701080

1081+
PARTITION_KEY_TAG = "milvus.partition_key"
1082+
PARTITION_KEY_TYPES = {DataType.INT64, DataType.VARCHAR}
1083+
1084+
1085+
def _partition_key_name(
1086+
online_config: MilvusOnlineStoreConfig, table: FeatureView
1087+
) -> Optional[str]:
1088+
"""Return the partition key field for a feature view, if one is configured.
1089+
1090+
The feature view's ``milvus.partition_key`` tag takes precedence over the
1091+
store-level ``partition_key``, which only applies to feature views that
1092+
contain that field.
1093+
"""
1094+
tagged = table.tags.get(PARTITION_KEY_TAG)
1095+
if tagged:
1096+
return tagged
1097+
configured = online_config.partition_key
1098+
if configured and any(field.name == configured for field in table.schema):
1099+
return configured
1100+
return None
1101+
1102+
1103+
def _mark_partition_key(
1104+
fields: List[FieldSchema], partition_key: str, table_name: str
1105+
) -> None:
1106+
field = next((f for f in fields if f.name == partition_key), None)
1107+
if field is None:
1108+
raise ValueError(
1109+
f"Partition key '{partition_key}' is not a field of feature view "
1110+
f"'{table_name}'."
1111+
)
1112+
if field.dtype not in PARTITION_KEY_TYPES:
1113+
raise ValueError(
1114+
f"Partition key '{partition_key}' of feature view '{table_name}' is "
1115+
f"stored as {field.dtype.name}, but Milvus partition keys must be "
1116+
"INT64 or VARCHAR."
1117+
)
1118+
field.is_partition_key = True
1119+
1120+
1121+
def _warn_if_partition_key_missing(
1122+
collection: Dict[str, Any], partition_key: str
1123+
) -> None:
1124+
has_key = any(
1125+
field["name"] == partition_key and field.get("is_partition_key")
1126+
for field in collection["fields"]
1127+
)
1128+
if not has_key:
1129+
logger.warning(
1130+
"Collection '%s' was created without partition key '%s'. The partition "
1131+
"key only applies when a collection is created: run `feast teardown` "
1132+
"and `feast apply`, then materialize again to use it.",
1133+
collection["collection_name"],
1134+
partition_key,
1135+
)
1136+
1137+
10711138
def _consistency_kwargs(online_config: MilvusOnlineStoreConfig) -> Dict[str, Any]:
10721139
"""Read and search kwargs; only pass a level when configured."""
10731140
if online_config.consistency_level:

‎sdk/python/tests/integration/online_store/test_milvus_remote.py‎

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222

2323
from feast import Entity, FeatureView
2424
from feast.field import Field
25+
from feast.filter_models import ComparisonFilter
2526
from feast.infra.online_stores.milvus_online_store.milvus import MilvusOnlineStore
2627
from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto
2728
from feast.protos.feast.types.Value_pb2 import Value as ValueProto
@@ -360,3 +361,49 @@ def test_consistency_level(
360361
)
361362
rows = _read(store, config, fv, [1], ["city"])
362363
assert rows[0] is not None and rows[0]["city"].string_val == "Oslo"
364+
365+
366+
def test_partition_key_filtering(
367+
tmp_path: Path, project: str, store: MilvusOnlineStore
368+
) -> None:
369+
config = _repo_config(tmp_path, project, consistency_level="Strong")
370+
fv = FeatureView(
371+
name="products",
372+
entities=[
373+
Entity(
374+
name="driver_id", join_keys=["driver_id"], value_type=ValueType.INT64
375+
)
376+
],
377+
ttl=timedelta(days=1),
378+
schema=[
379+
Field(name="driver_id", dtype=Int64),
380+
Field(name="brand_id", dtype=String),
381+
Field(
382+
name="embedding",
383+
dtype=Array(Float32),
384+
vector_index=True,
385+
vector_search_metric="COSINE",
386+
),
387+
Field(name="city", dtype=String),
388+
],
389+
tags={"milvus.partition_key": "brand_id"},
390+
)
391+
store.update(config, [], [fv], [], [], partial=False)
392+
rows = _vector_rows()
393+
rows[1]["brand_id"] = ValueProto(string_val="acme")
394+
rows[2]["brand_id"] = ValueProto(string_val="globex")
395+
_write_rows(store, config, fv, rows)
396+
397+
assert store.client is not None
398+
fields = store.client.describe_collection(f"{project}_{fv.name}")["fields"]
399+
assert [f["name"] for f in fields if f.get("is_partition_key")] == ["brand_id"]
400+
401+
hits = _search(
402+
store,
403+
config,
404+
fv,
405+
[1.0, 0.0],
406+
top_k=5,
407+
filters=ComparisonFilter(type="eq", key="brand_id", value="globex"),
408+
)
409+
assert [hit["city"].string_val for hit in hits] == ["Rome"]

‎sdk/python/tests/unit/infra/online_store/test_milvus_online_store.py‎

Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212

1313
from feast import Entity, FeatureView
1414
from feast.field import Field
15+
from feast.filter_models import ComparisonFilter
1516
from feast.infra.online_stores.milvus_online_store.milvus import (
1617
PLACEHOLDER_VECTOR_DIM,
1718
PLACEHOLDER_VECTOR_FIELD,
@@ -579,3 +580,169 @@ def test_collection_and_read_consistency_levels_independent() -> None:
579580
def test_invalid_consistency_level_rejected(field: str) -> None:
580581
with pytest.raises(ValidationError):
581582
MilvusOnlineStoreConfig(**{field: "Immediate"})
583+
584+
585+
def _catalog_feature_view(tags: Optional[Dict[str, str]] = None) -> FeatureView:
586+
return FeatureView(
587+
name="products",
588+
entities=[
589+
Entity(
590+
name="product_id", join_keys=["product_id"], value_type=ValueType.INT64
591+
)
592+
],
593+
ttl=timedelta(days=1),
594+
schema=[
595+
Field(name="product_id", dtype=Int64),
596+
Field(name="brand_id", dtype=String),
597+
Field(
598+
name="embedding",
599+
dtype=Array(Float32),
600+
vector_index=True,
601+
vector_search_metric="COSINE",
602+
),
603+
Field(name="title", dtype=String),
604+
Field(name="price", dtype=Float32),
605+
],
606+
tags=tags or {},
607+
)
608+
609+
610+
def _product_key(product_id: int) -> EntityKeyProto:
611+
return EntityKeyProto(
612+
join_keys=["product_id"], entity_values=[ValueProto(int64_val=product_id)]
613+
)
614+
615+
616+
def _write_products(
617+
store: MilvusOnlineStore, config: RepoConfig, fv: FeatureView
618+
) -> None:
619+
def embedding(x: float, y: float) -> ValueProto:
620+
value = ValueProto()
621+
value.float_list_val.val.extend([x, y])
622+
return value
623+
624+
now = datetime.now(timezone.utc)
625+
products = [
626+
(1, "acme", (1.0, 0.0), "Acme kettle"),
627+
(2, "acme", (0.9, 0.1), "Acme toaster"),
628+
(3, "globex", (1.0, 0.0), "Globex kettle"),
629+
]
630+
store.online_write_batch(
631+
config,
632+
fv,
633+
[
634+
(
635+
_product_key(product_id),
636+
{
637+
"brand_id": ValueProto(string_val=brand),
638+
"embedding": embedding(*vector),
639+
"title": ValueProto(string_val=title),
640+
},
641+
now,
642+
now,
643+
)
644+
for product_id, brand, vector, title in products
645+
],
646+
progress=None,
647+
)
648+
649+
650+
def _partition_key_field(
651+
store: MilvusOnlineStore, collection_name: str
652+
) -> Optional[str]:
653+
assert store.client is not None
654+
fields = store.client.describe_collection(collection_name)["fields"]
655+
return next((f["name"] for f in fields if f.get("is_partition_key")), None)
656+
657+
658+
@pytest.mark.parametrize(
659+
"online_store, tags",
660+
[
661+
({}, {"milvus.partition_key": "brand_id"}),
662+
({"partition_key": "brand_id"}, {}),
663+
],
664+
ids=["feature_view_tag", "store_config"],
665+
)
666+
def test_partition_key_filtering(
667+
tmp_path: Path, online_store: Dict[str, Any], tags: Dict[str, str]
668+
) -> None:
669+
config = _lite_config(tmp_path, **online_store)
670+
fv = _catalog_feature_view(tags)
671+
store = MilvusOnlineStore()
672+
store.update(config, [], [fv], [], [], partial=False)
673+
_write_products(store, config, fv)
674+
675+
assert _partition_key_field(store, "test_milvus_products") == "brand_id"
676+
results = store.retrieve_online_documents_v2(
677+
config,
678+
fv,
679+
["embedding", "brand_id", "title"],
680+
embedding=[1.0, 0.0],
681+
top_k=5,
682+
distance_metric="COSINE",
683+
filters=ComparisonFilter(type="eq", key="brand_id", value="acme"),
684+
)
685+
686+
titles = sorted(values["title"].string_val for _, _, values in results if values)
687+
assert titles == ["Acme kettle", "Acme toaster"]
688+
689+
690+
def test_store_partition_key_ignored_for_feature_views_without_the_field(
691+
tmp_path: Path,
692+
) -> None:
693+
config = _lite_config(tmp_path, partition_key="brand_id")
694+
fv = _scalar_feature_view()
695+
store = MilvusOnlineStore()
696+
store.update(config, [], [fv], [], [], partial=False)
697+
698+
assert _partition_key_field(store, "test_milvus_driver_stats") is None
699+
700+
701+
def test_feature_view_tag_overrides_store_partition_key(tmp_path: Path) -> None:
702+
config = _lite_config(tmp_path, partition_key="brand_id")
703+
fv = _catalog_feature_view({"milvus.partition_key": "title"})
704+
store = MilvusOnlineStore()
705+
store.update(config, [], [fv], [], [], partial=False)
706+
707+
assert _partition_key_field(store, "test_milvus_products") == "title"
708+
709+
710+
@pytest.mark.parametrize(
711+
"online_store, partition_key, error",
712+
[
713+
({}, "missing_field", "is not a field"),
714+
# Native numeric storage makes price a FLOAT, which can't be a partition key.
715+
({"enable_openai_compatible_store": True}, "price", "must be INT64 or VARCHAR"),
716+
],
717+
)
718+
@patch(f"{MILVUS_MODULE}.MilvusClient")
719+
def test_invalid_partition_key(
720+
mock_client_cls: MagicMock,
721+
online_store: Dict[str, Any],
722+
partition_key: str,
723+
error: str,
724+
) -> None:
725+
mock_client = _mock_client(mock_client_cls, has_collection=False)
726+
fv = _catalog_feature_view({"milvus.partition_key": partition_key})
727+
728+
with pytest.raises(ValueError, match=error):
729+
MilvusOnlineStore()._get_or_create_collection(_mock_config(**online_store), fv)
730+
mock_client.create_collection.assert_not_called()
731+
732+
733+
@patch(f"{MILVUS_MODULE}.MilvusClient")
734+
def test_warns_when_existing_collection_lacks_partition_key(
735+
mock_client_cls: MagicMock, caplog: pytest.LogCaptureFixture
736+
) -> None:
737+
mock_client = _mock_client(mock_client_cls, has_collection=True)
738+
mock_client.get_load_state.return_value = {"state": LoadState.Loaded}
739+
mock_client.describe_collection.return_value = {
740+
"collection_name": "test_milvus_products",
741+
"fields": [{"name": "brand_id", "type": DataType.VARCHAR, "params": {}}],
742+
}
743+
fv = _catalog_feature_view({"milvus.partition_key": "brand_id"})
744+
745+
MilvusOnlineStore()._get_or_create_collection(_mock_config(), fv)
746+
747+
assert "created without partition key 'brand_id'" in caplog.text
748+
mock_client.create_collection.assert_not_called()

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL