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.
This commit is contained in:
vitor-aignosi
2025-06-09 09:43:06 -03:00
parent deafb335a1
commit 7dbb9a29ea
16 changed files with 21 additions and 186 deletions

6
.env
View File

@@ -6,13 +6,13 @@ POSTGRES_DB=sientia
POSTGRES_MIN_CONNECTIONS=5 POSTGRES_MIN_CONNECTIONS=5
POSTGRES_MAX_CONNECTIONS=20 POSTGRES_MAX_CONNECTIONS=20
KAFKA_BOOTSTRAP_SERVERS=kafka:29092 KAFKA_BOOTSTRAP_SERVERS=localhost:9092
KAFKA_POLLING_TIME=1000 KAFKA_POLLING_TIME=1000
REDIS_HOST=redis REDIS_HOST=localhost
REDIS_PORT=6379 REDIS_PORT=6379
TEMPORAL_HOST=host.docker.internal:7233 TEMPORAL_HOST=localhost:7233
TEMPORAL_NAMESPACE=default TEMPORAL_NAMESPACE=default
LOG_LEVEL=INFO LOG_LEVEL=INFO

View File

@@ -3,5 +3,5 @@ psycopg2-binary
sqlalchemy sqlalchemy
asyncua asyncua
redis 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 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git

View File

@@ -1,12 +1,12 @@
from temporalio import activity, workflow from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through(): 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.redis import Redis
from scouter.activities.kafka import Kafka from scouter.activities.kafka import Kafka
from scouter.activities.gates import Gates from scouter.activities.gates import Gates
from logging import Logger from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler
from typing import Any from typing import Any

View File

@@ -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']

View File

@@ -7,7 +7,7 @@ from kafka import KafkaProducer
from temporalio import activity from temporalio import activity
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from scouter.activities.base import BaseActivity from sientia_do.temporal.activities.base import BaseActivity
class Faker(BaseActivity): class Faker(BaseActivity):

View File

@@ -2,7 +2,7 @@ from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
from sientia_do.notifications.models import NotificationLevel 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 scouter.utils.quality.filters import null_values_filter, out_of_bounds_filter
from typing import Any from typing import Any
import traceback import traceback

View File

@@ -3,7 +3,7 @@ from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
from logging import Logger from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler 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 typing import Any
from kafka import KafkaConsumer from kafka import KafkaConsumer
from pandas import DataFrame from pandas import DataFrame

View File

@@ -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()

View File

@@ -3,39 +3,19 @@ from temporalio import workflow, activity
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
from logging import Logger from logging import Logger
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from scouter.activities.base import BaseActivity from sientia_do.temporal.activities.redis_base import Redis as RedisBase
import redis
import json
from typing import Any from typing import Any
from pandas import DataFrame from pandas import DataFrame
from datetime import datetime from datetime import datetime
class Redis(BaseActivity): class Redis(RedisBase):
def __init__(self, host: str, port: int, def __init__(self, host: str, port: int,
username: str, password: str, username: str, password: str,
logger: Logger, notification_handler: NotificationHandler): logger: Logger, notification_handler: NotificationHandler):
self.host = host
self.port = port
self.username = username
self.password = password
self.redis_client = redis.Redis( RedisBase.__init__(self, host, port, username,
host=self.host, password, logger, notification_handler)
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)
@activity.defn(name="group_and_hold_data") @activity.defn(name="group_and_hold_data")
async def group_and_hold_data(self, input_data: dict[str, Any]): async def group_and_hold_data(self, input_data: dict[str, Any]):

View File

@@ -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

View File

@@ -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
)

View File

@@ -4,13 +4,13 @@ from temporalio.worker import Worker
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
import os import os
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.temporal.utils.logger import get_logger
from scouter.activities.activities import Activities from scouter.activities.activities import Activities
from scouter.workflow.scouter import Scouter from scouter.workflow.scouter import Scouter
from scouter.workflow.sub_workflows.core_scouter import CoreScouter from scouter.workflow.sub_workflows.core_scouter import CoreScouter
from scouter.workflow.fake_data import FakeData from scouter.workflow.fake_data import FakeData
from scouter.activities.faker import Faker from scouter.activities.faker import Faker
import asyncio import asyncio
from scouter.utils.logger import get_logger
from scouter.utils.connectors_config import ( from scouter.utils.connectors_config import (
build_postgres_config, build_postgres_config,
build_kafka_config, build_kafka_config,
@@ -30,11 +30,7 @@ async def main():
notification_handler = NotificationHandler( notification_handler = NotificationHandler(
servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'http://localhost:9092'), servers=os.getenv('KAFKA_BOOTSTRAP_SERVERS', 'http://localhost:9092'),
logger=logger, logger=logger,
project_name=os.getenv('PROJECT_NAME', 'scouter'), project_name=os.getenv('PROJECT_NAME', 'scouter')
pipeline_name='-',
trigger_name='-',
model_name='-',
model='-'
) )
logger.info('Starting Activities...') logger.info('Starting Activities...')

View File

@@ -4,7 +4,7 @@ with workflow.unsafe.imports_passed_through():
from scouter.activities.faker import Faker from scouter.activities.faker import Faker
from datetime import timedelta from datetime import timedelta
from typing import Dict, Any 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") @workflow.defn(name="fake_data")

View File

@@ -4,7 +4,7 @@ with workflow.unsafe.imports_passed_through():
from scouter.activities.activities import Activities from scouter.activities.activities import Activities
from typing import Any from typing import Any
from datetime import timedelta from datetime import timedelta
from scouter.utils.policies import retry_policy from sientia_do.temporal.utils.policies import retry_policy
@workflow.defn(name="scouter") @workflow.defn(name="scouter")

View File

@@ -4,7 +4,7 @@ with workflow.unsafe.imports_passed_through():
from scouter.activities.activities import Activities from scouter.activities.activities import Activities
from typing import Any from typing import Any
from datetime import timedelta 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") @workflow.defn(name="core_scouter")

View File

@@ -123,7 +123,7 @@ env:
- name: GITHUB_REPO_URL - name: GITHUB_REPO_URL
value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git" value: "git@github.com:Aignosi/sientia-dataops-scouter_temporal.git"
- name: GITHUB_BRANCH - name: GITHUB_BRANCH
value: "SIENTIAPDE-1005-implementar-os-workflows-mapeados-utilizando-as-workers-e-activities-apropriadas" value: "main"
- name: PYTHON_APP - name: PYTHON_APP
value: "scouter.worker.worker" value: "scouter.worker.worker"
@@ -164,7 +164,7 @@ env:
key: redis-password key: redis-password
- name: LOG_LEVEL - name: LOG_LEVEL
value: "INFO" value: "DEBUG"
- name: PROJECT_NAME - name: PROJECT_NAME
value: "sientia-scouter" value: "sientia-scouter"
@@ -180,6 +180,7 @@ ssh:
knownHostsPath: /mnt/known_hosts 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 # 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 # 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 \ # kubectl create secret generic git-ssh-key-sientia-scouter-worker \