Enhance Redis activity by formatting set and get operations for improved readability and consistency. Update tests to include metadata in assertions for better verification of Redis interactions.
314 lines
12 KiB
Python
314 lines
12 KiB
Python
from sientia_do.observability.metrics_controller import MetricsController
|
|
from temporalio import activity, workflow
|
|
|
|
with workflow.unsafe.imports_passed_through():
|
|
import traceback
|
|
from collections.abc import Hashable
|
|
from typing import Any
|
|
|
|
from pandas import DataFrame
|
|
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
|
|
from sientia_do.notifications.models import NotificationLevel
|
|
from sientia_do.observability.logger import Logger
|
|
from sientia_do.observability.sientia_monitoring import SientiaMonitoring
|
|
from sientia_do.repository.redis_repository import RedisRepository
|
|
from sientia_do.temporal.constants import DATETIME_FORMAT, now
|
|
|
|
|
|
class Redis(SientiaMonitoring):
|
|
"""
|
|
Redis operations for data caching and temporary storage.
|
|
|
|
This class extends the base Redis functionality to provide specialized
|
|
operations for the Scouter system, including:
|
|
- Data timestamp management for incremental processing
|
|
- Temporary data storage with configurable TTL
|
|
- Data grouping and holding for batch processing
|
|
- Error handling and notification integration
|
|
|
|
The class implements Temporal activities for Redis operations, enabling
|
|
distributed data processing with fault tolerance and monitoring.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
host: str,
|
|
port: int,
|
|
username: str,
|
|
password: str,
|
|
logger: Logger,
|
|
notification_handler: NotificationHandler,
|
|
metrics_controller: MetricsController,
|
|
):
|
|
"""
|
|
Initialize Redis connection and services.
|
|
|
|
Args:
|
|
host (str): Redis server hostname or IP address
|
|
port (int): Redis server port number
|
|
username (str): Redis authentication username
|
|
password (str): Redis authentication password
|
|
logger (Logger): Logger instance for operation logging
|
|
notification_handler (NotificationHandler): Handler for system notifications
|
|
"""
|
|
SientiaMonitoring.__init__(self, logger, notification_handler, metrics_controller)
|
|
self.redis_repository = RedisRepository(
|
|
host=host,
|
|
port=port,
|
|
username=username,
|
|
password=password,
|
|
logger=logger,
|
|
notification_handler=notification_handler,
|
|
metrics_controller=metrics_controller,
|
|
)
|
|
|
|
def close(self):
|
|
"""
|
|
Close the Redis connection.
|
|
"""
|
|
self.redis_repository.close()
|
|
SientiaMonitoring.shutdown(self)
|
|
|
|
@activity.defn(name='get_last_data_timestamp')
|
|
async def get_last_data_timestamp(self, input_data: dict[str, Any]) -> str | None:
|
|
"""
|
|
Retrieve the last processed data timestamp from Redis.
|
|
|
|
This activity retrieves the timestamp of the last successfully processed
|
|
data point for a specific workflow and schedule combination. It's used
|
|
for incremental data processing to avoid reprocessing the same data.
|
|
|
|
Args:
|
|
input_data (dict[str, Any]): Activity input parameters.
|
|
Required fields:
|
|
- metadata (dict[str, Any]): Workflow execution metadata
|
|
- workflow_name (str): Name of the workflow
|
|
- schedule_name (str): Name of the data collection schedule
|
|
|
|
Returns:
|
|
str | None: Last processed timestamp string, or None if no previous data exists
|
|
|
|
Raises:
|
|
Exception: If Redis operation fails
|
|
"""
|
|
metadata = input_data['metadata']
|
|
key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}'
|
|
|
|
self.info(f'Getting last data timestamp for {key}', metadata=metadata)
|
|
|
|
try:
|
|
data_hold = await self.redis_repository.get(key, metadata=metadata)
|
|
except Exception as e:
|
|
await self.send_notification_async(
|
|
metadata=metadata,
|
|
notification_id='REDIS_GET_ERROR',
|
|
message=f'Error getting last data timestamp: {e}',
|
|
block='get_last_data_timestamp',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback.format_exc(),
|
|
)
|
|
raise e
|
|
|
|
self.info(f'Last collected timestamp: {data_hold}', metadata=metadata)
|
|
|
|
if not data_hold:
|
|
return None
|
|
|
|
return data_hold
|
|
|
|
@activity.defn(name='put_last_data_timestamp')
|
|
async def put_last_data_timestamp(self, input_data: dict[str, Any]) -> str | None:
|
|
"""
|
|
Store the last processed data timestamp in Redis.
|
|
|
|
This activity stores the timestamp of the most recent data point that
|
|
has been successfully processed. The timestamp is used for incremental
|
|
data loading in subsequent workflow executions.
|
|
|
|
Args:
|
|
input_data (dict[str, Any]): Activity input parameters.
|
|
Required fields:
|
|
- metadata (dict[str, Any]): Workflow execution metadata
|
|
- data (dict[str, Any]): Processed data to extract timestamp from
|
|
- workflow_name (str): Name of the workflow
|
|
- schedule_name (str): Name of the data collection schedule
|
|
|
|
Returns:
|
|
str | None: The timestamp that was stored, or None if no data was processed
|
|
|
|
Raises:
|
|
Exception: If Redis operation fails
|
|
"""
|
|
metadata = input_data['metadata']
|
|
key = f'last_data_timestamp:{input_data["workflow_name"]}:{input_data["schedule_name"]}'
|
|
|
|
self.info(f'Putting last data timestamp for {key}', metadata=metadata)
|
|
|
|
data = DataFrame(input_data['data'])
|
|
|
|
if data.empty:
|
|
self.warning('No data to insert', metadata=metadata)
|
|
return None
|
|
|
|
last_data_timestamp = data['inserted_at'].max()
|
|
|
|
self.info(f'Last collected timestamp to insert: {last_data_timestamp}', metadata=metadata)
|
|
|
|
try:
|
|
await self.redis_repository.set(
|
|
key, last_data_timestamp, ttl=60 * 60 * 5, metadata=metadata
|
|
)
|
|
except Exception as e:
|
|
await self.send_notification_async(
|
|
metadata=metadata,
|
|
notification_id='REDIS_SET_ERROR',
|
|
message=f'Error setting last data timestamp: {e}',
|
|
block='put_last_data_timestamp',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback.format_exc(),
|
|
)
|
|
raise e
|
|
|
|
return last_data_timestamp
|
|
|
|
@activity.defn(name='group_and_hold_data')
|
|
async def group_and_hold_data(self, input_data: dict[str, Any]) -> dict[Hashable, Any]:
|
|
"""
|
|
Group data by tags and store temporarily in Redis with TTL.
|
|
|
|
This activity organizes processed data by tag names and stores it in Redis
|
|
with a configurable retention period. The data is grouped to enable
|
|
efficient batch processing and export operations.
|
|
|
|
Args:
|
|
input_data (dict[str, Any]): Activity input parameters.
|
|
Required fields:
|
|
- metadata (dict[str, Any]): Workflow execution metadata
|
|
- schedule_name (str): Name of the data collection schedule
|
|
- workflow_name (str): Name of the workflow
|
|
- data (dict[str, Any]): Data to group and store
|
|
- model_id (str): Unique model identifier
|
|
- model_tags (dict[str, Any]): Tag configuration
|
|
- retention_time (int): Data retention period in seconds
|
|
|
|
Returns:
|
|
dict[str, Any]: Grouped data organized by tag names
|
|
|
|
Raises:
|
|
Exception: If Redis operation fails
|
|
"""
|
|
metadata = input_data['metadata']
|
|
|
|
self.debug('Grouping and holding data...', metadata=metadata)
|
|
data = DataFrame(input_data['data'])
|
|
model_tags = input_data['model_tags']
|
|
retention_time = input_data['retention_time']
|
|
fill_missing_tags = input_data['fill_missing_tags']
|
|
|
|
key = f'held_data_{input_data["workflow_name"]}_{input_data["schedule_name"]}'
|
|
|
|
self.info(f'Getting held data for {key}', metadata=metadata)
|
|
|
|
try:
|
|
data_hold = await self.redis_repository.get(key, metadata=metadata)
|
|
except Exception as e:
|
|
await self.send_notification_async(
|
|
metadata=metadata,
|
|
notification_id='REDIS_GET_ERROR',
|
|
message=f'Error getting held data: {e}',
|
|
block='group_and_hold_data',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback.format_exc(),
|
|
)
|
|
raise e
|
|
|
|
if not data_hold:
|
|
data_hold = {}
|
|
if data.empty:
|
|
self.warning('No data to export', metadata=metadata)
|
|
return data_hold
|
|
|
|
self.info(f'Grouping and holding data for {len(data)} rows')
|
|
|
|
try:
|
|
# Remove possibly removed tags
|
|
tags = list(model_tags.keys())
|
|
tags.append('timestamp')
|
|
self.debug(f'Tags to keep: {tags}', metadata=metadata)
|
|
data_hold = {tag: content for tag, content in data_hold.items() if tag in tags}
|
|
self.debug(f'Data hold after removing removed tags: {data_hold}', metadata=metadata)
|
|
|
|
for _, row in data.iterrows():
|
|
value = row['value']
|
|
|
|
data_hold[row['name']] = value
|
|
|
|
if fill_missing_tags:
|
|
self.debug('Filling missing tags in data package', metadata=metadata)
|
|
missing_tags = [tag for tag in tags if tag not in list(data_hold.keys())]
|
|
|
|
for tag in missing_tags:
|
|
data_hold[tag] = None
|
|
|
|
data_hold['timestamp'] = (
|
|
data['timestamp'].max() if not data.empty else data_hold['timestamp']
|
|
)
|
|
|
|
await self.redis_repository.set(key, data_hold, ttl=retention_time, metadata=metadata)
|
|
|
|
data_hold_df = DataFrame(data_hold, index=[0])
|
|
data_hold_melted = data_hold_df.melt(
|
|
id_vars='timestamp', var_name='variable', value_name='value'
|
|
)
|
|
data_hold_melted['model_id'] = input_data['model_id']
|
|
|
|
data_hold_melted.reset_index(drop=True, inplace=True)
|
|
except Exception as e:
|
|
await self.send_notification_async(
|
|
metadata=metadata,
|
|
notification_id='REDIS_SET_ERROR',
|
|
message=f'Error setting held data: {e}',
|
|
block='group_and_hold_data',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback.format_exc(),
|
|
)
|
|
raise e
|
|
|
|
self.info(f'Data held and melted has {len(data_hold_melted)} rows')
|
|
|
|
self.debug(f'Data held and melted:\n {data_hold_melted.to_string()}', metadata=metadata)
|
|
|
|
return data_hold_melted.to_dict()
|
|
|
|
@activity.defn(name='store_data_package')
|
|
async def store_data_package(self, input_data: dict[str, Any]):
|
|
"""
|
|
Stores the data package in redis. It's a debug feature and must be toggled on.
|
|
input_data:
|
|
metadata: The metadata of the workflow.
|
|
workflow_name: The name of the workflow.
|
|
schedule_name: The name of the schedule.
|
|
held_data: The final scouter output.
|
|
data: The data used to collect the data.
|
|
"""
|
|
metadata = input_data['metadata']
|
|
key = f'data_package_{input_data["workflow_name"]}_{input_data["schedule_name"]}_{now().strftime(DATETIME_FORMAT)}'
|
|
|
|
data = DataFrame(input_data['data'])
|
|
held_data = DataFrame(input_data['held_data'])
|
|
|
|
cache = {'data': data.to_dict(), 'held_data': held_data.to_dict()}
|
|
|
|
try:
|
|
await self.redis_repository.set(key, cache, ttl=120, metadata=metadata)
|
|
except Exception as e:
|
|
await self.send_notification_async(
|
|
metadata=metadata,
|
|
notification_id='REDIS_SET_ERROR',
|
|
message=f'Error setting data package: {e}',
|
|
block='store_data_package',
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=traceback.format_exc(),
|
|
)
|
|
raise e
|