Update requirements to use sientia-dataops-library version 1.5.3 and enhance Redis activity logging by including metadata in get and set operations.
312 lines
12 KiB
Python
312 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
|