SIENTIAPDE-1350: Refactor: Improve logging, configuration, and resource management. Includes .env updates, client initialization logging, and resource closing.
This commit is contained in:
29
.env.example
29
.env.example
@@ -19,6 +19,7 @@ TEMPORAL_HOST=temporal-frontend.temporal.svc.cluster.local:7233
|
|||||||
TEMPORAL_NAMESPACE=model-manager
|
TEMPORAL_NAMESPACE=model-manager
|
||||||
TRAIN_TASK_QUEUE=train_model-queue
|
TRAIN_TASK_QUEUE=train_model-queue
|
||||||
CLEANUP_TASK_QUEUE=cleanup-queue
|
CLEANUP_TASK_QUEUE=cleanup-queue
|
||||||
|
TEMPORAL_USE_TLS=false
|
||||||
|
|
||||||
MONGODB_USERNAME=mongo_user
|
MONGODB_USERNAME=mongo_user
|
||||||
MONGODB_PASSWORD=mongo_db_password
|
MONGODB_PASSWORD=mongo_db_password
|
||||||
@@ -36,24 +37,16 @@ MINIO_RETRY_MODE=adaptive
|
|||||||
MINIO_CONNECT_TIMEOUT=10
|
MINIO_CONNECT_TIMEOUT=10
|
||||||
MINIO_READ_TIMEOUT=60
|
MINIO_READ_TIMEOUT=60
|
||||||
|
|
||||||
# Workflow Activity Timeouts (in seconds)
|
TIMEOUT_VALIDATE_PARAMS=30
|
||||||
# These timeouts are designed to handle large files (up to 200MB)
|
TIMEOUT_TRAIN_MODEL=2700
|
||||||
TIMEOUT_VALIDATE_PARAMS=30 # Parameter validation (fast operation)
|
TIMEOUT_DELETE_FILE=120
|
||||||
TIMEOUT_TRAIN_MODEL=2700 # Model training (30 min for large datasets)
|
TIMEOUT_UPDATE_DATABASE=30
|
||||||
TIMEOUT_DELETE_FILE=120 # Delete file from MinIO (1 min)
|
|
||||||
TIMEOUT_UPDATE_DATABASE=30 # Database update operations (30 sec)
|
|
||||||
|
|
||||||
# Cleanup Configuration
|
CLEANUP_RETENTION_HOURS=24
|
||||||
# Cleanup retention period in hours (files older than this will be deleted)
|
CLEANUP_DRY_RUN=false
|
||||||
CLEANUP_RETENTION_HOURS=24 # Default: 24 hours
|
TIMEOUT_CLEANUP_MINIO=300
|
||||||
# Enable dry run mode to test without actually deleting files
|
TIMEOUT_CLEANUP_LOCAL=120
|
||||||
CLEANUP_DRY_RUN=false # Set to true for testing without deletion
|
MAX_KEYS_CLEANUP=1000
|
||||||
# Cleanup operation timeouts
|
DEFAULT_CLEANUP_BUCKET=model-training
|
||||||
TIMEOUT_CLEANUP_MINIO=300 # MinIO cleanup timeout (5 minutes)
|
|
||||||
TIMEOUT_CLEANUP_LOCAL=120 # Local directory cleanup timeout (2 minutes)
|
|
||||||
# MinIO list operation page size for cleanup
|
|
||||||
MAX_KEYS_CLEANUP=1000 # Maximum keys per page when listing objects
|
|
||||||
# Default bucket for cleanup operations
|
|
||||||
DEFAULT_CLEANUP_BUCKET=model-training # Default bucket to clean
|
|
||||||
|
|
||||||
EXTRA_PIP_REQUIREMENTS=git+https://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@github.com/Aignosi/sientia-mlops-library.git
|
EXTRA_PIP_REQUIREMENTS=git+https://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@github.com/Aignosi/sientia-mlops-library.git
|
||||||
|
|||||||
@@ -152,3 +152,5 @@ class Activities(ExperimentTracking, Training, Cleanup):
|
|||||||
Prefer calling this method explicitly rather than relying on __del__.
|
Prefer calling this method explicitly rather than relying on __del__.
|
||||||
"""
|
"""
|
||||||
ExperimentTracking.close(self)
|
ExperimentTracking.close(self)
|
||||||
|
self.info('Postgres client closed')
|
||||||
|
self.storage_repository.close()
|
||||||
|
|||||||
@@ -90,6 +90,8 @@ class ExperimentTracking(Postgres):
|
|||||||
metrics_controller=metrics_controller,
|
metrics_controller=metrics_controller,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
self.info(f'Postgres client initialized at {host}:{port}')
|
||||||
|
|
||||||
def __del__(self):
|
def __del__(self):
|
||||||
"""
|
"""
|
||||||
Destructor to safely handle cleanup during garbage collection.
|
Destructor to safely handle cleanup during garbage collection.
|
||||||
|
|||||||
@@ -84,6 +84,7 @@ def build_mongodb_config() -> dict[str, Any]:
|
|||||||
'connection_string': connection_string,
|
'connection_string': connection_string,
|
||||||
'database_name': getenv('MONGODB_DATABASE_NAME', 'sientia'),
|
'database_name': getenv('MONGODB_DATABASE_NAME', 'sientia'),
|
||||||
'ttl_index_seconds': int(getenv('MONGODB_TTL_INDEX_HOURS', '1')) * 3600,
|
'ttl_index_seconds': int(getenv('MONGODB_TTL_INDEX_HOURS', '1')) * 3600,
|
||||||
|
'uri': uri,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -31,8 +31,8 @@ warnings.filterwarnings('ignore', category=FutureWarning, message=".*'squared' i
|
|||||||
class ModelRepository:
|
class ModelRepository:
|
||||||
def __init__(self, url, username, password, logger: Logger):
|
def __init__(self, url, username, password, logger: Logger):
|
||||||
self.model_serving = ModelServing(tracking_uri=url, username=username, password=password)
|
self.model_serving = ModelServing(tracking_uri=url, username=username, password=password)
|
||||||
|
|
||||||
self.logger = logger
|
self.logger = logger
|
||||||
|
self.logger.info(f'MLFlow client initialized at {url}')
|
||||||
|
|
||||||
def save_model(self, train_result: TrainModelResult) -> TrainModelResult:
|
def save_model(self, train_result: TrainModelResult) -> TrainModelResult:
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -86,7 +86,11 @@ class StorageRepository:
|
|||||||
use_ssl=use_ssl,
|
use_ssl=use_ssl,
|
||||||
)
|
)
|
||||||
|
|
||||||
self.logger.info(f'MinIO client initialized successfully: {endpoint_url}')
|
self.logger.info(f'MinIO client initialized at {endpoint_url}')
|
||||||
|
|
||||||
|
def close(self) -> None:
|
||||||
|
self.minio_client.close()
|
||||||
|
self.logger.info('MinIO client closed')
|
||||||
|
|
||||||
def fetch_file(self, bucket_name: str, file_name: str) -> BytesIO:
|
def fetch_file(self, bucket_name: str, file_name: str) -> BytesIO:
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -73,15 +73,14 @@ async def main():
|
|||||||
SystemExit: On graceful shutdown or error conditions
|
SystemExit: On graceful shutdown or error conditions
|
||||||
"""
|
"""
|
||||||
host = os.getenv('TEMPORAL_HOST', 'localhost:7233')
|
host = os.getenv('TEMPORAL_HOST', 'localhost:7233')
|
||||||
|
use_tls = os.getenv('TEMPORAL_USE_TLS', 'false').lower() == 'true'
|
||||||
logger = get_logger(__name__)
|
logger = get_logger(__name__)
|
||||||
|
|
||||||
metadata = {
|
metadata = {
|
||||||
'pod_id': POD_ID,
|
'pod_id': POD_ID,
|
||||||
}
|
}
|
||||||
|
|
||||||
logger.custom_info(f'Starting Worker with POD_ID: {POD_ID}', metadata)
|
|
||||||
start_prometheus_server(logger, metadata)
|
start_prometheus_server(logger, metadata)
|
||||||
logger.custom_info('Starting Notification Handler...', metadata)
|
|
||||||
mongo_config = build_mongodb_config()
|
mongo_config = build_mongodb_config()
|
||||||
|
|
||||||
notification_handler = NotificationHandler(
|
notification_handler = NotificationHandler(
|
||||||
@@ -91,7 +90,7 @@ async def main():
|
|||||||
project_name=os.getenv('PROJECT_NAME', 'model-manager'),
|
project_name=os.getenv('PROJECT_NAME', 'model-manager'),
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.custom_info('Starting Activities...', metadata)
|
logger.custom_info(f'MongoDB client initialized at {mongo_config["uri"]}', metadata)
|
||||||
|
|
||||||
activities = Activities(
|
activities = Activities(
|
||||||
postgres_config=build_postgres_config(),
|
postgres_config=build_postgres_config(),
|
||||||
@@ -101,23 +100,22 @@ async def main():
|
|||||||
notification_handler=notification_handler,
|
notification_handler=notification_handler,
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.custom_info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata)
|
|
||||||
|
|
||||||
new_runtime = Runtime(
|
new_runtime = Runtime(
|
||||||
telemetry=TelemetryConfig(
|
telemetry=TelemetryConfig(
|
||||||
metrics=PrometheusConfig(bind_address=f'0.0.0.0:{SDK_METRICS_PORT}')
|
metrics=PrometheusConfig(bind_address=f'0.0.0.0:{SDK_METRICS_PORT}')
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.custom_info(f'Starting Temporal Client at {host}...', metadata)
|
logger.custom_info(f'SDK metrics server initialized on port {SDK_METRICS_PORT}', metadata)
|
||||||
|
|
||||||
temporal_client = await client.Client.connect(
|
temporal_client = await client.Client.connect(
|
||||||
target_host=host,
|
target_host=host,
|
||||||
namespace=os.getenv('TEMPORAL_NAMESPACE', 'model-manager'),
|
namespace=os.getenv('TEMPORAL_NAMESPACE', 'model-manager'),
|
||||||
runtime=new_runtime,
|
runtime=new_runtime,
|
||||||
|
tls=use_tls,
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.custom_info('Starting Workers...', metadata)
|
logger.custom_info(f'Temporal client initialized at {host}', metadata)
|
||||||
|
|
||||||
workers = [
|
workers = [
|
||||||
Worker(
|
Worker(
|
||||||
@@ -159,7 +157,7 @@ async def main():
|
|||||||
for w in workers:
|
for w in workers:
|
||||||
handlers.append(w.run())
|
handlers.append(w.run())
|
||||||
|
|
||||||
logger.custom_info('Workers started successfully', metadata)
|
logger.custom_info('Model manager workers initialized', metadata)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# This will run the workers and wait for them to complete.
|
# This will run the workers and wait for them to complete.
|
||||||
@@ -169,6 +167,7 @@ async def main():
|
|||||||
logger.custom_error(f'An unhandled exception occurred: {e}', metadata)
|
logger.custom_error(f'An unhandled exception occurred: {e}', metadata)
|
||||||
finally:
|
finally:
|
||||||
notification_handler.shutdown()
|
notification_handler.shutdown()
|
||||||
|
logger.custom_info('MongoDB client closed', metadata)
|
||||||
await activities.shutdown()
|
await activities.shutdown()
|
||||||
# Exit with a non-zero status code to indicate failure to Kubernetes
|
# Exit with a non-zero status code to indicate failure to Kubernetes
|
||||||
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
|
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
|
||||||
@@ -195,7 +194,7 @@ def start_prometheus_server(logger: SientiaLogger, metadata: dict[str, str | Non
|
|||||||
try:
|
try:
|
||||||
port = int(os.getenv('HTTP_METRICS_PORT', 9090))
|
port = int(os.getenv('HTTP_METRICS_PORT', 9090))
|
||||||
start_http_server(port)
|
start_http_server(port)
|
||||||
logger.custom_info(f'Prometheus server started on port {port}.', metadata)
|
logger.custom_info(f'Prometheus server initialized on port {port}.', metadata)
|
||||||
metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP
|
metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP
|
||||||
except Exception as e: # noqa: BLE001
|
except Exception as e: # noqa: BLE001
|
||||||
logger.custom_critical(f'Failed to start Prometheus server: {e}', metadata)
|
logger.custom_critical(f'Failed to start Prometheus server: {e}', metadata)
|
||||||
|
|||||||
@@ -3,12 +3,6 @@
|
|||||||
# Exit on any error
|
# Exit on any error
|
||||||
set -e
|
set -e
|
||||||
|
|
||||||
#echo "Activating virtual environment..."
|
|
||||||
|
|
||||||
#conda activate ./venv
|
|
||||||
|
|
||||||
echo "Loading environment variables from .env..."
|
|
||||||
|
|
||||||
if [ -f .env ]; then
|
if [ -f .env ]; then
|
||||||
export $(cat .env | grep -v '^#' | xargs)
|
export $(cat .env | grep -v '^#' | xargs)
|
||||||
echo "Environment variables loaded from .env"
|
echo "Environment variables loaded from .env"
|
||||||
@@ -16,6 +10,4 @@ else
|
|||||||
echo "Warning: .env file not found. Continuing without environment variables."
|
echo "Warning: .env file not found. Continuing without environment variables."
|
||||||
fi
|
fi
|
||||||
|
|
||||||
echo "Starting ingestor application..."
|
|
||||||
|
|
||||||
python -m model_manager.worker.worker
|
python -m model_manager.worker.worker
|
||||||
|
|||||||
@@ -24,29 +24,35 @@ import uuid
|
|||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
|
from dotenv import load_dotenv
|
||||||
import psycopg2
|
import psycopg2
|
||||||
from psycopg2.extras import Json
|
from psycopg2.extras import Json
|
||||||
from temporalio import client
|
from temporalio import client
|
||||||
|
|
||||||
|
|
||||||
|
# Carrega variáveis de ambiente do arquivo .env na raiz do projeto
|
||||||
|
PROJECT_ROOT = Path(__file__).resolve().parent.parent
|
||||||
|
ENV_PATH = PROJECT_ROOT / '.env'
|
||||||
|
if ENV_PATH.exists():
|
||||||
|
load_dotenv(dotenv_path=ENV_PATH)
|
||||||
|
|
||||||
|
|
||||||
DOCS_PATH = Path('docs/test-model-data.csv')
|
DOCS_PATH = Path('docs/test-model-data.csv')
|
||||||
MINIO_ALIAS = os.getenv('MINIO_ALIAS', 'suse')
|
MINIO_ALIAS = 'suse'
|
||||||
MINIO_BUCKET = os.getenv('MINIO_BUCKET', 'model-training')
|
MINIO_BUCKET = 'model-training'
|
||||||
|
|
||||||
POSTGRES_CONFIG = {
|
POSTGRES_CONFIG = {
|
||||||
'host': os.getenv('POSTGRES_HOST', 'localhost'),
|
'host': os.getenv('POSTGRES_HOST'),
|
||||||
'port': os.getenv('POSTGRES_PORT', '55432'),
|
'port': os.getenv('POSTGRES_PORT'),
|
||||||
'user': os.getenv('POSTGRES_USER', 'postgres'),
|
'user': os.getenv('POSTGRES_USER'),
|
||||||
'password': os.getenv(
|
'password': os.getenv('POSTGRES_PASSWORD'),
|
||||||
'POSTGRES_PASSWORD',
|
'dbname': os.getenv('POSTGRES_DBNAME'),
|
||||||
'nFqc81y6kwmr2zuAIx43DhiOosFCVPpeEfTtTWZflkNjB2j1KtEeIANkhFR9mAX3',
|
|
||||||
),
|
|
||||||
'dbname': os.getenv('POSTGRES_DBNAME', 'sientia-core-mlops-bff'),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
TEMPORAL_HOST = os.getenv('TEMPORAL_HOST', 'localhost:37463')
|
TEMPORAL_HOST = os.getenv('TEMPORAL_HOST')
|
||||||
TEMPORAL_NAMESPACE = os.getenv('TEMPORAL_NAMESPACE', 'model-manager')
|
TEMPORAL_NAMESPACE = os.getenv('TEMPORAL_NAMESPACE')
|
||||||
TEMPORAL_TASK_QUEUE = os.getenv('TEMPORAL_TASK_QUEUE', 'train_model-queue')
|
TRAIN_TASK_QUEUE = os.getenv('TRAIN_TASK_QUEUE')
|
||||||
TEMPORAL_WORKFLOW = os.getenv('TEMPORAL_WORKFLOW', 'train_model')
|
TEMPORAL_WORKFLOW = 'train_model'
|
||||||
|
|
||||||
BASE_REQUEST_DATA = {
|
BASE_REQUEST_DATA = {
|
||||||
'experimentName': 'model-manager-test-01',
|
'experimentName': 'model-manager-test-01',
|
||||||
@@ -169,6 +175,7 @@ async def trigger_temporal_workflow(workflow_input: dict) -> str:
|
|||||||
temporal_client = await client.Client.connect(
|
temporal_client = await client.Client.connect(
|
||||||
target_host=TEMPORAL_HOST,
|
target_host=TEMPORAL_HOST,
|
||||||
namespace=TEMPORAL_NAMESPACE,
|
namespace=TEMPORAL_NAMESPACE,
|
||||||
|
tls=os.getenv('TEMPORAL_USE_TLS', False),
|
||||||
)
|
)
|
||||||
|
|
||||||
workflow_id = f'train-model-test-{uuid.uuid4()}'
|
workflow_id = f'train-model-test-{uuid.uuid4()}'
|
||||||
@@ -176,7 +183,7 @@ async def trigger_temporal_workflow(workflow_input: dict) -> str:
|
|||||||
TEMPORAL_WORKFLOW,
|
TEMPORAL_WORKFLOW,
|
||||||
workflow_input,
|
workflow_input,
|
||||||
id=workflow_id,
|
id=workflow_id,
|
||||||
task_queue=TEMPORAL_TASK_QUEUE,
|
task_queue=TRAIN_TASK_QUEUE,
|
||||||
execution_timeout=timedelta(minutes=5),
|
execution_timeout=timedelta(minutes=5),
|
||||||
run_timeout=timedelta(minutes=5),
|
run_timeout=timedelta(minutes=5),
|
||||||
task_timeout=timedelta(minutes=5),
|
task_timeout=timedelta(minutes=5),
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
- Atualizar o sientia-dataops-library para a versão 1.6.1
|
|
||||||
- Refatorar o arquivo .dockerignore para só deixar copiar os arquivos que forem necessários para a execução do container, pois ele está copiando muitos arquivos desnecessários.
|
- Refatorar o arquivo .dockerignore para só deixar copiar os arquivos que forem necessários para a execução do container, pois ele está copiando muitos arquivos desnecessários.
|
||||||
- atualizar as variáveis de ambiente no helm chart
|
- atualizar as variáveis de ambiente no helm chart
|
||||||
- Criar um gráfico no grafana para cada nova atividade.
|
- Criar um gráfico no grafana para cada nova atividade.
|
||||||
|
- Atualizar a documentação dos métodos alterados.
|
||||||
- Atualizar a documentação do projeto.
|
- Atualizar a documentação do projeto.
|
||||||
|
|
||||||
- Atualizar o .github/workflows/quality-gate.yml para usar os pipelines genéricos do github;
|
- Atualizar o .github/workflows/quality-gate.yml para usar os pipelines genéricos do github;
|
||||||
|
|||||||
Reference in New Issue
Block a user