Skip to content
Merged
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
24 changes: 23 additions & 1 deletion docs/reference/feature-servers/python-feature-server.md
Original file line number Diff line number Diff line change
Expand Up @@ -368,7 +368,7 @@ thread from starting). All categories default to `true` except
| `feast_feature_server_cpu_usage` | Gauge | — | `resource` | Process CPU usage % |
| `feast_feature_server_memory_usage` | Gauge | — | `resource` | Process memory usage % |
| `feast_feature_server_request_total` | Counter | `endpoint`, `status` | `request` | Total requests per endpoint |
| `feast_feature_server_request_latency_seconds` | Histogram | `endpoint`, `feature_count`, `feature_view_count` | `request` | Request latency with p50/p95/p99 support |
| `feast_feature_server_request_latency_seconds` | Histogram | `endpoint`, `feature_count`, `feature_view_count` | `request` | Request latency with p50/p95/p99 support (`feature_count` is bucketed, see below) |
| `feast_online_features_request_total` | Counter | — | `online_features` | Total online feature retrieval requests |
| `feast_online_features_entity_count` | Histogram | — | `online_features` | Entity rows per online feature request |
| `feast_feature_server_online_store_read_duration_seconds` | Histogram | — | `online_features` | Online store read phase duration (sync and async) |
Expand All @@ -382,6 +382,28 @@ thread from starting). All categories default to `true` except
| `feast_offline_store_request_latency_seconds` | Histogram | `method` | `offline_features` | Latency of offline store retrieval operations |
| `feast_offline_store_row_count` | Histogram | `method` | `offline_features` | Rows returned by offline store retrieval |

### Feature count bucketing

The `feature_count` label on `feast_feature_server_request_latency_seconds`
is bucketed rather than exact, to keep cardinality bounded for feature
services that vary widely in how many features they request. By default,
counts are grouped into `0`, `1-10`, `11-50`, `51-200`, and `201+`.

Customize the bucket boundaries with `feature_count_bins` in the `metrics`
block:

```yaml
feature_server:
type: local
metrics:
enabled: true
feature_count_bins: [5, 20]
```

This produces the labels `0`, `1-5`, `6-20`, and `21+`. Note that this only
affects the Prometheus label; the `feature_count` field in audit logs (see
below) always reports the exact count.

### Per-ODFV transformation metrics

