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

feat: Add optional GLIDE client for Redis online reads by chandlerok · Pull Request #6870 · feast-dev/feast · GitHub

Repository navigation

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

Filter by extension

Filter by extension .md  (1) .py  (3) .toml  (1) All 3 file types selected
Only manifest files
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
30 changes: 30 additions & 0 deletions docs/reference/online-stores/redis.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
Expand Up @@ -112,6 +112,34 @@ Benchmark results against Redis 8.6.2 (localhost, 50 entities, 3 features/FV, 30

The speedup grows with the number of feature views and is most pronounced in production environments with non-trivial network RTT to Redis.

### GLIDE client (`client: glide`)

By default the Redis online store uses redis-py, which packs every command and parses every reply in Python. Each socket receive releases and reacquires the GIL, so in a process running other Python threads each reacquire waits on the switch interval. With many entities and feature views per read, that cost dominates.

Setting `client: glide` runs the same `HMGET` commands as one non-atomic [GLIDE](https://github.com/valkey-io/valkey-glide) batch. GLIDE has a Rust core, so the batch crosses the FFI boundary once and the fetch runs with the GIL released.

```yaml
online_store:
type: redis
connection_string: "localhost:6379"
client: glide
```

Install the extra to use it:

```bash
pip install 'feast[glide]'
```

Notes:

* GLIDE reads the same keys and hash fields Feast already writes, so an existing store needs no migration and no rewrite of data.
* The default stays `client: redis`. Writes and `get_online_features_async` keep using redis-py, so only the synchronous read path changes.
* GLIDE supports `host:port`, `db`, `password`, `username`, `ssl`, `socket_timeout`, and `socket_connect_timeout` from `connection_string`. Other parameters are logged and ignored.
* GLIDE has no Sentinel support, so `client: glide` cannot be combined with `redis_type: redis_sentinel`. Use `client: redis` for Sentinel deployments.

If `client: glide` is set without the extra installed, Feast raises an import error naming the extra to install.

### Write path: `skip_dedup` for bulk loads

By default, `online_write_batch()` checks existing timestamps before writing (to avoid overwriting newer data with older data). This requires two pipeline round trips per batch: one to read existing timestamps, one to write new values.
Expand Down Expand Up @@ -144,6 +172,7 @@ The Redis online store implements `online_write_batch_async()` using the async R
| `key_ttl_seconds` | `null` | Redis `EXPIRE` TTL in seconds applied to the entity hash key after each write. Expires all feature views for that entity together. |
| `full_scan_for_deletion` | `true` | When `true`, deleting or renaming a feature view scans Redis to remove its hash fields. Set `false` to skip deletion scans (faster `feast apply`, but leaves orphaned data). |
| `skip_dedup` | `false` | When `true`, skips the existing-timestamp read before each write, halving write round trips. Suitable for initial bulk loads; may cause older values to overwrite newer ones under concurrent writers. |
| `client` | `redis` | Read client: `redis` for redis-py, or `glide` to run the synchronous batched read as one non-atomic GLIDE batch. Requires `pip install 'feast[glide]'`. |

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

Expand Down Expand Up @@ -172,5 +201,6 @@ Below is a matrix indicating which functionality is supported by the Redis onlin
| collocated by entity key | yes |
| async batch writes | yes |
| batched multi-feature-view reads (single pipeline) | yes |
| optional GLIDE client for synchronous reads | yes |

To compare this set of functionality against other online stores, please see the full [functionality matrix](overview.md#functionality-matrix).
3 changes: 2 additions & 1 deletion pyproject.toml
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 @@ -80,6 +80,7 @@ gcp = [
"fsspec>=2024.1.0",
]
ge = ["great_expectations>=0.15.41,<1"]
glide = ["valkey-glide-sync>=2.5,<3"]
go = ["cffi>=1.15.0"]
grpcio = [
"grpcio>=1.56.2",
Expand Down Expand Up @@ -166,7 +167,7 @@ test = [
]

ci = [
"feast[test, aws, azure, cassandra, clickhouse, couchbase, delta, docling, duckdb, elasticsearch, faiss, gcp, ge, go, grpcio, hazelcast, hbase, ibis, iceberg, image, k8s, mcp, milvus, mlflow, mongodb, mssql, mysql, openlineage, opentelemetry, oracle, scylladb, spark, trino, postgres, pytorch, qdrant, rag, ray, redis, singlestore, snowflake, sqlite_vec]",
"feast[test, aws, azure, cassandra, clickhouse, couchbase, delta, docling, duckdb, elasticsearch, faiss, gcp, ge, glide, go, grpcio, hazelcast, hbase, ibis, iceberg, image, k8s, mcp, milvus, mlflow, mongodb, mssql, mysql, openlineage, opentelemetry, oracle, scylladb, spark, trino, postgres, pytorch, qdrant, rag, ray, redis, singlestore, snowflake, sqlite_vec]",
"build",
"virtualenv==20.23.0",
"dbt-artifacts-parser",
Expand Down
213 changes: 192 additions & 21 deletions sdk/python/feast/infra/online_stores/redis.py
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 @@ -32,6 +32,7 @@
from pydantic import StrictStr

from feast import Entity, FeatureView, RepoConfig, utils
from feast.errors import FeastExtrasDependencyImportError
from feast.infra.key_encoding_utils import serialize_entity_key
from feast.infra.online_stores.helpers import (
_mmh3,
Expand All @@ -51,8 +52,6 @@
from redis.cluster import ClusterNode, RedisCluster
from redis.sentinel import Sentinel
except ImportError as e:
from feast.errors import FeastExtrasDependencyImportError

raise FeastExtrasDependencyImportError("redis", str(e))

logger = logging.getLogger(__name__)
Expand All @@ -71,6 +70,11 @@ class RedisType(str, Enum):
redis_sentinel = "redis_sentinel"


class RedisClient(str, Enum):
redis = "redis"
glide = "glide"


class RedisOnlineStoreConfig(FeastConfigBaseModel):
"""Online store config for Redis store"""

Expand All @@ -80,6 +84,15 @@ class RedisOnlineStoreConfig(FeastConfigBaseModel):
redis_type: RedisType = RedisType.redis
"""Redis type: redis or redis_cluster"""

client: RedisClient = RedisClient.redis
"""Redis client used for online reads: ``redis`` (redis-py, default) or ``glide``.

With ``glide``, the synchronous read paths (``online_read`` and the batched read
behind ``get_online_features``) issue their HMGET commands as one non-atomic GLIDE
batch, so the fetch runs off the GIL. Requires ``pip install 'feast[glide]'``.
Writes and async reads always use redis-py. GLIDE has no Sentinel support, so
``client: glide`` cannot be combined with ``redis_type: redis_sentinel``."""

sentinel_master: StrictStr = "mymaster"
"""Sentinel's master name"""

Expand All @@ -99,6 +112,115 @@ class RedisOnlineStoreConfig(FeastConfigBaseModel):
This may cause older feature values to overwrite newer ones under concurrent writers."""


def _seconds_to_millis(value: Any) -> Optional[int]:
"""Convert a redis-py timeout in seconds to GLIDE's millisecond unit."""
if value is None:
return None
return int(float(value) * 1000)


def _load_glide_sync():
"""Import ``valkey-glide-sync`` on demand.

The native extension is only loaded when ``client: glide`` is configured, so the
redis extra keeps working for users who did not opt in.
"""
try:
import glide_sync
except ImportError as e:
raise FeastExtrasDependencyImportError("glide", str(e)) from e
return glide_sync


_GLIDE_SUPPORTED_CONNECTION_PARAMS = frozenset(
{"ssl", "db", "password", "username", "socket_timeout", "socket_connect_timeout"}
)


def _glide_client_config(glide_sync, online_store_config: RedisOnlineStoreConfig):
"""Translate a Feast Redis connection string into a GLIDE client configuration."""
if online_store_config.redis_type == RedisType.redis_sentinel:
raise ValueError(
"client: glide does not support redis_type: redis_sentinel. "
"Use client: redis for Sentinel deployments."
)

startup_nodes, params = RedisOnlineStore._parse_connection_string(
online_store_config.connection_string
)
ignored = sorted(set(params) - _GLIDE_SUPPORTED_CONNECTION_PARAMS)
if ignored:
logger.warning(
"client: glide ignores these Redis connection_string parameters: %s",
", ".join(ignored),
)

credentials = None
password = params.get("password")
username = params.get("username")
if password is not None or username is not None:
credentials = glide_sync.ServerCredentials(
# _parse_connection_string json-decodes values, so a numeric password can
# arrive as an int while GLIDE expects a str.
password=None if password is None else str(password),
username=None if username is None else str(username),
)

connection_timeout = _seconds_to_millis(params.get("socket_connect_timeout"))
kwargs: Dict[str, Any] = {
"addresses": [
glide_sync.NodeAddress(host=node["host"], port=int(node["port"]))
for node in startup_nodes
],
"use_tls": bool(params.get("ssl", False)),
"credentials": credentials,
"request_timeout": _seconds_to_millis(params.get("socket_timeout")),
"advanced_config": glide_sync.AdvancedGlideClientConfiguration(
connection_timeout=connection_timeout
)
if connection_timeout is not None
else None,
}

if online_store_config.redis_type == RedisType.redis_cluster:
# Cluster nodes do not support SELECT, so `db` is not passed along.
return glide_sync.GlideClusterClientConfiguration(**kwargs)

if "db" in params:
kwargs["database_id"] = int(params["db"])
return glide_sync.GlideClientConfiguration(**kwargs)


def _create_glide_client(glide_sync, online_store_config: RedisOnlineStoreConfig):
"""Connect a GLIDE client built from the Redis online store config."""
config = _glide_client_config(glide_sync, online_store_config)
if online_store_config.redis_type == RedisType.redis_cluster:
return glide_sync.GlideClusterClient.create(config)
return glide_sync.GlideClient.create(config)


def _glide_hmget_batch(
glide_sync,
client,
commands: List[Tuple[bytes, List[Union[str, bytes]]]],
) -> List[List[Optional[bytes]]]:
"""Run one HMGET per ``(key, fields)`` pair in a single non-atomic GLIDE batch.

Replies keep the order of ``commands``, so callers can slice them exactly like a
redis-py pipeline result.
"""
batch = (
glide_sync.ClusterBatch(is_atomic=False)
if isinstance(client, glide_sync.GlideClusterClient)
else glide_sync.Batch(is_atomic=False)
)
for key, fields in commands:
batch.hmget(key, fields)

result = client.exec(batch, raise_on_error=True)
return list(result) if result else []


class RedisOnlineStore(OnlineStore):
"""
Redis implementation of the online store interface.
Expand All @@ -114,6 +236,10 @@ class RedisOnlineStore(OnlineStore):
_client_async: Optional[Union[redis_asyncio.Redis, redis_asyncio.RedisCluster]] = (
None
)
# Only set when `client: glide` is configured, so the native extension stays
# untouched for users who never opt in.
_glide_client: Optional[Any] = None
_glide_sync: Optional[Any] = None

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Nit: _glide_client is cached here but never cleaned up. The existing teardown() closes the redis-py client — the GLIDE client should get the same treatment to release native (Rust-side) resources, file descriptors, and the connection pool.

Something like:

def teardown(self, config, tables, entities):
    ...
    if self._glide_client:
        self._glide_client.close()
        self._glide_client = None

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Done in d12d035. teardown() now closes self._glide_client and sets it back to None, so the native client is released as you suggested.

One note: the redis-py client was not actually being closed in teardown() before this change either (the method only deletes keys), so I scoped this to the GLIDE client rather than also changing the redis-py lifecycle in this PR. Happy to close that one too if you'd prefer it here.

@property
def async_supported(self) -> SupportedAsyncMethods:
Expand Down Expand Up @@ -225,6 +351,12 @@ def teardown(
for join_keys in join_keys_to_delete:
self.delete_entity_values(config, list(join_keys))

# Release the native GLIDE client, if one was ever opened, so teardown does
# not leave its Rust-side resources and file descriptors behind.
if self._glide_client:
self._glide_client.close()
self._glide_client = None

@staticmethod
def _parse_connection_string(connection_string: str):
"""
Expand Down Expand Up @@ -305,6 +437,20 @@ async def _get_client_async(self, online_store_config: RedisOnlineStoreConfig):
self._client_async = redis_asyncio.Redis(**kwargs)
return self._client_async

def _get_glide_sync(self):
"""Returns the cached ``glide_sync`` module, importing it on first use."""
if self._glide_sync is None:
self._glide_sync = _load_glide_sync()
return self._glide_sync

def _get_glide_client(self, online_store_config: RedisOnlineStoreConfig):
"""Creates and caches a GLIDE client built from the online store config."""
if not self._glide_client:
self._glide_client = _create_glide_client(
self._get_glide_sync(), online_store_config
)
return self._glide_client

def online_write_batch(
self,
config: RepoConfig,
Expand Down Expand Up @@ -538,12 +684,14 @@ def _generate_hset_keys_for_features(
feature_view: FeatureView,
requested_features: Optional[List[str]] = None,
fv_name_override: Optional[str] = None,
) -> Tuple[List[str], List[str]]:
) -> Tuple[List[str], List[Union[str, bytes]]]:
if not requested_features:
requested_features = [f.name for f in feature_view.features]

fv_name = fv_name_override or feature_view.name
hset_keys = [_mmh3(f"{fv_name}:{k}") for k in requested_features]
hset_keys: List[Union[str, bytes]] = [
_mmh3(f"{fv_name}:{k}") for k in requested_features
]

ts_key = f"_ts:{fv_name}"
hset_keys.append(ts_key)
Expand All @@ -553,7 +701,7 @@ def _generate_hset_keys_for_features(

def _convert_redis_values_to_protobuf(
self,
redis_values: List[List[ByteString]],
redis_values: Sequence[Sequence[Optional[ByteString]]],
feature_view: str,
requested_features: List[str],
) -> List[Tuple[Optional[datetime], Optional[Dict[str, ValueProto]]]]:
Expand All @@ -562,17 +710,40 @@ def _convert_redis_values_to_protobuf(
for values in redis_values
]

def _read_hash_fields(
self,
config: RepoConfig,
commands: List[Tuple[bytes, List[Union[str, bytes]]]],
) -> List[List[Optional[ByteString]]]:
"""Fetch one HMGET reply per ``(key, fields)`` pair, in the given order.

Defaults to a single redis-py pipeline. When ``client: glide`` is configured
the same commands run as one non-atomic GLIDE batch, so the fetch runs off the
GIL instead of in Python per command.
"""
online_store_config = config.online_store
assert isinstance(online_store_config, RedisOnlineStoreConfig)

if online_store_config.client == RedisClient.glide:
return _glide_hmget_batch(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

_load_glide_sync() is called on every read when client == RedisClient.glide. While Python caches imported modules in sys.modules, this still incurs a function-call + dict-lookup overhead per read that is easy to avoid.

Consider caching the module reference on the instance alongside _glide_client, e.g.:

_glide_sync: Optional[Any] = None

def _get_glide_sync(self):
    if self._glide_sync is None:
        self._glide_sync = _load_glide_sync()
    return self._glide_sync

Then _read_hash_fields and _get_glide_client can both use self._get_glide_sync() instead of calling _load_glide_sync() each time.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Done in d12d035. Added _get_glide_sync(), which imports the module once and caches it on the instance; both _read_hash_fields and _get_glide_client now go through it instead of calling _load_glide_sync() each read.

self._get_glide_sync(),
self._get_glide_client(online_store_config),
commands,
)

client = self._get_client(online_store_config)
with client.pipeline(transaction=False) as pipe:
for redis_key, fields in commands:
pipe.hmget(redis_key, fields)
return pipe.execute()

def online_read(
self,
config: RepoConfig,
table: FeatureView,
entity_keys: List[EntityKeyProto],
requested_features: Optional[List[str]] = None,
) -> List[Tuple[Optional[datetime], Optional[Dict[str, ValueProto]]]]:
online_store_config = config.online_store
assert isinstance(online_store_config, RedisOnlineStoreConfig)

client = self._get_client(online_store_config)
feature_view = table
fv_name = _versioned_fv_name(table, config)

Expand All @@ -581,11 +752,9 @@ def online_read(
)
keys = self._generate_redis_keys_for_entities(config, entity_keys)

with client.pipeline(transaction=False) as pipe:
for redis_key_bin in keys:
pipe.hmget(redis_key_bin, hset_keys)

redis_values = pipe.execute()
redis_values = self._read_hash_fields(
config, [(redis_key, hset_keys) for redis_key in keys]
)

return self._convert_redis_values_to_protobuf(
redis_values, fv_name, requested_features
Expand Down Expand Up @@ -648,12 +817,14 @@ def _read_features_per_fv(
)

if work_items:
client = self._get_client(config.online_store)
with client.pipeline(transaction=False) as pipe:
for _, _, _, hset_keys, redis_keys, _, _ in work_items:
for redis_key in redis_keys:
pipe.hmget(redis_key, hset_keys)
all_results = pipe.execute()
all_results = self._read_hash_fields(
config,
[
(redis_key, hset_keys)
for _, _, _, hset_keys, redis_keys, _, _ in work_items
for redis_key in redis_keys
],
)

offset = 0
for (
Expand Down Expand Up @@ -751,7 +922,7 @@ async def _read_features_per_fv_async(

def _get_features_for_entity(
self,
values: List[ByteString],
values: Sequence[Optional[ByteString]],
feature_view: str,
requested_features: List[str],
) -> Tuple[Optional[datetime], Optional[Dict[str, ValueProto]]]:
Expand Down
Loading
Loading

Back | FazBrowse Home | New Git URL