From 7dbb9a29ea36bd267cef8e9a4c67d62398e73bc0 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 9 Jun 2025 09:43:06 -0300 Subject: [PATCH] SIENTIAPDE-1094 Update environment configuration and refactor activity imports - Changed Kafka, Redis, and Temporal host configurations to use localhost. - Updated the version reference for the sientia-dataops-library in requirements.txt. - Refactored import paths for activities to align with new module structure. - Removed unused base.py and postgres.py files. - Updated logger and policies imports to reflect new module locations. - Adjusted values.yaml for branch and log level settings. --- .env | 6 +- requirements.txt | 2 +- scouter/activities/activities.py | 4 +- scouter/activities/base.py | 26 ------ scouter/activities/faker.py | 2 +- scouter/activities/gates.py | 2 +- scouter/activities/kafka.py | 2 +- scouter/activities/postgres.py | 85 ------------------- scouter/activities/redis.py | 28 +----- scouter/utils/logger.py | 22 ----- scouter/utils/policies.py | 9 -- scouter/worker/worker.py | 8 +- scouter/workflow/fake_data.py | 2 +- scouter/workflow/scouter.py | 2 +- .../workflow/sub_workflows/core_scouter.py | 2 +- values.yaml | 5 +- 16 files changed, 21 insertions(+), 186 deletions(-) delete mode 100644 scouter/activities/base.py delete mode 100644 scouter/activities/postgres.py delete mode 100644 scouter/utils/logger.py delete mode 100644 scouter/utils/policies.py diff --git a/.env b/.env index 432cc6e..28055e7 100644 --- a/.env +++ b/.env @@ -6,13 +6,13 @@ POSTGRES_DB=sientia POSTGRES_MIN_CONNECTIONS=5 POSTGRES_MAX_CONNECTIONS=20 -KAFKA_BOOTSTRAP_SERVERS=kafka:29092 +KAFKA_BOOTSTRAP_SERVERS=localhost:9092 KAFKA_POLLING_TIME=1000 -REDIS_HOST=redis +REDIS_HOST=localhost REDIS_PORT=6379 -TEMPORAL_HOST=host.docker.internal:7233 +TEMPORAL_HOST=localhost:7233 TEMPORAL_NAMESPACE=default LOG_LEVEL=INFO diff --git a/requirements.txt b/requirements.txt index 90e049b..5604fc7 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,5 +3,5 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.1.14 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git diff --git a/scouter/activities/activities.py b/scouter/activities/activities.py index fa4d9c1..1fba641 100644 --- a/scouter/activities/activities.py +++ b/scouter/activities/activities.py @@ -1,12 +1,12 @@ from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): - from scouter.activities.postgres import Postgres + from sientia_do.temporal.activities.postgres import Postgres + from sientia_do.notifications.handlers import NotificationHandler from scouter.activities.redis import Redis from scouter.activities.kafka import Kafka from scouter.activities.gates import Gates from logging import Logger - from sientia_do.notifications.handlers import NotificationHandler from typing import Any diff --git a/scouter/activities/base.py b/scouter/activities/base.py deleted file mode 100644 index 3adb0e4..0000000 --- a/scouter/activities/base.py +++ /dev/null @@ -1,26 +0,0 @@ -from typing import Any -from logging import Logger -from temporalio import activity -from sientia_do.notifications.handlers import NotificationHandler - - -class BaseActivity: - def __init__(self, logger: Logger, notification_handler: NotificationHandler): - self.logger = logger - self.notification_handler = notification_handler - - @activity.defn(name="prepare_activity") - async def prepare_activity(self, input_data: dict[str, Any]): - """ - Prepare the activity for the notification handler. - - Args: - workflow_name (str): The name of the workflow. - schedule_name (str): The name of the schedule. - model_name (str): The name of the model. - model_id (str): The id of the model. - """ - self.notification_handler.base_notification.pipeline_name = input_data['workflow_name'] - self.notification_handler.base_notification.schedule_name = input_data['schedule_name'] - self.notification_handler.base_notification.model_name = input_data['model_name'] - self.notification_handler.base_notification.model_id = input_data['model_id'] diff --git a/scouter/activities/faker.py b/scouter/activities/faker.py index deb89ad..1ccc4ec 100644 --- a/scouter/activities/faker.py +++ b/scouter/activities/faker.py @@ -7,7 +7,7 @@ from kafka import KafkaProducer from temporalio import activity from sientia_do.notifications.handlers import NotificationHandler -from scouter.activities.base import BaseActivity +from sientia_do.temporal.activities.base import BaseActivity class Faker(BaseActivity): diff --git a/scouter/activities/gates.py b/scouter/activities/gates.py index c9f1395..918555a 100644 --- a/scouter/activities/gates.py +++ b/scouter/activities/gates.py @@ -2,7 +2,7 @@ from temporalio import workflow, activity with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.models import NotificationLevel - from scouter.activities.base import BaseActivity + from sientia_do.temporal.activities.base import BaseActivity from scouter.utils.quality.filters import null_values_filter, out_of_bounds_filter from typing import Any import traceback diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index 80ddd70..55cdee4 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -3,7 +3,7 @@ from temporalio import workflow, activity with workflow.unsafe.imports_passed_through(): from logging import Logger from sientia_do.notifications.handlers import NotificationHandler - from scouter.activities.base import BaseActivity + from sientia_do.temporal.activities.base import BaseActivity from typing import Any from kafka import KafkaConsumer from pandas import DataFrame diff --git a/scouter/activities/postgres.py b/scouter/activities/postgres.py deleted file mode 100644 index 53ca32e..0000000 --- a/scouter/activities/postgres.py +++ /dev/null @@ -1,85 +0,0 @@ -import traceback -from temporalio import workflow, activity - -with workflow.unsafe.imports_passed_through(): - from sqlalchemy import create_engine - from sqlalchemy.orm import sessionmaker - from sqlalchemy.pool import QueuePool - from pandas import DataFrame - from logging import Logger - from sientia_do.notifications.handlers import NotificationHandler - from sientia_do.notifications.models import NotificationLevel - from scouter.activities.base import BaseActivity - from typing import Any - - -class Postgres(BaseActivity): - def __init__(self, host: str, port: int, - user: str, password: str, dbname: str, - min_connections: int, max_connections: int, - logger: Logger, notification_handler: NotificationHandler): - self.host = host - self.port = port - self.user = user - self.password = password - self.dbname = dbname - - # Create SQLAlchemy engine with connection pooling - self.engine = create_engine( - f'postgresql://{user}:{password}@{host}:{port}/{dbname}', - poolclass=QueuePool, - pool_size=min_connections, - max_overflow=max_connections - min_connections, - pool_pre_ping=True - ) - self.session_factory = sessionmaker(bind=self.engine) - - BaseActivity.__init__(self, logger, notification_handler) - - def close(self): - self.engine.dispose() - - def __del__(self): - self.close() - - @activity.defn(name="export_data_to_postgres") - async def export_data_to_postgres(self, input_data: dict[str, Any]): - """ - Exports data to a postgres table. - - Args: - input_data (dict[str, Any]): The data to export. Contains: - schema (str): The schema of the table. - table_name (str): The name of the table. - data (DataFrame): The data to export. - """ - - self.logger.debug( - f"Exporting data to postgres: {input_data['data']}") - - schema = input_data["schema"] - table_name = input_data["table_name"] - data = DataFrame(input_data["data"]) - - with self.session_factory() as session: - try: - data.to_sql(table_name, self.engine, schema=schema, - if_exists="append", index=False) - session.commit() - - except Exception as e: - trace = traceback.format_exc() - self.notification_handler.build_and_send_notification( - notification_id="ERROR_EXPORTING_DATA_TO_POSTGRES", - message=f"Error exporting data to postgres: {e}", - block="export_data_to_postgres", - level=NotificationLevel.ERROR, - attachment_content=trace - ) - - self.logger.error(trace) - - else: - self.logger.debug("Data exported to postgres") - finally: - session.close() diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index 6af1225..13b4f54 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -3,39 +3,19 @@ from temporalio import workflow, activity with workflow.unsafe.imports_passed_through(): from logging import Logger from sientia_do.notifications.handlers import NotificationHandler - from scouter.activities.base import BaseActivity - import redis - import json + from sientia_do.temporal.activities.redis_base import Redis as RedisBase from typing import Any from pandas import DataFrame from datetime import datetime -class Redis(BaseActivity): +class Redis(RedisBase): def __init__(self, host: str, port: int, username: str, password: str, logger: Logger, notification_handler: NotificationHandler): - self.host = host - self.port = port - self.username = username - self.password = password - self.redis_client = redis.Redis( - host=self.host, - port=self.port, - decode_responses=True, - username=self.username, - password=self.password - ) - - BaseActivity.__init__(self, logger, notification_handler) - - def get(self, key: str): - history = self.redis_client.get(key) - return json.loads(history) if history else None - - def set(self, key: str, data: dict, ttl=600): - self.redis_client.set(key, json.dumps(data), ex=ttl) + RedisBase.__init__(self, host, port, username, + password, logger, notification_handler) @activity.defn(name="group_and_hold_data") async def group_and_hold_data(self, input_data: dict[str, Any]): diff --git a/scouter/utils/logger.py b/scouter/utils/logger.py deleted file mode 100644 index 42a9cfd..0000000 --- a/scouter/utils/logger.py +++ /dev/null @@ -1,22 +0,0 @@ -from os import getenv -import logging -import sys - - -def get_logger(name: str): - log_level = getenv('LOG_LEVEL', 'INFO').upper() - - logger = logging.getLogger(name) - logger.setLevel(log_level) - stream_handler = logging.StreamHandler(sys.stdout) - stream_handler.setLevel(log_level) - - stream_handler.setFormatter( - logging.Formatter( - '%(asctime)s - %(name)s - %(levelname)s - %(message)s' - ) - ) - - logger.addHandler(stream_handler) - - return logger diff --git a/scouter/utils/policies.py b/scouter/utils/policies.py deleted file mode 100644 index 9230f31..0000000 --- a/scouter/utils/policies.py +++ /dev/null @@ -1,9 +0,0 @@ -from temporalio.common import RetryPolicy -from datetime import timedelta - -retry_policy = RetryPolicy( - initial_interval=timedelta(seconds=1), - backoff_coefficient=2.0, - maximum_interval=timedelta(minutes=1), - maximum_attempts=1 -) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index bc2eaab..c49ca64 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -4,13 +4,13 @@ from temporalio.worker import Worker with workflow.unsafe.imports_passed_through(): import os from sientia_do.notifications.handlers import NotificationHandler + from sientia_do.temporal.utils.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 from scouter.workflow.fake_data import FakeData from scouter.activities.faker import Faker import asyncio - from scouter.utils.logger import get_logger from scouter.utils.connectors_config import ( build_postgres_config, build_kafka_config, @@ -30,11 +30,7 @@ async def main(): notification_handler = NotificationHandler( servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'http://localhost:9092'), logger=logger, - project_name=os.getenv('PROJECT_NAME', 'scouter'), - pipeline_name='-', - trigger_name='-', - model_name='-', - model='-' + project_name=os.getenv('PROJECT_NAME', 'scouter') ) logger.info('Starting Activities...') diff --git a/scouter/workflow/fake_data.py b/scouter/workflow/fake_data.py index 6ff347a..63c724f 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 scouter.utils.policies import retry_policy + from sientia_do.temporal.utils.policies import retry_policy @workflow.defn(name="fake_data") diff --git a/scouter/workflow/scouter.py b/scouter/workflow/scouter.py index 72596db..f3c98fa 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 scouter.utils.policies import retry_policy + from sientia_do.temporal.utils.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 f1d3dd6..c1a7d2c 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 scouter.utils.policies import retry_policy + from sientia_do.temporal.utils.policies import retry_policy @workflow.defn(name="core_scouter") diff --git a/values.yaml b/values.yaml index ca9f0bd..507eeb9 100644 --- a/values.yaml +++ b/values.yaml @@ -123,7 +123,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git" - name: GITHUB_BRANCH - value: "SIENTIAPDE-1005-implementar-os-workflows-mapeados-utilizando-as-workers-e-activities-apropriadas" + value: "main" - name: PYTHON_APP value: "scouter.worker.worker" @@ -164,7 +164,7 @@ env: key: redis-password - name: LOG_LEVEL - value: "INFO" + value: "DEBUG" - name: PROJECT_NAME value: "sientia-scouter" @@ -180,6 +180,7 @@ ssh: knownHostsPath: /mnt/known_hosts # kubectl create secret docker-registry docker-hub-secret --namespace sientia --docker-server=http://aignosi.azurecr.io --docker-username=aignosi --docker-password=5I5zpQ6sRaHqX1hD3dr+2mo647yO3FRc359/wu6gsP+ACRDRz5mp + # helm upgrade --install sientia-scouter-worker sientia/sientia-module -n sientia --create-namespace -f ./values.yaml --version 0.1.0-uat # kubectl create secret generic git-ssh-key-sientia-scouter-worker \