From 2a76f2e3c6943205822fe70d9eb55585371d3e44 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 20 Aug 2025 12:39:43 -0300 Subject: [PATCH] SIENTIAPDE-1199 Enhance logging in Gates, MongoDB, and Redis activities by replacing debug statements with info level logs, improving observability of data processing steps. Update worker configuration to adjust concurrency settings and enable autoscaling for task polling. --- scouter/activities/gates.py | 32 ++++++++++++++++++++++++-------- scouter/activities/mongodb.py | 2 +- scouter/activities/redis.py | 16 +++++++++++++--- scouter/worker/worker.py | 13 +++++++------ 4 files changed, 45 insertions(+), 18 deletions(-) diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index bc6f471..3815063 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -89,8 +89,8 @@ class Gates(BaseActivity): # Convert input data to DataFrame df = DataFrame(input_data['data']) - self.debug( - f"Aggregating time series data: {df.to_string()}", + self.info( + f"Aggregating time series data for {len(df)} rows", metadata=metadata ) @@ -147,10 +147,16 @@ class Gates(BaseActivity): } result_df = DataFrame(list(result.values())) - self.debug( - f"Aggregated data:\n {result_df.to_string()}", + self.info( + f"Aggregated data has {len(result_df)} rows", metadata=metadata ) + + self.debug( + f"Aggregated data: {result_df.to_string()}", + metadata=metadata + ) + return result_df.to_dict() except Exception as e: @@ -195,8 +201,8 @@ class Gates(BaseActivity): data = DataFrame(input_data['data']) model_tags = input_data['model_tags'] - self.debug( - f"Applying quality gate to data: {data.to_string()}", + self.info( + f"Applying quality gate to data to {len(data)} rows", metadata=metadata ) @@ -249,8 +255,8 @@ class Gates(BaseActivity): if policy == "DISCARD": data = data[~data.index.isin(filtered_data.index)] - self.debug( - "Data quality gate applied", + self.info( + f"Data quality gate applied, final data has {len(data)} rows", metadata=metadata ) @@ -265,8 +271,18 @@ class Gates(BaseActivity): """ metadata = input_data['metadata'] + self.info( + f"Writing metrics for {metadata['model_name']}", + metadata=metadata + ) + metrics.LABORIOUS_DATA_WRITTEN_COUNT.labels( pod_id=self.pod_id, model_name=metadata['model_name'], pipeline_name=metadata['workflow_name'] ).inc() + + self.info( + f"Metrics written for {metadata['model_name']}", + metadata=metadata + ) diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 206df33..e36c563 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -87,7 +87,7 @@ class MongoDB(BaseActivity): collection_name = input_data['collection_name'] last_data_timestamp = input_data['last_data_timestamp'] - self.debug( + self.info( f"Loading data from MongoDB: {input_data}", metadata=metadata ) diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index c3c504f..fb3f556 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -29,6 +29,8 @@ class Redis(RedisBase): 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: @@ -42,7 +44,7 @@ class Redis(RedisBase): ) raise e - self.debug( + self.info( f"Last collected timestamp: {data_hold}", metadata=metadata ) @@ -60,6 +62,8 @@ class Redis(RedisBase): 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: @@ -70,7 +74,7 @@ class Redis(RedisBase): last_data_timestamp = data['inserted_at'].max() - self.debug( + self.info( f"Last collected timestamp to insert: {last_data_timestamp}", metadata=metadata ) @@ -116,6 +120,8 @@ class Redis(RedisBase): 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: @@ -137,6 +143,8 @@ class Redis(RedisBase): ) return data_hold + self.info(f"Grouping and holding data for {len(data)} rows") + try: # Remove possibly removed tags @@ -195,8 +203,10 @@ class Redis(RedisBase): ) raise e + self.info(f"Data held and melted has {len(data_hold_melted)} rows") + self.debug( - f"Data grouped and held successfully:\n {data_hold_melted.to_string()}", + f"Data held and melted:\n {data_hold_melted.to_string()}", metadata=metadata ) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index cb1013d..9bdfcc2 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -1,5 +1,5 @@ from temporalio import workflow, client -from temporalio.worker import Worker +from temporalio.worker import Worker, PollerBehaviorAutoscaling from temporalio.runtime import Runtime, TelemetryConfig, PrometheusConfig with workflow.unsafe.imports_passed_through(): @@ -99,11 +99,12 @@ async def main(): activities.write_metrics, activities.store_data_package, ], - max_concurrent_workflow_tasks=100, - max_concurrent_activities=100, - max_concurrent_local_activities=100, - max_concurrent_workflow_task_polls=100, - max_cached_workflows=50, + max_concurrent_workflow_tasks=50, + max_concurrent_activities=50, + max_concurrent_local_activities=50, + max_cached_workflows=200, + workflow_task_poller_behavior=PollerBehaviorAutoscaling(), + activity_task_poller_behavior=PollerBehaviorAutoscaling() ), Worker( temporal_client,