842 lines
30 KiB
Markdown
842 lines
30 KiB
Markdown
# Sientia DataOps Model Manager
|
|
|
|
A comprehensive AI model management platform for the complete machine learning lifecycle. Handles model training, versioning, deployment, monitoring, and governance. Streamlines MLOps workflows with centralized model registry, automated pipelines, performance tracking, and enterprise-grade compliance features.
|
|
|
|
## Features
|
|
|
|
### Core Functionality
|
|
- **Batch Prediction Processing**: High-throughput ML model inference using MLFlow models
|
|
- **Temporal Workflow Orchestration**: Robust workflow management with automatic retry policies and fault tolerance
|
|
- **Data Quality Gates**: Configurable filtering for data validation, MLFlow API responses, and custom validation rules
|
|
- **Multi-Model Support**: Flexible ML model management with retention policies and versioning
|
|
- **Real-time Data Export**: PostgreSQL persistence and OPC server integration for industrial systems
|
|
- **Comprehensive Monitoring**: Prometheus metrics and detailed logging for operational visibility
|
|
|
|
### Advanced Capabilities
|
|
- **Incremental Data Processing**: Timestamp-based data loading to avoid reprocessing
|
|
- **Configurable Data Retention**: Model retention policies with automatic cleanup
|
|
- **Notification System**: Integrated alerting and notification management via MongoDB
|
|
- **Scalable Architecture**: Kubernetes-ready deployment with horizontal scaling support
|
|
- **Model Retraining**: Automated model retraining workflows with production model updates
|
|
|
|
## Architecture
|
|
|
|
The Laborious system uses a Temporal-based workflow architecture with clear separation of concerns and robust error handling. The architecture is designed for high availability, scalability, and operational excellence in production ML environments.
|
|
|
|
|
|
### Architecture Principles
|
|
|
|
#### 1. **Separation of Concerns**
|
|
- **Worker Layer**: Manages Temporal workers, task queues, and application lifecycle
|
|
- **Workflow Layer**: Orchestrates business logic and process coordination
|
|
- **Activity Layer**: Implements specific operations and external system interactions
|
|
- **Data Layer**: Handles data persistence, caching, and external service connections
|
|
|
|
#### 2. **Fault Tolerance & Resilience**
|
|
- **Automatic Retry Policies**: Configurable retry strategies for transient failures
|
|
- **Circuit Breaker Pattern**: Prevents cascading failures in external service calls
|
|
- **Graceful Degradation**: System continues operating with reduced functionality
|
|
- **Comprehensive Error Handling**: Detailed error reporting and notification integration
|
|
|
|
#### 3. **Scalability & Performance**
|
|
- **Horizontal Scaling**: Multiple worker instances for load distribution
|
|
- **Task Queue Isolation**: Separate queues for different workflow types
|
|
- **Connection Pooling**: Optimized database and external service connections
|
|
- **Asynchronous Processing**: Non-blocking operations for improved throughput
|
|
|
|
#### 4. **Observability & Monitoring**
|
|
- **Prometheus Metrics**: Comprehensive system and business metrics
|
|
- **Structured Logging**: Consistent log format with correlation IDs
|
|
- **Health Checks**: Endpoint health monitoring and alerting
|
|
- **Performance Tracing**: Request flow tracking and bottleneck identification
|
|
|
|
### Key Components
|
|
|
|
#### **Worker (`laborious/worker/worker.py`)**
|
|
- **Purpose**: Main application orchestrator managing Temporal workers and task queues
|
|
- **Responsibilities**:
|
|
- Temporal client initialization and connection management
|
|
- Worker lifecycle management and graceful shutdown
|
|
- Task queue configuration and load balancing
|
|
- Prometheus metrics server initialization
|
|
- Notification handler setup and configuration
|
|
- OPC server connection management
|
|
- **Key Features**:
|
|
- Automatic scaling with `PollerBehaviorAutoscaling`
|
|
- Health check endpoints for Kubernetes liveness/readiness probes
|
|
- Graceful shutdown with cleanup procedures
|
|
- Multi-instance deployment support
|
|
- Two dedicated task queues: `predictions_batch-queue` and `minimal_retrain-queue`
|
|
|
|
#### **Workflows (`laborious/workflows/`)**
|
|
- **PredictionsBatch**: Main entry point for batch prediction pipelines
|
|
- **PredictionProcess**: Core prediction pipeline with MLFlow integration
|
|
- **FormatAndExportPrediction**: Data formatting and export operations
|
|
- **MinimalRetrain**: Automated model retraining and deployment
|
|
- **Key Features**:
|
|
- Temporal workflow definitions with retry policies
|
|
- Child workflow orchestration and delegation
|
|
- Comprehensive error handling and recovery
|
|
- Configurable timeout and retry strategies
|
|
|
|
#### **Activities (`laborious/activities/`)**
|
|
- **Activities**: Main activity orchestrator combining all functionality through multiple inheritance
|
|
- **Gates**: Data quality validation and filtering mechanisms
|
|
- **MLFlow**: Model transformation and prediction operations
|
|
- **OPC**: Real-time data export to industrial OPC servers
|
|
- **Key Features**:
|
|
- Multiple inheritance pattern for unified activity interface
|
|
- Configurable filter policies and validation rules
|
|
- MLFlow model serving integration with configurable flavors
|
|
- OPC UA client with certificate-based authentication
|
|
- Comprehensive error handling and notification integration
|
|
- Support for multiple OPC servers with independent configurations
|
|
|
|
#### **Data Services (`laborious/utils/`)**
|
|
- **Connectors Config**: Environment variable-based configuration management
|
|
- **Repository**: Data access layer for MLFlow and OPC operations
|
|
- `model_repository.py`: MLFlow model operations and retraining
|
|
- `opc_repository.py`: OPC server communication and data writing
|
|
- **Filters**: Data quality validation and MLFlow response filtering
|
|
- `conditional_filters.py`: Input data validation filters
|
|
- `mlflow_filters.py`: MLFlow API response validation filters
|
|
- **Key Features**:
|
|
- Environment variable-based configuration with sensible defaults
|
|
- Connection pool management and optimization
|
|
- Security credential management
|
|
- Configuration validation and error handling
|
|
- Support for multiple OPC servers and MLFlow model flavors
|
|
|
|
### Data Flow Architecture
|
|
|
|
#### **1. Batch Prediction Pipeline**
|
|
```
|
|
Input Data (PostgreSQL) → Data Quality Gates → MLFlow Transform →
|
|
MLFlow Prediction → Response Validation → Export (PostgreSQL + OPC)
|
|
```
|
|
|
|
#### **2. Model Retraining Pipeline**
|
|
```
|
|
Training Data → Model Retraining → Quality Validation →
|
|
Production Update → Notification & Monitoring
|
|
```
|
|
|
|
#### **3. Real-time Export Pipeline**
|
|
```
|
|
Prediction Results → Data Formatting → OPC Server Write →
|
|
Success/Failure Metrics → Notification System
|
|
```
|
|
|
|
### Security Architecture
|
|
|
|
#### **Authentication & Authorization**
|
|
- **Certificate-based OPC Authentication**: Secure industrial communication
|
|
- **MLFlow API Authentication**: Username/password with secure transmission
|
|
- **Database Connection Security**: Encrypted connections with credential management
|
|
- **Kubernetes Secrets Integration**: Secure credential storage and access
|
|
|
|
#### **Network Security**
|
|
- **TLS/SSL Encryption**: Secure communication channels
|
|
- **Network Isolation**: Kubernetes network policies and service mesh
|
|
- **Firewall Rules**: Controlled access to external services
|
|
- **VPN Integration**: Secure remote access and management
|
|
|
|
#### **Data Security**
|
|
- **Data Encryption**: At-rest and in-transit encryption
|
|
- **Access Control**: Role-based access control (RBAC)
|
|
- **Audit Logging**: Comprehensive access and operation logging
|
|
- **Data Retention**: Configurable data lifecycle management
|
|
|
|
## 🔄 Workflows
|
|
|
|
### 1. Predictions Batch Workflow (`predictions_batch.py`)
|
|
|
|
The **PredictionsBatch** workflow is the main entry point for batch prediction pipelines. It orchestrates the complete prediction process and implements a robust data loading and processing pattern.
|
|
|
|
#### Purpose
|
|
- **Batch Prediction Orchestration**: Coordinates data loading and prediction processing
|
|
- **Data Preparation**: Loads data using custom SQL queries with configurable schemas
|
|
- **Workflow Delegation**: Delegates actual prediction processing to the PredictionProcess workflow
|
|
- **Configuration Management**: Handles model configuration, filters, and retention policies
|
|
|
|
#### Execution Flow
|
|
1. **Data Loading**: Executes custom SQL query to load data from PostgreSQL
|
|
2. **Input Preparation**: Prepares prediction input with metadata and configuration
|
|
3. **Workflow Delegation**: Spawns PredictionProcess child workflow for actual processing
|
|
4. **Error Handling**: Implements comprehensive error handling with retry policies
|
|
|
|
#### Key Features
|
|
- **Custom Query Support**: Flexible SQL-based data loading
|
|
- **Schema Configuration**: Configurable data schema definitions
|
|
- **Automatic Retry**: Implements Temporal retry policies for fault tolerance
|
|
- **Timeout Management**: 60-second timeout for all activities
|
|
- **Comprehensive Error Handling**: Detailed error reporting and notification integration
|
|
|
|
#### Input Parameters
|
|
```json
|
|
{
|
|
"schedule_name": "hourly_predictions",
|
|
"model_name": "temperature_prediction_model",
|
|
"model_id": "temp_pred_001",
|
|
"query": "SELECT * FROM sensor_data WHERE timestamp > NOW() - INTERVAL '1 hour'",
|
|
"schema": {
|
|
"timestamp": "datetime",
|
|
"temperature": "float",
|
|
"humidity": "float"
|
|
},
|
|
"table_name": "predictions",
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"model_retention": 60,
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"opc_output_config": {
|
|
"server_id": "opc_server_1",
|
|
"tags": ["prediction_output"]
|
|
}
|
|
}
|
|
```
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. load_custom_query] --> B[2. prediction_process 🔃]
|
|
|
|
A -.-> Database[(Database)]
|
|
```
|
|
|
|
### 2. Prediction Process Workflow (`prediction_process.py`)
|
|
|
|
The **PredictionProcess** workflow implements the core prediction pipeline for ML model inference. It handles data quality validation, MLFlow model interactions, and prediction processing.
|
|
|
|
#### Purpose
|
|
- **Data Quality Validation**: Applies configurable filters for data integrity
|
|
- **MLFlow Integration**: Manages model transformation and prediction requests
|
|
- **Response Validation**: Filters MLFlow API responses for quality assurance
|
|
- **Prediction Export**: Delegates prediction formatting and export operations
|
|
|
|
#### Execution Flow
|
|
1. **Timestamp Retrieval**: Gets the last processed timestamp for incremental processing
|
|
2. **Input Data Gate**: Applies configured filters for data quality validation
|
|
3. **Path Decision**: Determines processing path based on filter results
|
|
4. **MLFlow Transform**: Requests data transformation using MLFlow models
|
|
5. **Response Validation**: Filters transform responses for quality assurance
|
|
6. **MLFlow Prediction**: Executes prediction using transformed data
|
|
7. **Content Validation**: Filters prediction responses for final quality check
|
|
8. **Export Delegation**: Delegates to FormatAndExportPrediction workflow
|
|
|
|
#### Key Features
|
|
- **Configurable Quality Gates**: Multiple filter types with policy-based configuration
|
|
- **Flexible Path Handling**: Configurable decision paths (STOP, CONTINUE, REPEAT)
|
|
- **MLFlow Integration**: Comprehensive model management and inference
|
|
- **Incremental Processing**: Timestamp-based data processing optimization
|
|
- **Comprehensive Monitoring**: Detailed metrics and error reporting
|
|
|
|
#### Input Parameters
|
|
```json
|
|
{
|
|
"metadata": {
|
|
"schedule_name": "hourly_predictions",
|
|
"model_name": "temperature_prediction_model",
|
|
"model_id": "temp_pred_001",
|
|
"workflow_name": "predictions_batch"
|
|
},
|
|
"data": {...},
|
|
"schema": {...},
|
|
"table_name": "predictions",
|
|
"model_id": "temp_pred_001",
|
|
"model_name": "temperature_prediction_model",
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"},
|
|
"SPECIFIC_VARIABLES_NULL_VALUES": {
|
|
"POLICY": "STOP",
|
|
"config": {"variables": ["temperature", "humidity"]}
|
|
}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"},
|
|
"NAN_VALUES": {"POLICY": "STOP"}
|
|
},
|
|
"model_retention": 60,
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"opc_output_config": {...}
|
|
}
|
|
```
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. get_last_timestamp] --> B[2. input_gate] --> C[3. request_transform] --> D[4. mlflow_response_gate] --> E[5. mlflow_content_gate] --> F[6. request_predict] --> G[7. mlflow_response_gate] --> H[8. format_and_export_prediction🔃]
|
|
|
|
A -.-> Redis[(Redis)]
|
|
C -.-> MLFlow[MLFlow]
|
|
F -.-> MLFlow[MLFlow]
|
|
G -.-> Filters[MLFlow Filters]
|
|
```
|
|
|
|
### 3. Format and Export Prediction Workflow (`format_and_export_prediction.py`)
|
|
|
|
The **FormatAndExportPrediction** workflow handles prediction data formatting and export operations to multiple destinations.
|
|
|
|
#### Purpose
|
|
- **Data Formatting**: Formats prediction data for different output destinations
|
|
- **PostgreSQL Export**: Persists predictions to database with metrics
|
|
- **OPC Integration**: Writes predictions to OPC servers for real-time access
|
|
- **Metrics Recording**: Tracks export operations and performance metrics
|
|
|
|
#### Execution Flow
|
|
1. **Path Decision**: Determines formatting path based on configuration
|
|
2. **Data Formatting**: Formats prediction data for specific output requirements
|
|
3. **PostgreSQL Export**: Writes formatted predictions to database
|
|
4. **OPC Export**: Writes predictions to OPC servers
|
|
5. **Metrics Recording**: Records export performance and success metrics
|
|
|
|
#### Key Features
|
|
- **Flexible Formatting**: Configurable output formats for different destinations
|
|
- **Multi-Destination Export**: PostgreSQL and OPC server integration
|
|
- **Performance Monitoring**: Comprehensive metrics for export operations
|
|
- **Error Handling**: Robust error handling with notification integration
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. format_prediction/format_default_prediction] --> B[2. write_opc_data] --> C[3. export_data_to_postgres] --> D[4. write_metrics]
|
|
|
|
A -.-> Format[Data Formatting]
|
|
B -.-> OPC[OPC Servers]
|
|
C -.-> PostgreSQL[(PostgreSQL)]
|
|
D -.-> Prometheus[Prometheus]
|
|
```
|
|
|
|
### 4. Minimal Retrain Workflow (`minimal_retrain.py`)
|
|
|
|
The **MinimalRetrain** workflow handles automated model retraining and production model updates.
|
|
|
|
#### Purpose
|
|
- **Model Retraining**: Automates ML model retraining processes
|
|
- **Production Updates**: Manages production model version updates
|
|
- **Data Export**: Exports training data for model development
|
|
- **Quality Assurance**: Ensures model quality before production deployment
|
|
|
|
#### Execution Flow
|
|
1. **Data Loading**: Loads training data using custom queries
|
|
2. **Model Retraining**: Executes model retraining process
|
|
3. **Quality Validation**: Validates retrained model performance
|
|
4. **Production Update**: Updates production model if quality criteria met
|
|
5. **Data Export**: Exports training data for analysis
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. load_custom_query] --> B[2. retrain_model] --> C[3. update_production_model] --> D[4. export_data_to_postgres]
|
|
|
|
A -.-> Database[(Database)]
|
|
B -.-> MLFlow[MLFlow]
|
|
C -.-> MLFlow[MLFlow]
|
|
D -.-> PostgreSQL[(PostgreSQL)]
|
|
```
|
|
|
|
## 📋 Prerequisites
|
|
|
|
- Python 3.11+
|
|
- Temporal server/cluster
|
|
- PostgreSQL database
|
|
- MLFlow server
|
|
- OPC server(s)
|
|
- MongoDB server (for notifications)
|
|
|
|
**Note**: External dependencies must be available either through:
|
|
- Kubernetes cluster deployment
|
|
- Docker Compose setup
|
|
- Cloud-managed services
|
|
- Local installations
|
|
|
|
## 🚀 Installation
|
|
|
|
### Local Development Setup
|
|
|
|
1. **Clone the repository**
|
|
```bash
|
|
git clone <repository-url>
|
|
cd sientia-dataops-model-manager
|
|
```
|
|
|
|
2. **Create virtual environment**
|
|
```bash
|
|
conda create -p ./venv python=3.11
|
|
conda activate ./venv
|
|
```
|
|
|
|
3. **Install dependencies**
|
|
|
|
1. **Install github cli**
|
|
```bash
|
|
sudo apt update
|
|
sudo apt install gh -y
|
|
```
|
|
|
|
2. **Authenticate with github**
|
|
```bash
|
|
gh auth login
|
|
```
|
|
|
|
3. **Run the install_dependencies.sh script**
|
|
```bash
|
|
chmod +x install_dependencies.sh
|
|
./install_dependencies.sh
|
|
```
|
|
|
|
4. **Install Python dependencies**
|
|
```bash
|
|
python -m pip install --upgrade pip
|
|
pip install -r requirements.txt
|
|
```
|
|
|
|
5. **Install test libraries**
|
|
```bash
|
|
pip install pytest pytest-cov pytest-asyncio
|
|
```
|
|
|
|
4. **Create environment configuration file**
|
|
```bash
|
|
cp .env.example .env
|
|
# Edit .env with your connection details
|
|
```
|
|
|
|
5. **Configure external dependencies**
|
|
|
|
You'll need to set up port forwarding or connections to external services. For example:
|
|
|
|
```bash
|
|
# Port forwarding from Kubernetes cluster
|
|
kubectl port-forward svc/postgresql 5432:5432
|
|
kubectl port-forward svc/mlflow 5000:5000
|
|
kubectl port-forward svc/mongodb 27017:27017
|
|
|
|
# Or connect to external services
|
|
# Ensure services are accessible on localhost with appropriate ports
|
|
```
|
|
|
|
## 📦 How to Run
|
|
|
|
### Running the Laborious Application
|
|
|
|
Use the provided script to run the application locally:
|
|
|
|
```bash
|
|
# Make script executable (first time only)
|
|
chmod +x run_local.sh
|
|
|
|
# Run the application
|
|
./run_local.sh
|
|
```
|
|
|
|
The script will:
|
|
- Activate the virtual environment
|
|
- Load environment variables from `.env`
|
|
- Start the laborious worker application
|
|
|
|
### Running Tests and Coverage
|
|
|
|
Use the provided script to run tests with coverage:
|
|
|
|
```bash
|
|
# Make script executable (first time only)
|
|
chmod +x run_coverage.sh
|
|
|
|
# Run tests with coverage
|
|
./run_coverage.sh
|
|
```
|
|
|
|
The script will:
|
|
- Activate the virtual environment
|
|
- Run pytest with coverage reporting
|
|
- Generate HTML coverage report
|
|
- Open the coverage report in your browser
|
|
|
|
### Manual Test Execution
|
|
|
|
You can also run tests manually:
|
|
|
|
```bash
|
|
# Activate virtual environment
|
|
source ./venv/bin/activate
|
|
|
|
# Run all tests
|
|
pytest
|
|
|
|
# Run with coverage
|
|
pytest --cov=laborious --cov-report=html
|
|
|
|
# Run specific test categories
|
|
pytest tests/activities/
|
|
pytest tests/workflow/
|
|
```
|
|
|
|
### Manual Application Execution
|
|
|
|
For manual execution without scripts:
|
|
|
|
```bash
|
|
# Activate virtual environment
|
|
source ./venv/bin/activate
|
|
|
|
# Load environment variables (if using .env file)
|
|
if [ -f .env ]; then
|
|
export $(cat .env | grep -v '^#' | xargs)
|
|
fi
|
|
|
|
# Start the laborious worker
|
|
python -m laborious.worker.worker
|
|
```
|
|
|
|
## 🧪 Testing
|
|
|
|
### Test Structure
|
|
```
|
|
tests/
|
|
├── activities/ # Activity implementation tests
|
|
├── workflow/ # Workflow orchestration tests
|
|
├── utils/ # Utility function tests
|
|
└── integration/ # End-to-end workflow tests
|
|
```
|
|
|
|
### Test Execution
|
|
```bash
|
|
# Install test dependencies
|
|
pip install pytest pytest-cov pytest-asyncio
|
|
|
|
# Run tests with coverage
|
|
pytest --cov=laborious --cov-report=html
|
|
|
|
# Run specific test modules
|
|
pytest tests/activities/test_gates.py
|
|
pytest tests/workflow/test_predictions_batch.py
|
|
```
|
|
|
|
## 📊 Monitoring and Metrics
|
|
|
|
The Laborious system exposes comprehensive Prometheus metrics for operational visibility and performance monitoring:
|
|
|
|
### Application Health Metrics
|
|
- `app_up`: Application health status (1=healthy, 0=unhealthy)
|
|
- Labels: `pod_id`
|
|
|
|
### Prediction Operation Metrics
|
|
- `laborious_predictions_written_count`: Counter for successful prediction exports
|
|
- Labels: `pod_id`, `model_name`, `pipeline_name`
|
|
- `laborious_prediction_confidence_monitor`: Gauge for current prediction confidence levels
|
|
- Labels: `pod_id`, `model_name`, `pipeline_name`
|
|
- `laborious_prediction_response_time_monitor`: Histogram for prediction response times
|
|
- Labels: `pod_id`, `model_name`, `pipeline_name`
|
|
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
|
|
|
### OPC Export Metrics
|
|
- `laborious_prediction_opc_writing_count`: Counter for OPC server write operations
|
|
- Labels: `pod_id`, `model_name`, `pipeline_name`, `opc_server_id`
|
|
- `laborious_prediction_opc_writing_response_time_monitor`: Histogram for OPC write response times
|
|
- Labels: `pod_id`, `model_name`, `pipeline_name`, `opc_server_id`
|
|
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
|
|
|
### Data Quality Metrics
|
|
- Filter pass/fail rates through notification system
|
|
- MLFlow API response validation metrics
|
|
- Data quality gate performance tracking
|
|
|
|
## ⚙️ Configuration
|
|
|
|
### Environment Variables
|
|
|
|
| Variable | Description | Default | Required |
|
|
|----------|-------------|---------|----------|
|
|
| `TEMPORAL_HOST` | Temporal server address | `localhost:7233` | Yes |
|
|
| `TEMPORAL_NAMESPACE` | Temporal namespace | `laborious` | No |
|
|
| `POSTGRES_HOST` | PostgreSQL hostname | `localhost` | Yes |
|
|
| `POSTGRES_PORT` | PostgreSQL port | `5432` | Yes |
|
|
| `POSTGRES_USER` | PostgreSQL username | `sientia` | Yes |
|
|
| `POSTGRES_PASSWORD` | PostgreSQL password | `sientia` | Yes |
|
|
| `POSTGRES_DBNAME` | PostgreSQL database | `sientia` | Yes |
|
|
| `POSTGRES_MIN_CONNECTIONS` | Minimum PostgreSQL connections | `5` | No |
|
|
| `POSTGRES_MAX_CONNECTIONS` | Maximum PostgreSQL connections | `20` | No |
|
|
| `MLFLOW_HOST` | MLFlow server hostname | `http://localhost` | Yes |
|
|
| `MLFLOW_PORT` | MLFlow server port | `5080` | Yes |
|
|
| `MLFLOW_USERNAME` | MLFlow username | `aignosi` | Yes |
|
|
| `MLFLOW_PASSWORD` | MLFlow password | `aignosi` | Yes |
|
|
| `OPC_CONFIG` | OPC server configuration (JSON) | `{}` | No |
|
|
| `OPC_ID` | OPC server identifier | `1` | No |
|
|
| `OPC_URL` | OPC server URL | `opc.tcp://localhost:4840` | No |
|
|
| `OPC_SERVER_URI` | OPC server URI | `opc.tcp://localhost:4840` | No |
|
|
| `OPC_CERT_PATH` | OPC client certificate path | `None` | No |
|
|
| `OPC_PRIVATE_KEY_PATH` | OPC private key path | `None` | No |
|
|
| `OPC_SERVER_CERT_PATH` | OPC server certificate path | `None` | No |
|
|
| `OPC_RECONNECTION_INTERVAL` | OPC reconnection interval (ms) | `120` | No |
|
|
| `MONGODB_URL` | MongoDB connection URI | `localhost:27018` | Yes |
|
|
| `MONGODB_USERNAME` | MongoDB username | `root` | Yes |
|
|
| `MONGODB_PASSWORD` | MongoDB password | `wKZDbMNU1c` | Yes |
|
|
| `MONGODB_DATABASE_NAME` | MongoDB database name | `sientia` | Yes |
|
|
| `MONGODB_TTL_INDEX_HOURS` | MongoDB TTL index hours | `1` | No |
|
|
| `LOG_LEVEL` | Application log level | `INFO` | No |
|
|
| `PROJECT_NAME` | Project name for metrics | `laborious` | No |
|
|
| `HTTP_METRICS_PORT` | Prometheus metrics port | `9090` | No |
|
|
| `HTTP_SDK_METRICS_PORT` | Temporal SDK metrics port | `9091` | No |
|
|
| `POD_ID` | Kubernetes pod identifier | `None` | No |
|
|
|
|
|
|
|
|
|
|
### OPC Configuration
|
|
|
|
For multiple OPC servers, use the `OPC_CONFIG` environment variable:
|
|
|
|
```json
|
|
{
|
|
"opc_server_1": {
|
|
"url": "opc.tcp://server1:4840",
|
|
"name": "Server1",
|
|
"server_uri": "urn:server1:opcua",
|
|
"cert_path": "/path/to/cert.pem",
|
|
"private_key_path": "/path/to/key.pem",
|
|
"server_cert_path": "/path/to/server_cert.pem",
|
|
"reconnection_interval": 5000
|
|
},
|
|
"opc_server_2": {
|
|
"url": "opc.tcp://server2:4840",
|
|
"name": "Server2",
|
|
"server_uri": "urn:server2:opcua",
|
|
"cert_path": "/path/to/cert.pem",
|
|
"private_key_path": "/path/to/key.pem",
|
|
"server_cert_path": "/path/to/server_cert.pem",
|
|
"reconnection_interval": 5000
|
|
}
|
|
}
|
|
```
|
|
|
|
For single OPC server, use individual environment variables:
|
|
- `OPC_URL`
|
|
- `OPC_NAME`
|
|
- `OPC_SERVER_URI`
|
|
- `OPC_CERT_PATH`
|
|
- `OPC_PRIVATE_KEY_PATH`
|
|
- `OPC_SERVER_CERT_PATH`
|
|
- `OPC_RECONNECTION_INTERVAL`
|
|
|
|
### Workflow Configuration
|
|
|
|
MongoDB pipeline configuration:
|
|
|
|
#### Predictions Batch Workflow configuration sample
|
|
|
|
This is the configuration for the Predictions Batch Workflow, to be inserted into the MongoDB pipeline collection.
|
|
|
|
```json
|
|
{
|
|
"schedule_name": "laborious-orchestrated-pipeline",
|
|
"model_id": "1",
|
|
"workflow_type": "predictions_batch",
|
|
"frequency": "30s",
|
|
"max_retry_policy": 1,
|
|
"query": "select * from sientia_data.laborious_data where model_id = 1 and \"timestamp\" > NOW() - INTERVAL '5 minutes' order by \"timestamp\" desc limit 30;",
|
|
"write_tags": [
|
|
{
|
|
"server_id": "server1",
|
|
"type": "prediction",
|
|
"addr": "ns=2;i=5",
|
|
"data_type": "double"
|
|
},
|
|
{
|
|
"server_id": "server1",
|
|
"type": "confidence",
|
|
"addr": "ns=2;i=6",
|
|
"data_type": "double"
|
|
}
|
|
],
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"},
|
|
"SPECIFIC_VARIABLES_NULL_VALUES": {
|
|
"POLICY": "CONTINUE",
|
|
"config": {"variables": ["Counter"]}
|
|
}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "REPEAT"},
|
|
"NAN_VALUES": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "CONTINUE"}
|
|
},
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"active": true,
|
|
"updated_at": {
|
|
"$date": "2025-09-16T10:00:00.000Z"
|
|
},
|
|
"datetime_columns": ["timestamp", "created_at"],
|
|
"predictions_storage_policy": "lts:1"
|
|
}
|
|
```
|
|
|
|
This is the configuration created by the Orchestrator in Temporal.
|
|
|
|
```json
|
|
{
|
|
"datetime_columns":["timestamp","created_at"],
|
|
"frequency":"15m",
|
|
"input_filters":{"EMPTY_DATA":{"config":{},"policy":"STOP"}},
|
|
"max_retry_policy":1,
|
|
"mlflow_predict_filters":{"API_ERROR":{"config":{},"policy":"CONTINUE"}},
|
|
"mlflow_transform_filters":{
|
|
"API_ERROR":{"config":{},"policy":"CONTINUE"},
|
|
"EMPTY_DATA":{"config":{},"policy":"STOP"}
|
|
},
|
|
"model_config":{
|
|
"is_compressed":true,
|
|
"predict_flavor":"pyfunc",
|
|
"retention_minutes":60,
|
|
"retention_target":"artifact",
|
|
"transform_function_keyword":"transform"
|
|
},
|
|
"model_id":"352",
|
|
"model_name":"courier",
|
|
"opc_output_config":{},
|
|
"path_priority":["STOP","CONTINUE","REPEAT"],
|
|
"predictions_storage_policy":"lts:1",
|
|
"query":"select * from sientia_data.laborious_data where model_id = 352 order by \"timestamp\" desc limit 300;",
|
|
"retention_time":3600,
|
|
"schedule_name":"laborious-courier",
|
|
"schema":"sientia_data",
|
|
"table_name":"predictions",
|
|
"updated_at":"2025-09-12 19:35:01.600000+0000",
|
|
"workflow_type":"predictions_batch"
|
|
}
|
|
```
|
|
|
|
## 🔧 Development
|
|
|
|
### Project Structure
|
|
```
|
|
laborious/
|
|
├── activities/ # Temporal activity implementations
|
|
│ ├── activities.py # Main activities orchestrator
|
|
│ ├── gates.py # Data quality gates and filtering
|
|
│ ├── mlflow.py # MLFlow model operations
|
|
│ └── opc.py # OPC server operations
|
|
├── workflows/ # Temporal workflow definitions
|
|
│ ├── predictions_batch.py # Main batch prediction workflow
|
|
│ ├── minimal_retrain.py # Model retraining workflow
|
|
│ └── sub_workflows/ # Sub-workflow implementations
|
|
│ ├── prediction_process.py # Core prediction workflow
|
|
│ └── format_and_export_prediction.py # Export workflow
|
|
├── worker/ # Worker implementation
|
|
│ └── worker.py # Main worker orchestrator
|
|
├── utils/ # Utility functions
|
|
│ ├── connectors_config.py # Database configuration
|
|
│ ├── filters/ # Data quality filters
|
|
│ │ ├── conditional_filters.py # Conditional data filters
|
|
│ │ └── mlflow_filters.py # MLFlow response filters
|
|
│ └── repository/ # Data access layer
|
|
│ ├── model_repository.py # MLFlow model operations
|
|
│ └── opc_repository.py # OPC server operations
|
|
├── metrics.py # Prometheus metrics definitions
|
|
└── __init__.py
|
|
```
|
|
|
|
### Adding New Features
|
|
|
|
1. **Follow Temporal patterns** for new workflows and activities
|
|
2. **Add comprehensive docstrings** for all public methods
|
|
3. **Include Prometheus metrics** for monitoring
|
|
4. **Add unit tests** for new functionality
|
|
5. **Update this README** with new features and configuration
|
|
|
|
## 🐛 Troubleshooting
|
|
|
|
### Common Issues
|
|
|
|
1. **Temporal Connection Failures**
|
|
- Verify Temporal server is running and accessible
|
|
- Check namespace configuration and permissions
|
|
- Review server logs for connection issues
|
|
|
|
2. **MLFlow Connection Issues**
|
|
- Verify MLFlow server is running and accessible
|
|
- Check authentication credentials and permissions
|
|
- Ensure model names and versions exist
|
|
|
|
3. **Database Connection Issues**
|
|
- Verify PostgreSQL service is running
|
|
- Check connection credentials and network access
|
|
- Ensure proper connection pool configuration
|
|
|
|
4. **OPC Connection Failures**
|
|
- Verify OPC server is accessible
|
|
- Check certificate and key file paths
|
|
- Review OPC server logs for connection issues
|
|
|
|
5. **Workflow Execution Failures**
|
|
- Review activity error logs and notifications
|
|
- Check data quality filter configurations
|
|
- Verify input data format and required fields
|
|
|
|
### Debug Mode
|
|
|
|
Enable debug logging by setting the log level:
|
|
```bash
|
|
export LOG_LEVEL=DEBUG
|
|
```
|
|
|
|
## ⚡ Performance Tuning
|
|
|
|
### Key Parameters
|
|
|
|
- **Worker Concurrency**: Adjust `max_concurrent_workflow_tasks` and `max_concurrent_activities`
|
|
- **Connection Pools**: Optimize database connection pool sizes
|
|
- **Model Retention**: Configure MLFlow model retention based on requirements
|
|
- **Batch Sizes**: Adjust data processing batch sizes for optimal throughput
|
|
|
|
### Scaling Considerations
|
|
|
|
- **Horizontal Scaling**: Deploy multiple worker instances
|
|
- **Task Queue Distribution**: Use multiple task queues for different workflow types
|
|
- **Database Performance**: Optimize indexes and connection pooling
|
|
- **MLFlow Performance**: Configure appropriate model serving resources
|
|
|
|
## 🤝 Contributing
|
|
|
|
1. Fork the repository
|
|
2. Create a feature branch
|
|
3. Make your changes with comprehensive testing
|
|
4. Update documentation and docstrings
|
|
5. Submit a pull request
|
|
|
|
### Code Quality Standards
|
|
|
|
- Follow PEP 8 style guidelines
|
|
- Include comprehensive docstrings for all public methods
|
|
- Maintain test coverage above 80%
|
|
- Use type hints where appropriate
|
|
- Follow Temporal.io best practices
|
|
|
|
## 📄 License
|
|
|
|
This project is licensed under the terms specified in the LICENSE file.
|
|
|
|
## 🆘 Support
|
|
|
|
For support and questions:
|
|
- Check the troubleshooting section above
|
|
- Review the metrics and logs for error patterns
|
|
- Open an issue in the project repository
|
|
- Contact the development team
|
|
|
|
---
|
|
|
|
**Note**: The Laborious system is designed for production use in industrial ML environments. Ensure proper security configuration and network isolation for production deployments.
|