Files
sientia-dataops-scouter_tem…/scouter/activities/redis.py
vitor-aignosi 0b337d00f9 SIENTIAPDE-1193
Refactor Redis activity to utilize the new 'now' function for timestamp generation, ensuring consistency with updated datetime handling. This change replaces the direct use of 'datetime.now()' with 'now()' from the temporal constants.
2025-08-22 11:24:26 -03:00

249 lines
8.7 KiB
Python

from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through():
from logging import Logger
import traceback
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
from sientia_do.notifications.models import NotificationLevel
from sientia_do.temporal.activities.redis_base import Redis as RedisBase
from sientia_do.observability.logger import Logger
from typing import Any
from pandas import DataFrame
from scouter import metrics
from sientia_do.temporal.constants import DATETIME_FORMAT, now
class Redis(RedisBase):
def __init__(self, host: str, port: int,
username: str, password: str,
logger: Logger, notification_handler: NotificationHandler):
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:
"""
Gets the last data timestamp from redis.
"""
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]):
"""
Puts the last data timestamp into redis.
"""
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]):
"""
Groups and holds data in redis. Keep a copy of the most recent
received data for a given pipeline and schedule. This activity updates
the data in redis and return the full keeped data.
Args:
input_data (dict[str, Any]): The data to group and hold.
workflow_name (str): The name of the workflow.
schedule_name (str): The name of the schedule.
data (dict[str, Any]): The data to group and hold.
retention_time (int): The retention time for data in redis in seconds.
"""
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']
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())
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))
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