Skip to content
14 changes: 8 additions & 6 deletions docs/how-to-guides/online-server-performance-tuning.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,9 @@ When the server processes a `get_online_features()` call, it groups the requeste
**Guideline:** For features that share the same entity key and are frequently requested together, consolidate them into a **single feature view**. This reduces the number of store round-trips per request. Split feature views only when features have different entities, different materialization schedules, or different data source update frequencies.

{% hint style="info" %}
**Redis exception:** The Redis online store overrides `get_online_features()` to batch all `HMGET` commands across every feature view into a **single pipeline execution**. Because all feature views for the same entity share one Redis hash key, the number of Redis round trips is always **1**, regardless of how many feature views the request touches. This means the "fewer feature views" guideline is less critical for Redis than for other stores — but consolidating feature views still reduces serialization and protobuf overhead at the application layer.
**Redis exception:** The Redis online store overrides `_read_features_per_fv()` to batch all `HMGET` commands across every feature view into a **single pipeline execution**. Because all feature views for the same entity share one Redis hash key, the number of Redis round trips is always **1**, regardless of how many feature views the request touches. This means the "fewer feature views" guideline is less critical for Redis than for other stores — but consolidating feature views still reduces serialization and protobuf overhead at the application layer.

**PostgreSQL exception:** PostgreSQL stores each feature view in its own table and overrides `_read_features_per_fv()` to combine the per-view reads into a **single `UNION ALL` statement**, so a request touching any number of feature views costs **1** query on both the sync and async paths. This helps most where round trips dominate — a handful of entities against a remote database. Requests larger than `max_batched_result_rows` (entity rows × feature views, default 2048) fall back to one query per view, because at that size the saved round trips no longer offset holding every view's rows at once. A single-view request is never batched, since there is nothing to combine. Set `max_batched_result_rows: 0` in the online store config to turn batching off. As with Redis, consolidating feature views still reduces application-layer overhead.
{% endhint %}

### Feature services are free (and can be faster)
Expand Down Expand Up @@ -87,7 +89,7 @@ Requesting just `combined_score` triggers reads from **both** `driver_stats_fv`

## Pre-computed feature vectors

When a `get_online_features()` request touches multiple feature views, the server issues a separate store read per feature view. For services spanning 5–15+ feature views, this fan-out dominates latency — even with Redis pipeline batching, the protobuf deserialization and response-building overhead grows linearly with the number of views.
When a `get_online_features()` request touches multiple feature views, the server issues a separate store read per feature view. For services spanning 5–15+ feature views, this fan-out dominates latency — even where the store batches its reads (Redis, PostgreSQL), the protobuf deserialization and response-building overhead grows linearly with the number of views.

**Pre-computed feature vectors** solve this by storing all of a feature service's features for each entity as a single serialized blob. At read time, the server fetches one blob per entity instead of N reads per feature view, reducing the operation to O(1).

Expand Down Expand Up @@ -276,7 +278,7 @@ The online store is the single largest factor in `get_online_features()` latency
| ----- | ------------------- | ---------- | -------- | ------------- |
| **Redis / Dragonfly** | < 1 ms | No (threadpool) | Ultra-low latency, high throughput; all FV reads batched into 1 pipeline | Requires in-memory capacity for your dataset |
| **DynamoDB** | 2–5 ms | Yes | Serverless, auto-scaling on AWS | Pay-per-request cost; batch API limits (100 items) |
| **PostgreSQL** | 3–10 ms | No (threadpool) | Teams with existing Postgres infra | Connection pooling needed at scale |
| **PostgreSQL** | 3–10 ms | No (threadpool) | Teams with existing Postgres infra; all FV reads batched into 1 query | Connection pooling needed at scale |
| **MongoDB** | 2–5 ms | Yes | Flexible schema, async-native | Requires index tuning for large datasets |
| **Aerospike** | < 1 ms | No (threadpool) | Ultra-low latency, hybrid memory (RAM + SSD), large datasets | Namespace must be pre-configured on the cluster |
| **Bigtable** | 3–8 ms | No (threadpool) | Large-scale GCP workloads | Row-key design affects read performance |
Expand Down Expand Up @@ -304,8 +306,8 @@ The feature server can read from the online store using either an **async** or *
| ----- | ---------- | ----------- | ----- |
| **DynamoDB** | Yes | Yes | Uses `aiobotocore` for non-blocking I/O |
| **MongoDB** | Yes | Yes | Uses `motor` (async MongoDB driver) |
| **PostgreSQL** | Implemented | No | Has `online_read_async` but does not yet advertise via `async_supported`; uses sync/threadpool path |
| **Redis** | Implemented | **Yes** | `online_read_async` and `online_write_batch_async` both implemented; uses sync/threadpool path for `get_online_features` (overridden with batched single pipeline) |
| **PostgreSQL** | Implemented | No | Has `online_read_async` but does not yet advertise via `async_supported`; uses sync/threadpool path. `_read_features_per_fv` is overridden to batch all feature view reads into a single `UNION ALL` query |
| **Redis** | Implemented | **Yes** | `online_read_async` and `online_write_batch_async` both implemented; uses sync/threadpool path for `get_online_features` (`_read_features_per_fv` overridden with batched single pipeline) |
| **Aerospike** | Implemented | No | Async methods wrap the blocking C client via `run_in_executor`; does not yet advertise via `async_supported`, so the server still uses the threadpool path |
| All others | No | No | Fall back to sync with `run_in_threadpool()` |

