From 3bbc2c49935c86a51fb2e32369aa4ebf99ccac3e Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 19 Aug 2025 16:29:30 -0300 Subject: [PATCH 1/3] SIENTIAPDE-1199 Update dependencies and improve observability by changing logger imports. Bump sientia-dataops-library version to 1.4.1 and update image tag to 0.4.2 in values.yaml. --- requirements.txt | 2 +- scouter/activities/activities.py | 2 +- scouter/activities/faker.py | 2 +- scouter/activities/gates.py | 2 +- scouter/activities/mongodb.py | 2 +- scouter/activities/redis.py | 2 +- scouter/worker/worker.py | 2 +- scouter/workflow/fake_data.py | 2 +- scouter/workflow/scouter.py | 2 +- scouter/workflow/sub_workflows/core_scouter.py | 2 +- values.yaml | 4 ++-- 11 files changed, 12 insertions(+), 12 deletions(-) diff --git a/requirements.txt b/requirements.txt index 1b3009f..91864f7 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,7 +5,7 @@ asyncua redis aiokafka pymongo -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.7 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.1 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index 9ada6d3..875a739 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -3,7 +3,7 @@ from temporalio import workflow with workflow.unsafe.imports_passed_through(): from sientia_do.temporal.activities.postgres import Postgres from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler - from sientia_do.temporal.utils.logger import Logger + from sientia_do.observability.logger import Logger from scouter.activities.redis import Redis from scouter.activities.gates import Gates from scouter.activities.mongodb import MongoDB diff --git a/scouter/activities/faker.py b/scouter/activities/faker.py index f8e5253..92bb953 100644 --- a/scouter/activities/faker.py +++ b/scouter/activities/faker.py @@ -7,7 +7,7 @@ from temporalio import activity from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.temporal.activities.base import BaseActivity -from sientia_do.temporal.utils.logger import Logger +from sientia_do.observability.logger import Logger class Faker(BaseActivity): diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index c7840f4..bc6f471 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -6,7 +6,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.activities.base import BaseActivity from sientia_do.notifications.handlers import NotificationHandler - from sientia_do.temporal.utils.logger import Logger + from sientia_do.observability.logger import Logger from scouter.utils.quality.filters import null_values_filter, out_of_bounds_filter from typing import Any import traceback diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index 3a80436..206df33 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -9,7 +9,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.activities.base import BaseActivity - from sientia_do.temporal.utils.logger import Logger + from sientia_do.observability.logger import Logger def clear_mongo_id(docs: list) -> list: diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index f72059d..c3c504f 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -6,7 +6,7 @@ with workflow.unsafe.imports_passed_through(): 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 sientia_do.observability.logger import Logger from typing import Any from pandas import DataFrame from datetime import datetime diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 6aa61ea..cb1013d 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -6,7 +6,7 @@ with workflow.unsafe.imports_passed_through(): import sys import os from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler - from sientia_do.temporal.utils.logger import get_logger + from sientia_do.observability.logger import get_logger from scouter.activities.activities import Activities from scouter.workflow.scouter import Scouter from scouter.workflow.sub_workflows.core_scouter import CoreScouter diff --git a/scouter/workflow/fake_data.py b/scouter/workflow/fake_data.py index 63c724f..a71c900 100644 --- a/scouter/workflow/fake_data.py +++ b/scouter/workflow/fake_data.py @@ -4,7 +4,7 @@ with workflow.unsafe.imports_passed_through(): from scouter.activities.faker import Faker from datetime import timedelta from typing import Dict, Any - from sientia_do.temporal.utils.policies import retry_policy + from sientia_do.temporal.policies import retry_policy @workflow.defn(name="fake_data") diff --git a/scouter/workflow/scouter.py b/scouter/workflow/scouter.py index ffcb625..b1c4938 100644 --- a/scouter/workflow/scouter.py +++ b/scouter/workflow/scouter.py @@ -4,7 +4,7 @@ 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 + from sientia_do.temporal.policies import retry_policy @workflow.defn(name="scouter") diff --git a/scouter/workflow/sub_workflows/core_scouter.py b/scouter/workflow/sub_workflows/core_scouter.py index d686d7f..d6c83f4 100644 --- a/scouter/workflow/sub_workflows/core_scouter.py +++ b/scouter/workflow/sub_workflows/core_scouter.py @@ -4,7 +4,7 @@ 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 + from sientia_do.temporal.policies import retry_policy @workflow.defn(name="core_scouter") diff --git a/values.yaml b/values.yaml index b54199c..5d413f0 100644 --- a/values.yaml +++ b/values.yaml @@ -11,7 +11,7 @@ image: # This sets the pull policy for images. pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. - tag: "0.4.0" + tag: "0.4.2" # This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ imagePullSecrets: @@ -150,7 +150,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git" - name: GITHUB_BRANCH - value: "SIENTIAPDE-1169-pensar-e-projetar-testes-de-breakdown-e-performance" + value: "SIENTIAPDE-1199-revisar-e-testar-observabilidade" - name: PYTHON_APP value: "scouter.worker.worker" From 2a76f2e3c6943205822fe70d9eb55585371d3e44 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 20 Aug 2025 12:39:43 -0300 Subject: [PATCH 2/3] 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, From 799aa5af3e17e3f5052ba05a7d3b33879ca176a6 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 21 Aug 2025 11:35:26 -0300 Subject: [PATCH 3/3] SIENTIAPDE-1199 Refactor logging in worker.py to use custom_info and custom_error methods, enhancing log metadata with pod and workflow details for improved observability. --- scouter/worker/worker.py | 30 ++++++++++++++++++++---------- 1 file changed, 20 insertions(+), 10 deletions(-) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index 9bdfcc2..f636d59 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -30,12 +30,20 @@ async def main(): logger = get_logger(__name__) - logger.info(f"Starting Worker with pod_id: {POD_ID}") + metadata = { + 'pod_id': POD_ID, + 'model_name': '-', + 'model_id': '-', + 'workflow_name': '-', + 'schedule_name': '-', + } - logger.info("Starting prometheus client...") + logger.custom_info(f"Starting Worker with pod_id: {POD_ID}", metadata) + + logger.custom_info("Starting prometheus client...", metadata) start_prometheus_server() - logger.info('Starting Notification Handler...') + logger.custom_info('Starting Notification Handler...', metadata) mongo_config = build_mongodb_config() notification_handler = NotificationHandler( @@ -45,7 +53,7 @@ async def main(): project_name=os.getenv('PROJECT_NAME', 'scouter') ) - logger.info('Starting Activities...') + logger.custom_info('Starting Activities...', metadata) activities = Activities( logger=logger, @@ -55,7 +63,7 @@ async def main(): mongodb_config=build_mongodb_config() ) - logger.info('Starting Faker Activities...') + logger.custom_info('Starting Faker Activities...', metadata) faker_activities = Faker( logger=logger, @@ -64,7 +72,8 @@ async def main(): 'KAFKA_BOOTSTRAP_SERVERS', 'localhost:9092') ) - logger.info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...') + logger.custom_info( + f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata) new_runtime = Runtime( telemetry=TelemetryConfig( @@ -73,7 +82,7 @@ async def main(): ) ) - logger.info('Starting Temporal Client...') + logger.custom_info('Starting Temporal Client...', metadata) temporal_client = await client.Client.connect( target_host=host, @@ -81,7 +90,7 @@ async def main(): runtime=new_runtime ) - logger.info('Starting Workers...') + logger.custom_info('Starting Workers...', metadata) workers = [ Worker( @@ -120,13 +129,14 @@ async def main(): for w in workers: handlers.append(w.run()) - logger.info('Workers started successfully') + logger.custom_info('Workers started successfully', metadata) try: await asyncio.gather(*handlers) except BaseException as e: # NOSONAR - logger.error("An unhandled exception occurred: %s", e, exc_info=True) + logger.custom_error("An unhandled exception occurred: %s", + e, exc_info=True, metadata=metadata) finally: if notification_handler: notification_handler.shutdown()