The `transformation_duration_seconds` and `write_transformation_duration_seconds`
Expand Down
33 changes: 32 additions & 1 deletion sdk/python/feast/feature_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@
from feast.feature_view import FeatureViewState
from feast.feature_view_utils import get_feature_view_from_feature_store
from feast.filter_models import ComparisonFilter, CompoundFilter
from feast.infra.feature_servers.base_config import MetricsConfig
from feast.permissions.action import WRITE, AuthzedAction
from feast.permissions.security_manager import (
assert_permissions,
Expand Down Expand Up @@ -238,6 +239,21 @@ def _resolve_feature_counts(
return str(feat_count), str(len(fv_names))


def bin_feature_count(count: int, bins: List[int]) -> str:
"""Map a raw feature count to an inclusive range label."""
if count == 0:
return "0"

lower = 1

for upper in bins:
if count <= upper:
return f"{lower}-{upper}"
lower = upper + 1

return f"{lower}+"


def _emit_online_audit(
request: GetOnlineFeaturesRequest,
features: Union[List[str], "feast.FeatureService"],
Expand Down Expand Up @@ -632,6 +648,16 @@ async def lifespan(app: FastAPI):

app.add_middleware(AuditLoggingMiddleware)

fs_cfg = getattr(store.config, "feature_server", None)
metrics_cfg = getattr(fs_cfg, "metrics", None)

default_feature_count_bins = MetricsConfig().feature_count_bins
feature_count_bins = (
getattr(metrics_cfg, "feature_count_bins", default_feature_count_bins)
if metrics_cfg is not None
else default_feature_count_bins
)

@app.post(
"/get-online-features",
dependencies=[Depends(inject_user_details)],
Expand All @@ -643,7 +669,11 @@ async def get_online_features(request: GetOnlineFeaturesRequest) -> Any:
) as metrics_ctx:
features = await _get_features(request, store)
feat_count, fv_count = _resolve_feature_counts(features)
metrics_ctx.feature_count = feat_count

metrics_ctx.feature_count = bin_feature_count(
int(feat_count),
feature_count_bins,
)
metrics_ctx.feature_view_count = fv_count

entity_count = len(next(iter(request.entities.values()), []))
Expand Down Expand Up @@ -1288,6 +1318,7 @@ def start_server(

fs_cfg = getattr(store.config, "feature_server", None)
metrics_cfg = getattr(fs_cfg, "metrics", None)

metrics_from_config = getattr(metrics_cfg, "enabled", False)
metrics_active = metrics or metrics_from_config
uses_gunicorn = sys.platform != "win32"
Expand Down
20 changes: 19 additions & 1 deletion sdk/python/feast/infra/feature_servers/base_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
# limitations under the License.
from typing import Literal, Optional

from pydantic import StrictBool, StrictInt
from pydantic import Field, StrictBool, StrictInt, field_validator

from feast.repo_config import FeastConfigBaseModel

Expand Down Expand Up @@ -59,6 +59,24 @@ class MetricsConfig(FeastConfigBaseModel):
(feast_feature_server_request_total,
feast_feature_server_request_latency_seconds)."""

feature_count_bins: list[int] = Field(default_factory=lambda: [10, 50, 200])
"""Upper bounds used to bucket the ``feature_count`` label in request latency metrics.

For example, ``[10, 50, 200]`` produces labels
``1-10``, ``11-50``, ``51-200``, and ``201+``.
"""

@field_validator("feature_count_bins")
@classmethod
def validate_feature_count_bins(cls, bins: list[int]) -> list[int]:
if any(bound <= 0 for bound in bins):
raise ValueError("feature_count_bins must contain only positive integers")

if any(lower >= upper for lower, upper in zip(bins, bins[1:])):
raise ValueError("feature_count_bins must be strictly increasing")

return bins

online_features: StrictBool = True
"""Emit online feature retrieval metrics
(feast_online_features_request_total,
Expand Down
130 changes: 125 additions & 5 deletions sdk/python/tests/unit/test_metrics.py
Original file line number Diff line number Diff line change
Expand Up @@ -741,6 +741,82 @@ def test_feature_service(self):
assert fv_count == "2"


class TestBinFeatureCount:
@pytest.mark.parametrize(
("count", "expected"),
[
(0, "0"),
(1, "1-10"),
(10, "1-10"),
(11, "11-50"),
(50, "11-50"),
(51, "51-200"),
(200, "51-200"),
(201, "201+"),
],
)
def test_default_boundaries(self, count, expected):
from feast.feature_server import bin_feature_count

assert bin_feature_count(count, [10, 50, 200]) == expected

@pytest.mark.parametrize(
("count", "expected"),
[
(1, "1-5"),
(5, "1-5"),
(6, "6-20"),
(20, "6-20"),
(21, "21+"),
],
)
def test_custom_boundaries(self, count, expected):

from feast.feature_server import bin_feature_count

assert bin_feature_count(count, [5, 20]) == expected


class TestMetricsConfig:
def test_feature_count_bins_default(self):
from feast.infra.feature_servers.base_config import MetricsConfig

config = MetricsConfig()
assert config.feature_count_bins == [10, 50, 200]

def test_feature_count_bins_custom(self):
from feast.infra.feature_servers.base_config import MetricsConfig

config = MetricsConfig(feature_count_bins=[5, 20])
assert config.feature_count_bins == [5, 20]

@pytest.mark.parametrize(
"bins",
[
[0, 10, 50],
[-1, 10, 50],
],
)
def test_feature_count_bins_must_be_positive(self, bins):
from feast.infra.feature_servers.base_config import MetricsConfig

with pytest.raises(ValueError, match="positive"):
MetricsConfig(feature_count_bins=bins)

@pytest.mark.parametrize(
"bins",
[
[50, 10, 200],
[10, 10, 200],
],
)
def test_feature_count_bins_must_be_strictly_increasing(self, bins):
from feast.infra.feature_servers.base_config import MetricsConfig

with pytest.raises(ValueError, match="strictly increasing"):
MetricsConfig(feature_count_bins=bins)


class TestFeatureServerMetricsIntegration:
"""Test that feature server endpoints record metrics."""

Expand All @@ -751,6 +827,11 @@ def mock_fs_factory(self):
def builder(**async_support):
provider = FooProvider.with_async_support(**async_support)
fs = MagicMock()

from feast.infra.feature_servers.base_config import MetricsConfig

fs.config.feature_server.metrics = MetricsConfig()

fs._get_provider.return_value = provider
from feast.online_response import OnlineResponse
from feast.protos.feast.serving.ServingService_pb2 import (
Expand Down Expand Up @@ -798,20 +879,30 @@ def test_get_online_features_records_metrics(self, mock_fs_factory):
@pytest.mark.parametrize(
"features,expected_feat_count,expected_fv_count",
[
(["fv1:a"], "1", "1"),
(["fv1:a", "fv1:b", "fv2:c"], "3", "2"),
(["fv1:a"], "1-10", "1"),
(["fv1:a", "fv1:b", "fv2:c"], "1-10", "2"),
(
["fv1:a", "fv1:b", "fv2:c", "fv2:d", "fv3:e"],
"5",
"1-10",
"3",
),
(
[f"fv1:f{i}" for i in range(11)],
"11-50",
"1",
),
],
ids=[
"1_feat_1_fv",
"3_feats_2_fvs",
"5_feats_3_fvs",
"11_feats_1_fv",
],
ids=["1_feat_1_fv", "3_feats_2_fvs", "5_feats_3_fvs"],
)
def test_latency_labels_with_varying_request_sizes(
self, mock_fs_factory, features, expected_feat_count, expected_fv_count
):
"""Verify feature_count and feature_view_count labels change with request size."""
"""Verify feature_count is bucketed while feature_view_count remains exact."""
from fastapi.testclient import TestClient

from feast.feature_server import get_app
Expand All @@ -837,6 +928,35 @@ def test_latency_labels_with_varying_request_sizes(
after_sum = request_latency.labels(**label_set)._sum.get()
assert after_sum > before_sum

def test_latency_labels_use_custom_feature_count_bins(self, mock_fs_factory):
from fastapi.testclient import TestClient

from feast.feature_server import get_app
from feast.infra.feature_servers.base_config import MetricsConfig

fs = mock_fs_factory(online_read=False)
fs.config.feature_server.metrics = MetricsConfig(feature_count_bins=[2, 4])

client = TestClient(get_app(fs))

label_set = dict(
endpoint="/get-online-features",
feature_count="3-4",
feature_view_count="1",
)
before_sum = request_latency.labels(**label_set)._sum.get()

client.post(
"/get-online-features",
json={
"features": ["fv:a", "fv:b", "fv:c"],
"entities": {"id": [1]},
},
)

after_sum = request_latency.labels(**label_set)._sum.get()
assert after_sum > before_sum

def test_push_records_metrics(self, mock_fs_factory):
from fastapi.testclient import TestClient

Expand Down