Expand Down Expand Up @@ -387,7 +389,7 @@ online_store:

#### Batched multi-feature-view reads

The Redis online store overrides `get_online_features()` to issue all `HMGET` commands — across every feature view in the request — in a **single pipeline execution**. This reduces Redis round trips from `N` (one per feature view) to `1` regardless of request size.
The Redis online store overrides `_read_features_per_fv()` to issue all `HMGET` commands — across every feature view in the request — in a **single pipeline execution**. This reduces Redis round trips from `N` (one per feature view) to `1` regardless of request size.

| Feature views | Round trips (other stores) | Round trips (Redis) |
| :---: | :---: | :---: |
Expand Down
2 changes: 2 additions & 0 deletions docs/reference/online-stores/postgres.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ The PostgreSQL online store provides support for materializing feature values in

* `sslmode` defaults to `require`, which encrypts the connection without certificate verification. To disable SSL (e.g. for local development), set `sslmode: disable`. For certificate verification, set `sslmode` to `verify-ca` or `verify-full` and provide the corresponding `sslrootcert_path` (and optionally `sslcert_path` and `sslkey_path` for mutual TLS)

* Reads that span several feature views are combined into one query. `max_batched_result_rows` (default 2048) caps the request size, in entity rows × feature views, that is batched this way; set it to `0` to always issue one query per feature view

## Getting started
In order to use this online store, you'll need to run `pip install 'feast[postgres]'`. You can get started by then running `feast init -t postgres`.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
Generator,
List,
Literal,
NamedTuple,
Optional,
Sequence,
Tuple,
Expand All @@ -19,8 +20,9 @@
from psycopg import AsyncConnection, sql
from psycopg.connection import Connection
from psycopg_pool import AsyncConnectionPool, ConnectionPool
from pydantic import NonNegativeInt

from feast import Entity, FeatureView, ValueType
from feast import Entity, FeatureView, ValueType, utils
from feast.filter_models import (
ComparisonFilter,
CompoundFilter,
Expand Down Expand Up @@ -55,6 +57,18 @@
"inner_product": "<#>",
}

# Default for ``max_batched_result_rows``: the largest request, measured as entity
# rows x feature views, that is read in one batched query. Past that, batching stops
# paying for itself: the saved round trips are amortized away by the payload, while
# holding every feature view's rows at once makes peak memory grow with the number of
# views. The measure is a proxy taken from the request shape; the rows actually
# returned also scale with the number of requested features. Measured against
# Postgres 16, a 10-view request is ~3-5x faster at 1-100 entities and a wash at
# 5000, where the combined result costs tens of MB per in-flight request. Past the
# limit the generic one-query-per-view path is used instead, which processes and
# releases a single view at a time.
MAX_BATCHED_RESULT_ROWS = 2048

_PG_COMPARISON_OPS: Dict[str, str] = {
"eq": "=",
"ne": "!=",
Expand All @@ -65,6 +79,16 @@
}


class _BatchedRead(NamedTuple):
"""One feature view's share of a batched online read."""

table: FeatureView
requested_features: List[str]
keys: List[bytes]
idxs: Tuple[List[int], ...]
output_len: int


