Skip to content

Commit 176ed0d

Browse files
anshishrivastavaShizoqua
authored andcommitted
feat: Add DynamoDB in-place list update support for array-based features (feast-dev#5916)
Signed-off-by: Shizoqua <hr.lanreshittu@gmail.com>
1 parent 5d63092 commit 176ed0d

4 files changed

Lines changed: 777 additions & 1 deletion

File tree

‎sdk/python/feast/feature_store.py‎

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2013,6 +2013,91 @@ async def write_to_online_store_async(
20132013
provider = self._get_provider()
20142014
await provider.ingest_df_async(feature_view, df)
20152015

2016+
async def update_online_store(
2017+
self,
2018+
feature_view_name: str,
2019+
df: pd.DataFrame,
2020+
update_expressions: Dict[str, str],
2021+
allow_registry_cache: bool = True,
2022+
) -> None:
2023+
"""
2024+
Update features using DynamoDB-specific list operations.
2025+
2026+
This method provides efficient in-place list updates using DynamoDB's native
2027+
UpdateItem operations with list_append and other expressions. This is more
2028+
efficient than the read-modify-write pattern for array-based features.
2029+
2030+
Args:
2031+
feature_view_name: The feature view to update.
2032+
df: DataFrame with new values to append/prepend to existing lists.
2033+
update_expressions: Dict mapping feature names to DynamoDB update expressions.
2034+
Examples:
2035+
- {"transactions": "list_append(transactions, :new_val)"} # append
2036+
- {"recent_items": "list_append(:new_val, recent_items)"} # prepend
2037+
allow_registry_cache: Whether to allow cached registry.
2038+
2039+
Raises:
2040+
NotImplementedError: If online store doesn't support update expressions.
2041+
ValueError: If the feature view or update expressions are invalid.
2042+
2043+
Example:
2044+
# Append new transactions to existing transaction history
2045+
await store.update_online_store(
2046+
feature_view_name="user_transactions",
2047+
df=new_transactions_df,
2048+
update_expressions={
2049+
"transaction_history": "list_append(transaction_history, :new_val)",
2050+
"recent_amounts": "list_append(:new_val, recent_amounts)" # prepend
2051+
}
2052+
)
2053+
"""
2054+
# Check if online store supports update expressions
2055+
provider = self._get_provider()
2056+
if not hasattr(provider.online_store, "update_online_store_async"):
2057+
raise NotImplementedError(
2058+
f"Online store {type(provider.online_store).__name__} "
2059+
"does not support async update expressions. This feature is only available "
2060+
"with DynamoDB online store."
2061+
)
2062+
2063+
feature_view, df = self._get_feature_view_and_df_for_online_write(
2064+
feature_view_name=feature_view_name,
2065+
df=df,
2066+
allow_registry_cache=allow_registry_cache,
2067+
transform_on_write=False, # Don't transform for updates
2068+
)
2069+
2070+
# Validate that the dataframe has meaningful feature data
2071+
if df is not None:
2072+
if df.empty:
2073+
warnings.warn("Cannot update with empty dataframe")
2074+
return
2075+
2076+
# Check if feature columns are empty
2077+
feature_column_names = [f.name for f in feature_view.features]
2078+
if feature_column_names:
2079+
feature_df = df[feature_column_names]
2080+
if feature_df.empty or feature_df.isnull().all().all():
2081+
warnings.warn("Cannot update with empty feature columns")
2082+
return
2083+
2084+
# Prepare data for online store
2085+
from feast.infra.passthrough_provider import PassthroughProvider
2086+
2087+
rows_to_write = PassthroughProvider._prep_rows_to_write_for_ingestion(
2088+
feature_view=feature_view,
2089+
df=df,
2090+
)
2091+
2092+
# Call DynamoDB-specific async method
2093+
await provider.online_store.update_online_store_async(
2094+
config=self.config,
2095+
table=feature_view,
2096+
data=rows_to_write,
2097+
update_expressions=update_expressions,
2098+
progress=None,
2099+
)
2100+
20162101
def write_to_offline_store(
20172102
self,
20182103
feature_view_name: str,

0 commit comments

Comments
 (0)