Files
sientia-dataops-scouter_tem…/scouter/activities/redis.py
vitor-aignosi 4bd29ae0d7 SIENTIAPDE-1318
SIENTIAPDE-1318: Add support for filling missing tags in Redis data processing. Enhanced the Redis class to include a new parameter for handling missing tags, ensuring that absent tags are filled with None in the data package.
2025-10-20 16:24:36 -03:00

309 lines
12 KiB
Python

from collections.abc import Hashable
from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
import traceback
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.temporal.activities.redis_base import Redis as RedisBase
from sientia_do.temporal.constants import DATETIME_FORMAT, now
from scouter import metrics
class Redis(RedisBase):
"""
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,
):
"""
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
"""
RedisBase.__init__(self, host, port, username, password, logger, notification_handler)
@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}')
try:
data_hold = self.get(key)
except Exception as e:
self.send_notification(
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}')
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:
self.set(key, last_data_timestamp, ttl=60 * 60 * 5)
except Exception as e:
self.send_notification(
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}')
try:
data_hold = self.get(key)
except Exception as e:
self.send_notification(
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)
to_register_metrics = []
for _, row in data.iterrows():
value = row['value']
data_hold[row['name']] = value
to_register_metrics.append((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']
)
self.set(key, data_hold, ttl=retention_time)
# Register metrics
self.debug(f'Metrics to register: {to_register_metrics}', metadata=metadata)
for metric in to_register_metrics:
metrics.TAG_CHANGES_MONITOR.labels(
pod_id=self.pod_id,
model_name=metadata['model_name'],
pipeline_name=metadata['workflow_name'],
tag_name=metric[0],
).set(metric[1])
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:
self.send_notification(
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:
self.set(key, cache, ttl=120)
except Exception as e:
self.send_notification(
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