class PostgresFilterTranslator(FilterTranslator):
"""Translates Feast filters into Postgres SQL WHERE clause fragments."""

Expand Down Expand Up @@ -159,6 +183,9 @@ def _pg_filter_col_and_val(value: Any) -> Tuple[str, Any]:
class PostgreSQLOnlineStoreConfig(PostgreSQLConfig, VectorStoreConfig):
type: Literal["postgres"] = "postgres"
enable_openai_compatible_store: Optional[bool] = False
max_batched_result_rows: NonNegativeInt = MAX_BATCHED_RESULT_ROWS
"""Largest request, as entity rows x feature views, read in one batched query.
Set to 0 to always read one feature view per query."""


class PostgreSQLOnlineStore(OnlineStore):
Expand Down Expand Up @@ -424,6 +451,13 @@ def _process_rows(
row[0] if isinstance(row[0], bytes) else row[0].tobytes()
].append(row[1:])

return PostgreSQLOnlineStore._result_from_values_dict(keys, values_dict)

@staticmethod
def _result_from_values_dict(
keys: List[bytes], values_dict: Dict[bytes, List[Tuple]]
) -> List[Tuple[Optional[datetime], Optional[Dict[str, ValueProto]]]]:
"""Assemble per-entity rows in ``keys`` order from a feature-value mapping."""
result: List[Tuple[Optional[datetime], Optional[Dict[str, ValueProto]]]] = []
for key in keys:
if key in values_dict:
Expand All @@ -438,6 +472,225 @@ def _process_rows(
result.append((None, None))
return result

def _read_features_per_fv(
self,
config: RepoConfig,
grouped_refs: List,
join_key_values: Dict,
entity_name_to_join_key_map: Dict,
online_features_response,
full_feature_names: bool,
include_feature_view_version_metadata: bool,
) -> None:
"""Read every requested feature view in a single round trip.

The generic path issues one query per feature view. Each feature view lives
in its own table here, so the per-view reads can be combined with UNION ALL
and split apart again afterwards.
"""
if not self._should_batch(config, grouped_refs, join_key_values):
return super()._read_features_per_fv(
config,
grouped_refs,
join_key_values,
entity_name_to_join_key_map,
online_features_response,
full_feature_names,
include_feature_view_version_metadata,
)

reads = self._prepare_batched_reads(
config, grouped_refs, join_key_values, entity_name_to_join_key_map
)

query, params = self._construct_batched_query_and_params(config, reads)
with self._get_conn(config, autocommit=True) as conn, conn.cursor() as cur:
cur.execute(query, params)
buckets = self._bucket_rows(reads, cur)

self._populate_from_buckets(
reads,
buckets,
online_features_response,
full_feature_names,
include_feature_view_version_metadata,
)

async def _read_features_per_fv_async(
self,
config: RepoConfig,
grouped_refs: List,
join_key_values: Dict,
entity_name_to_join_key_map: Dict,
online_features_response,
full_feature_names: bool,
include_feature_view_version_metadata: bool,
) -> None:
"""Async version of :meth:`_read_features_per_fv`.

The generic async path issues the per-view queries concurrently. Combining
them still replaces those N queries with one.
"""
if not self._should_batch(config, grouped_refs, join_key_values):
return await super()._read_features_per_fv_async(
config,
grouped_refs,
join_key_values,
entity_name_to_join_key_map,
online_features_response,
full_feature_names,
include_feature_view_version_metadata,
)

reads = self._prepare_batched_reads(
config, grouped_refs, join_key_values, entity_name_to_join_key_map
)

query, params = self._construct_batched_query_and_params(config, reads)
async with self._get_conn_async(config, autocommit=True) as conn:
async with conn.cursor() as cur:
await cur.execute(query, params)
buckets = [defaultdict(list) for _ in reads] # type: List[Dict[bytes, List[Tuple]]]
async for row in cur:
self._bucket_row(buckets, row)

self._populate_from_buckets(
reads,
buckets,
online_features_response,
full_feature_names,
include_feature_view_version_metadata,
)

def _prepare_batched_reads(
self,
config: RepoConfig,
grouped_refs: List,
join_key_values: Dict,
entity_name_to_join_key_map: Dict,
) -> List["_BatchedRead"]:
"""Resolve the entity keys each feature view needs before querying.

Feature views can be keyed on different entities, so the serialized keys are
computed per view rather than shared.
"""
reads = []
for table, requested_features in grouped_refs:
table_entity_values, idxs, output_len = utils._get_unique_entities(
table, join_key_values, entity_name_to_join_key_map
)
entity_key_protos = utils._get_entity_key_protos(table_entity_values)
reads.append(
_BatchedRead(
table=table,
requested_features=requested_features,
keys=self._prepare_keys(
entity_key_protos, config.entity_key_serialization_version
),
idxs=idxs,
output_len=output_len,
)
)
return reads

@staticmethod
def _construct_batched_query_and_params(
config: RepoConfig, reads: List["_BatchedRead"]
) -> Tuple[sql.Composed, List[Any]]:
"""UNION ALL the per feature view reads into one statement.

Each branch selects a constant tag so the combined result can be split back
apart. The table name itself is not in the result set, and two feature views
can return the same entity key, so the tag is what keeps them separable.
"""
versioning = config.registry.enable_online_feature_view_versioning
branches: List[sql.Composed] = []
params: List[Any] = []
for tag, read in enumerate(reads):
table_name = _table_id(config.project, read.table, versioning)
if read.requested_features:
branch = sql.SQL(
"SELECT {tag} AS fv_tag, entity_key, feature_name, value, event_ts "
"FROM {table} WHERE entity_key = ANY(%s) AND feature_name = ANY(%s)"
).format(tag=sql.Literal(tag), table=sql.Identifier(table_name))
params.extend([read.keys, list(read.requested_features)])
else:
branch = sql.SQL(
"SELECT {tag} AS fv_tag, entity_key, feature_name, value, event_ts "
"FROM {table} WHERE entity_key = ANY(%s)"
).format(tag=sql.Literal(tag), table=sql.Identifier(table_name))
params.append(read.keys)
branches.append(branch)

return sql.SQL(" UNION ALL ").join(branches), params

@staticmethod
def _should_batch(
config: RepoConfig, grouped_refs: List, join_key_values: Dict
) -> bool:
"""Whether to read this request's feature views in one batched query.

A single view has no round trips to save, so it keeps the generic path. The
size check is estimated from the request shape rather than from the resolved
keys, so that deciding against batching costs nothing. The number of entity
rows in the request is an upper bound on the unique keys any one view will
read. See :data:`MAX_BATCHED_RESULT_ROWS`.
"""
if len(grouped_refs) < 2:
return False
limit = getattr(
config.online_store, "max_batched_result_rows", MAX_BATCHED_RESULT_ROWS
)
entity_rows = max(
(len(values) for values in join_key_values.values()), default=0
)
return entity_rows * len(grouped_refs) <= limit

@staticmethod
def _bucket_row(buckets: List[Dict[bytes, List[Tuple]]], row: Tuple) -> None:
"""File one row under its feature view's bucket, keyed by entity key.

Only the three payload columns are kept. Slicing the row instead would copy
it a second time while the combined result is still referenced.
"""
entity_key = row[1] if isinstance(row[1], bytes) else row[1].tobytes()
buckets[row[0]][entity_key].append((row[2], row[3], row[4]))

def _bucket_rows(
self, reads: List["_BatchedRead"], row_iter
) -> List[Dict[bytes, List[Tuple]]]:
"""Consume the combined result one row at a time.

The result spans every requested feature view, so calling ``fetchall()`` and
then slicing each row would hold two or three copies of it at peak. Bucketing
as rows arrive keeps a single copy.
"""
buckets: List[Dict[bytes, List[Tuple]]] = [defaultdict(list) for _ in reads]
for row in row_iter:
self._bucket_row(buckets, row)
return buckets

def _populate_from_buckets(
self,
reads: List["_BatchedRead"],
buckets: List[Dict[bytes, List[Tuple]]],
online_features_response,
full_feature_names: bool,
include_feature_view_version_metadata: bool,
) -> None:
"""Hand each feature view's bucket to the response in grouped_refs order."""
for read, bucket in zip(reads, buckets):
utils._populate_response_from_feature_data(
read.requested_features,
self._result_from_values_dict(read.keys, bucket),
read.idxs,
online_features_response,
full_feature_names,
read.table,
read.output_len,
include_feature_view_version_metadata,
)

def update(
self,
config: RepoConfig,
Expand Down
Loading
Loading