Files
sientia-dataops-scouter_tem…/scouter/workflow/scouter.py
vitor-aignosi d08d1b1337 SIENTIAPDE-1110
Refactor Activities class to remove Kafka and Druid dependencies, simplifying initialization. Update values.yaml to set replica count to 1 for reduced resource usage. Adjust Redis activity to set TTL to None for better data retention. Remove unused Kafka and Druid activity files and their associated tests, streamlining the codebase.
2025-07-03 16:25:55 -03:00

88 lines
3.0 KiB
Python

from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from scouter.activities.activities import Activities
from typing import Any
from datetime import timedelta
from sientia_do.temporal.utils.policies import retry_policy
@workflow.defn(name="scouter")
class Scouter:
@workflow.run
async def run(self, input_data: dict[str, Any]):
"""
Scouter workflow. Loads data from kafka and sends it to the core_scouter
workflow.
Args:
input_data (dict[str, Any]): The data to process. Contains:
topic (str): The topic to load data from.
schedule_name (str): The name of the schedule.
model_name (str): The name of the model.
model_id (str): The id of the model.
trigger_laborious (bool): Whether to trigger laborious.
filters (dict[str, str]): The filters to apply.
schema (str): The schema of the table to export data to.
table_name (str): The name of the table to export data to.
retention_time (int): The retention time for data in redis in seconds.
model_tags (dict[str, Any]): The tags of the model.
And it's respective configuration.
"""
input_data['workflow_name'] = 'scouter'
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']
}
}
last_data_timestamp = await workflow.execute_local_activity_method(
Activities.get_last_data_timestamp,
{
**metadata,
'workflow_name': input_data['workflow_name'],
'schedule_name': input_data['schedule_name']
},
start_to_close_timeout=timedelta(seconds=60),
retry_policy=retry_policy
)
data = await workflow.execute_local_activity_method(
Activities.load_latest_data,
{
**metadata,
'collection_name': f"raw_{input_data['schedule_name']}",
'last_data_timestamp': last_data_timestamp
},
start_to_close_timeout=timedelta(seconds=60),
retry_policy=retry_policy
)
if data == {}:
return
await workflow.execute_activity_method(
Activities.put_last_data_timestamp,
{
**metadata,
'data': data,
'workflow_name': input_data['workflow_name'],
'schedule_name': input_data['schedule_name']
},
start_to_close_timeout=timedelta(seconds=60),
retry_policy=retry_policy
)
input_data['data'] = data
input_data['metadata'] = metadata
await workflow.execute_child_workflow(
'core_scouter',
input_data
)