Update image tag to 0.4.0 in values.yaml, modify GITHUB_BRANCH for SIENTIAPDE-1169, and enhance Redis activity to handle model tags for improved data retention.
227 lines
8.0 KiB
Python
227 lines
8.0 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.temporal.utils.logger import Logger
|
|
from typing import Any
|
|
from pandas import DataFrame
|
|
from datetime import datetime
|
|
from scouter import metrics
|
|
|
|
|
|
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']}"
|
|
|
|
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.debug(
|
|
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']}"
|
|
|
|
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.debug(
|
|
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']}"
|
|
|
|
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
|
|
|
|
try:
|
|
|
|
# Remove possibly removed tags
|
|
tags = list(model_tags.keys())
|
|
data_hold = {tag: content for tag,
|
|
content in data_hold.items() if tag in tags}
|
|
|
|
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 \
|
|
datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
|
|
self.set(key, data_hold, ttl=retention_time)
|
|
|
|
# Register metrics
|
|
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.debug(
|
|
f"Data grouped and held successfully:\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']}_{datetime.now().strftime('%Y-%m-%d_%H-%M-%S')}"
|
|
|
|
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
|