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,