| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
| 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 | ||
|
|
||
| @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( | ||
|
Comment thread
Copy link
Copy Markdown
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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_syncThen _read_hash_fields and _get_glide_client can both use self._get_glide_sync() instead of calling _load_glide_sync() each time.
Sorry, something went wrong.
All reactions
Copy link
Copy Markdown
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low QualityDone 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.
Sorry, something went wrong.
All reactions
|
||
| 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 | ||
| Back | FazBrowse Home | New Git URL |
There was a problem hiding this comment.
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 QualityNit: _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:
Sorry, something went wrong.
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
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 QualityDone 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.
Sorry, something went wrong.
Uh oh!
There was an error while loading. Please reload this page.