SIENTIAPDE-1110

Update sientia-dataops-library version to 1.2.0 in requirements.txt; refactor logging in activities to use a unified Logger instance and include metadata in log messages across various activities.
This commit is contained in:
vitor-aignosi
2025-06-26 15:15:31 -03:00
parent 9355af1709
commit 4a90b77786
15 changed files with 238 additions and 111 deletions

View File

@@ -3,10 +3,10 @@ from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
from sientia_do.temporal.activities.postgres import Postgres
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.utils.logger import Logger
from scouter.activities.redis import Redis
from scouter.activities.kafka import Kafka
from scouter.activities.gates import Gates
from logging import Logger
from typing import Any
@@ -62,10 +62,6 @@ class Activities(Postgres, Redis, Kafka, Gates):
notification_handler=notification_handler
)
@activity.defn(name="prepare_activity")
async def prepare_activity(self, input_data: dict[str, Any]):
await super().prepare_activity(input_data)
def shutdown(self):
Postgres.close(self)
Kafka.close(self)

View File

@@ -2,12 +2,12 @@ import random
from datetime import datetime, timezone
from typing import Any
import json
from logging import Logger
from kafka import KafkaProducer
from temporalio import activity
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.activities.base import BaseActivity
from sientia_do.temporal.utils.logger import Logger
class Faker(BaseActivity):
@@ -42,6 +42,7 @@ class Faker(BaseActivity):
Defaults to random.randint(1, len(self.tags)).
"""
metadata = input_data['metadata']
topic = input_data.get('topic')
num_messages = input_data.get(
'num_messages', random.randint(1, len(self.tags))) # NOSONAR
@@ -49,8 +50,10 @@ class Faker(BaseActivity):
if not topic:
raise ValueError("Topic must be specified in input_data")
self.logger.info(
f"Generating {num_messages} messages for topic {topic}")
self.info(
f"Generating {num_messages} messages for topic {topic}",
metadata=metadata
)
for _ in range(num_messages):
# Select random tag and name
@@ -77,4 +80,4 @@ class Faker(BaseActivity):
# Ensure all messages are sent
self.producer.flush()
self.logger.info("Success")
self.info("Success", metadata=metadata)

View File

@@ -73,12 +73,16 @@ class Gates(BaseActivity):
dict[str, Any]: The aggregated data.
"""
metadata = input_data['metadata']
try:
# Convert input data to DataFrame
df = DataFrame(input_data['data'])
self.logger.debug(
f"Aggregating time series data: {df.to_string()}")
self.debug(
f"Aggregating time series data: {df.to_string()}",
metadata=metadata
)
# Initialize result dictionary
result = {}
@@ -101,16 +105,26 @@ class Gates(BaseActivity):
if aggr_value == 'continue':
continue
self.logger.debug(
f"Aggregated data: {aggr_value}")
self.logger.debug(
f"Latest timestamp: {latest_timestamp}")
self.logger.debug(
f"Groups: {group.to_string()}")
self.logger.debug(
f"group name: {name}")
self.logger.debug(
f"group tag: {tag}")
self.debug(
f"Aggregated data: {aggr_value}",
metadata=metadata
)
self.debug(
f"Latest timestamp: {latest_timestamp}",
metadata=metadata
)
self.debug(
f"Groups: {group.to_string()}",
metadata=metadata
)
self.debug(
f"group name: {name}",
metadata=metadata
)
self.debug(
f"group tag: {tag}",
metadata=metadata
)
# Store the result
result[f"{tag}_{name}"] = {
@@ -122,7 +136,10 @@ class Gates(BaseActivity):
}
result_df = DataFrame(list(result.values()))
self.logger.debug(f"Aggregated data:\n {result_df.to_string()}")
self.debug(
f"Aggregated data:\n {result_df.to_string()}",
metadata=metadata
)
return result_df.to_dict()
except Exception as e:
@@ -136,7 +153,7 @@ class Gates(BaseActivity):
attachment_content=trace
)
self.logger.error(trace)
self.error(trace, metadata=metadata)
raise e
@activity.defn(name="data_quality_gate")
@@ -160,16 +177,23 @@ class Gates(BaseActivity):
dict[str, Any]: The data validated.
"""
metadata = input_data['metadata']
filters = input_data['filters']
data = DataFrame(input_data['data'])
model_tags = input_data['model_tags']
self.logger.debug(
f"Applying quality gate to data: {data.to_string()}")
self.debug(
f"Applying quality gate to data: {data.to_string()}",
metadata=metadata
)
for filter_name, policy in filters.items():
if filter_name not in quality_gate_filters:
self.logger.warning(f"Filter {filter_name} not found")
self.warning(
f"Filter {filter_name} not found",
metadata=metadata
)
continue
try:
@@ -186,7 +210,7 @@ class Gates(BaseActivity):
attachment_content=trace
)
self.logger.error(trace)
self.error(trace, metadata=metadata)
else:
if filtered_data.empty:
@@ -206,6 +230,9 @@ class Gates(BaseActivity):
if policy == "DISCARD":
data = data[~data.index.isin(filtered_data.index)]
self.logger.debug("Data quality gate applied")
self.debug(
"Data quality gate applied",
metadata=metadata
)
return data.to_dict()

View File

@@ -4,6 +4,7 @@ with workflow.unsafe.imports_passed_through():
from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.activities.base import BaseActivity
from sientia_do.temporal.utils.logger import Logger
from typing import Any
from kafka import KafkaConsumer
from pandas import DataFrame
@@ -27,7 +28,7 @@ class Kafka(BaseActivity):
def close(self):
"""Closes the connector connection."""
self.logger.info("Closing Kafka connector...")
self.info("Closing Kafka connector...")
self.kafka_connector.close()
def __del__(self):
@@ -45,7 +46,12 @@ class Kafka(BaseActivity):
dict[str, Any]: The data loaded from the topic.
"""
self.logger.debug(f"Loading data from topic: {input_data['topic']}")
metadata = input_data['metadata']
self.debug(
f"Loading data from topic: {input_data['topic']}",
metadata=metadata
)
topic = input_data["topic"]
@@ -58,7 +64,10 @@ class Kafka(BaseActivity):
# Poll for messages
records = self.kafka_connector.poll(timeout_ms=self.polling_time)
self.logger.debug(f"Polled {len(records)} records from topic: {topic}")
self.debug(
f"Polled {len(records)} records from topic: {topic}",
metadata=metadata
)
# Process the polled records
for _topic_partition, msgs in records.items():
@@ -69,10 +78,14 @@ class Kafka(BaseActivity):
if not message_values:
return {}
self.logger.debug(
f"Loaded {len(message_values)} messages from topic: {topic}")
self.debug(
f"Loaded {len(message_values)} messages from topic: {topic}",
metadata=metadata
)
self.logger.debug(
f"Loaded data: {message_values}")
self.debug(
f"Loaded data: {message_values}",
metadata=metadata
)
return DataFrame(message_values).to_dict()

