diff --git a/docs/reference/feature-servers/python-feature-server.md b/docs/reference/feature-servers/python-feature-server.md index b1b873cc7d2..34d62a007d7 100644 --- a/docs/reference/feature-servers/python-feature-server.md +++ b/docs/reference/feature-servers/python-feature-server.md @@ -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) | @@ -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` diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index eb85234822f..d84278c7879 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -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, @@ -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"], @@ -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)], @@ -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()), [])) @@ -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" diff --git a/sdk/python/feast/infra/feature_servers/base_config.py b/sdk/python/feast/infra/feature_servers/base_config.py index 4a534210712..b16d14686bb 100644 --- a/sdk/python/feast/infra/feature_servers/base_config.py +++ b/sdk/python/feast/infra/feature_servers/base_config.py @@ -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 @@ -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, diff --git a/sdk/python/tests/unit/test_metrics.py b/sdk/python/tests/unit/test_metrics.py index a3a65ebe782..67b91d0b4ac 100644 --- a/sdk/python/tests/unit/test_metrics.py +++ b/sdk/python/tests/unit/test_metrics.py @@ -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.""" @@ -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 ( @@ -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 @@ -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