View File

@@ -4,6 +4,7 @@ with workflow.unsafe.imports_passed_through():
from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler
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
@@ -31,7 +32,11 @@ class Redis(RedisBase):
data (dict[str, Any]): The data to group and hold.
retention_time (int): The retention time for data in redis in seconds.
"""
self.logger.debug("Grouping and holding data...")
metadata = input_data['metadata']
self.debug("Grouping and holding data...",
metadata=metadata
)
data = DataFrame(input_data['data'])
retention_time = input_data['retention_time']
@@ -42,7 +47,9 @@ class Redis(RedisBase):
if not data_hold:
data_hold = {}
if data.empty:
self.logger.warning("No data to export")
self.warning("No data to export",
metadata=metadata
)
return data_hold
for _, row in data.iterrows():
@@ -62,7 +69,9 @@ class Redis(RedisBase):
data_hold_melted.reset_index(drop=True, inplace=True)
self.logger.debug(
f"Data grouped and held successfully:\n {data_hold_melted.to_string()}")
self.debug(
f"Data grouped and held successfully:\n {data_hold_melted.to_string()}",
metadata=metadata
)
return data_hold_melted.to_dict()

View File

@@ -32,21 +32,19 @@ class Scouter:
input_data['workflow_name'] = 'scouter'
await workflow.execute_local_activity_method(
Activities.prepare_activity,
{
'workflow_name': input_data['workflow_name'],
'schedule_name': input_data['schedule_name'],
metadata = {
'metadata': {
'model_id': input_data['model_id'],
'model_name': input_data['model_name'],
'model_id': input_data['model_id']
},
retry_policy=retry_policy,
start_to_close_timeout=timedelta(seconds=60)
)
'schedule_name': input_data['schedule_name'],
'workflow_name': input_data['workflow_name']
}
}
data = await workflow.execute_activity_method(
Activities.load_from_kafka,
{
**metadata,
'topic': input_data['topic']
},
retry_policy=retry_policy,

View File

@@ -19,6 +19,7 @@ class CoreScouter:
Args:
input_data (dict[str, Any]): The data to process. Contains:
metadata (dict[str, Any]): The metadata of the workflow.
workflow_name (str): The name of the workflow.
schedule_name (str): The name of the schedule.
model_name (str): The name of the model.
@@ -31,9 +32,19 @@ class CoreScouter:
retention_time (int): The retention time for data in redis in seconds.
"""
metadata = {
'metadata': {
'model_id': input_data['model_id'],
'model_name': input_data['model_name'],
'schedule_name': input_data['schedule_name'],
'workflow_name': input_data['workflow_name']
}
}
filtered_data = await workflow.execute_local_activity_method(
Activities.data_quality_gate,
{
**metadata,
'filters': input_data['filters'],
'data': input_data['data'],
'model_tags': input_data['model_tags']
@@ -45,6 +56,7 @@ class CoreScouter:
grouped_data = await workflow.execute_local_activity_method(
Activities.aggregate_data,
{
**metadata,
'data': filtered_data,
'model_tags': input_data['model_tags']
},
@@ -55,8 +67,9 @@ class CoreScouter:
held_data = await workflow.execute_local_activity_method(
Activities.group_and_hold_data,
{
'workflow_name': input_data['workflow_name'],
**metadata,
'schedule_name': input_data['schedule_name'],
'workflow_name': input_data['workflow_name'],
'data': grouped_data,
'model_id': input_data['model_id'],
'retention_time': input_data['retention_time']
@@ -71,6 +84,7 @@ class CoreScouter:
async_export = workflow.execute_activity_method(
Activities.export_data_to_postgres,
{
**metadata,
'schema': input_data['schema'],
'table_name': input_data['table_name'],
'data': held_data