SIENTIAPDE-1712 Enhance logging in Gates, MLFlow, and Storage classes by integrating logger parameter for improved traceability. This update allows for better monitoring of operations and data handling across these components.
Sientia DataOps Laborious
A comprehensive, Temporal-based ML orchestration system for industrial data processing and model inference. Laborious delivers enterprise-grade batch prediction, model management, optional real-time export (OPC and PI Web API), and automated retraining with strong data quality validation and observability.
📑 Table of Contents
- Features
- Architecture
- Workflows
- Installation & Setup
- How to Run
- Code Quality & Validation
- Testing
- Monitoring and Metrics
- Configuration
- Development
- Troubleshooting
- Performance Tuning
- Contributing
- License
- Support
Features
Core Functionality
- Batch Prediction Processing: High-throughput ML inference using MLFlow models
- Temporal Workflow Orchestration: Robust workflow management with retries and fault tolerance
- Data Quality Gates: Configurable filtering for input data and MLFlow API responses
- Multi-Model Support: Flexible model management with retention and versioning
- Optional Real-time Export: PostgreSQL persistence, OPC server integration, and PI Web API integration for industrial systems
- Comprehensive Monitoring: Prometheus metrics and structured logging for observability
Advanced Capabilities
- Incremental Data Processing: Timestamp-based loading to avoid reprocessing
- Configurable Data Retention: Model retention policies with automatic cleanup
- MinIO Payload Offload: Automatic offload of large DataFrames to MinIO with retention cleanup
- Data Drift Detection: Univariate and multivariate drift monitoring against reference data
- Regression Metrics: Automated RMSE, MSE, MAE, R² calculation and export
- Notification System: Integrated alerting via MongoDB
- Scalable Architecture: Kubernetes-ready with horizontal scaling
- Model Retraining: Automated retraining workflows with production model updates
Development & Quality Assurance
- Code Quality Tools: Ruff (lint/format), mypy (types), Bandit (security)
- Automated Validation: CI quality gates and individual tool commands
- Comprehensive Testing: pytest with async support and high coverage
- Type Safety: Static type checking with mypy
- Coverage Visualization: Coverage Gutters integration
Architecture
Laborious uses a Temporal-based architecture with strong separation of concerns and defensive error handling for production ML.
Architecture Principles
1. Separation of Concerns
- Worker Layer: Temporal workers, task queues, lifecycle
- Workflow Layer: Business orchestration and coordination
- Activity Layer: External system interactions and isolated operations
- Data Layer: Persistence, caching, connectors
2. Fault Tolerance & Resilience
- Automatic Retry Policies for transient failures
- Graceful Degradation and circuit breaking for dependencies
- Detailed Error Handling with notifications
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)
- Temporal client setup, worker lifecycle, task queues
- Metrics server initialization, notification handler setup
- Graceful shutdown and autoscaling-friendly behavior
Workflows (laborious/workflows/)
predictions_batch.py: Batch prediction entry pointsub_workflows/prediction_process.py: Core prediction pipelinesub_workflows/format_and_export_prediction.py: Formatting and exportminimal_retrain.py: Automated model retraining and production updatedrift.py: Data drift detection and monitoringsimple_metrics.py: Regression metrics calculation (RMSE, MSE, MAE, R²)
Activities (laborious/activities/)
gates.py: Data quality validation, filtering, and data formatting operations- Input/response/content gates for quality validation
- Prediction and transformed data formatting
- Retrain report formatting and metrics recording
mlflow.py: Transform, predict, and model management operations- MLFlow model transformation and prediction
- Model retraining and production updates
- Reference data retrieval from MLflow Model Registry
storage.py: PostgreSQL queries and MinIO-aware data loadingload_query_with_minio_offload: SQL load with automatic MinIO offloadexport_payload_to_postgres: Resolve MinIO payloads and export to Postgrescleanup_minio_objects_expired: Retention-based MinIO object cleanupquery_to_minio: Legacy parquet upload for retraining data
model_metrics.py: Drift detection and regression metrics- Univariate and multivariate drift calculation
- Simple metrics (RMSE, MSE, MAE, R²)
opc.py: OPC UA export to industrial systems (optional)api.py: PI Web API export operations (optional)- Prediction and confidence data writing to PI Web API
- Error handling and notification integration
activities.py: Aggregates all activity interfaces (Storage, MLFlow, Gates, OPC, ModelMetrics, API)
Data Services (laborious/utils/)
connectors_config.py: Env-driven configuration buildersmodels/minio_dataframe_payload.py: MinIO-offloaded DataFrame payload modelrepository/model_repository.py: MLFlow operations and retrainingrepository/opc_repository.py: OPC communication and writesrepository/minio_manager.py: MinIO object storage operationsfilters/conditional_filters.pyandfilters/mlflow_filters.py
Data Flow Architecture
1. Batch Prediction Pipeline
Input Data (PostgreSQL) → Data Quality Gates → MLFlow Transform →
MLFlow Prediction → Response Validation → Format & Export
├─→ Predictions → PostgreSQL [+ OPC] [+ PI Web API]
└─→ Transformed Data → PostgreSQL (optional)
2. Model Retraining Pipeline
Training Data → Model Retraining → Quality Validation →
Production Update → Notification & Monitoring
3. Drift Detection Pipeline
Target Data (PostgreSQL) + Reference Data (MLFlow) →
Drift Calculation (univariate + multivariate) → PostgreSQL Export
4. Simple Metrics Pipeline
Predictions + Targets (PostgreSQL JOIN) →
Metrics Calculation (RMSE, MSE, MAE, R²) → PostgreSQL Export
Security Architecture
Authentication & Authorization
- MLFlow API Authentication: Username/password
- Database Security: Encrypted connections and credential management
- OPC Certificates (if enabled): Client/server certs
- PI Web API Authentication: Bearer token or basic authentication
- Kubernetes Secrets: Secure secret storage
Network Security
- TLS/SSL, network policies, service mesh, firewalls, VPN
Data Security
- At-rest/in-transit encryption, RBAC, audit logging, 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
- Data Loading: Executes custom SQL query to load data from PostgreSQL
- Input Preparation: Prepares prediction input with metadata and configuration
- Workflow Delegation: Spawns PredictionProcess child workflow for actual processing
- 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
{
"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"]
},
"pi_web_api_output_config": {
"endpoint": "https://pi-server.com/piwebapi",
"prediction_tags": {"tag1": "web_id_1"},
"confidence_tags": {"tag2": "web_id_2"}
}
}
Architecture Diagram
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
- Input Data Gate: Applies configured filters for data quality validation
- Path Decision: Determines processing path based on filter results
- MLFlow Transform: Requests data transformation using MLFlow models
- Response Validation: Filters transform responses for quality assurance
- Content Validation: Filters transformed data content for quality check
- MLFlow Prediction: Executes prediction using transformed data
- Prediction Response Validation: Filters prediction responses for final quality check
- Export Delegation: Delegates to FormatAndExportPrediction workflow
- MinIO Cleanup: Cleans up expired offloaded payloads (if any, in
finallyblock)
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
- MinIO Cleanup: Automatic retention-based cleanup of offloaded payloads
- Comprehensive Monitoring: Detailed metrics and error reporting
Input Parameters
{
"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": {...},
"pi_web_api_output_config": {
"endpoint": "https://pi-server.com/piwebapi",
"prediction_tags": {"tag1": "web_id_1"},
"confidence_tags": {"tag2": "web_id_2"}
}
}
Architecture Diagram
flowchart LR
A[1. input_gate] --> B[2. request_transform] --> C[3. mlflow_response_gate] --> D[4. mlflow_content_gate] --> E[5. request_predict] --> F[6. mlflow_response_gate] --> G[7. format_and_export_prediction🔃]
G --> H[8. cleanup_minio_objects_expired]
B -.-> MLFlow[MLFlow]
E -.-> MLFlow[MLFlow]
H -.-> MinIO[(MinIO)]
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
- PI Web API Integration: Writes predictions and confidence to PI Web API for industrial systems
- Metrics Recording: Tracks export operations and performance metrics
Execution Flow
- Path Decision: Determines formatting path based on configuration
- Data Formatting: Formats prediction data for specific output requirements
- Transformed Data Processing: Optionally formats and exports transformed data separately
- PI Web API Export: Writes predictions and confidence to PI Web API (if configured)
- OPC Export: Writes predictions to OPC servers (if configured)
- PostgreSQL Export: Writes formatted predictions to database
- Metrics Recording: Records export performance and success metrics
Key Features
- Flexible Formatting: Configurable output formats for different destinations
- Multi-Destination Export: PostgreSQL, OPC server, and PI Web API integration
- Transformed Data Export: Optional separate export of MLFlow transformed data
- Performance Monitoring: Comprehensive metrics for export operations
- Error Handling: Robust error handling with notification integration
Architecture Diagram
flowchart LR
A[1. format_prediction/format_default_prediction] --> B[2. format_transformed_data] --> C[3. write_pi_web_api_data] --> D[4. write_opc_data] --> E[5. export_data_to_postgres] --> F[6. write_metrics]
A -.-> Format[Data Formatting]
B -.-> Transform[Transformed Data]
C -.-> PIWebAPI[PI Web API]
D -.-> OPC[OPC Servers]
E -.-> PostgreSQL[(PostgreSQL)]
F -.-> Prometheus[Prometheus]
Transformed Data Export
When transformed_data is provided in the input, the workflow will:
- Format the transformed data using
format_transformed_dataactivity - Export it to a separate table (
transform_table_name) asynchronously - Wait for both prediction and transformed data exports to complete
- This enables separate tracking of model transformations for analysis and debugging
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
- Data Loading: Loads training data using custom queries
- Model Retraining: Executes model retraining process
- Quality Validation: Validates retrained model performance
- Production Update: Updates production model if quality criteria met
- Data Export: Exports training data for analysis
Architecture Diagram
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)]
5. Drift Workflow (drift.py)
The Drift workflow detects data drift by comparing current data against a reference dataset from the MLflow Model Registry.
Execution Flow
- Data Loading: Loads target data and reference data in parallel
- Drift Calculation: Calculates univariate and multivariate drift metrics
- Data Export: Exports drift metrics to PostgreSQL
Architecture Diagram
flowchart LR
A[1. load_custom_query] --> C[3. calculate_drift] --> D[4. export_data_to_postgres]
B[2. get_reference_data] --> C
A -.-> Database[(Database)]
B -.-> MLFlow[MLFlow]
D -.-> PostgreSQL[(PostgreSQL)]
Input Parameters
{
"schedule_name": "hourly_drift",
"model_name": "temperature_model",
"model_id": "temp_001",
"schema": "sientia_data",
"source_table_name": "laborious_data",
"target_table_name": "drift_metrics",
"interval": 60,
"model_config": { "target": "temperature" },
"drift_metrics": ["kolmogorov_smirnov", "jensen_shannon", "wasserstein"],
"chunk_period": "min"
}
6. Simple Metrics Workflow (simple_metrics.py)
The SimpleMetrics workflow calculates regression metrics (RMSE, MSE, MAE, R²) by comparing predictions against actual target values.
Execution Flow
- Data Loading: Loads prediction vs target data via a JOIN query
- Metrics Calculation: Calculates configured regression metrics
- Data Export: Exports metrics to PostgreSQL
Architecture Diagram
flowchart LR
A[1. load_custom_query] --> B[2. calculate_simple_metrics] --> C[3. export_data_to_postgres]
A -.-> Database[(Database)]
C -.-> PostgreSQL[(PostgreSQL)]
Input Parameters
{
"schedule_name": "hourly_metrics",
"model_name": "temperature_model",
"model_id": "temp_001",
"schema": "sientia_data",
"predictions_table_name": "predictions",
"data_table_name": "laborious_data",
"target_table_name": "simple_metrics",
"interval_minutes": 60,
"model_config": { "target": "temperature" },
"metrics": ["rmse", "mse", "mae", "r2"]
}
📋 Prerequisites
- Python 3.11+
- Temporal server/cluster
- PostgreSQL database
- MLFlow server
- MinIO object storage (for MLFlow artifacts)
- MongoDB server (for notifications)
- OPC server(s) if using OPC export
- PI Web API server if using PI Web API export
Note: External dependencies must be available either through:
- Kubernetes cluster deployment
- Docker Compose setup
- Cloud-managed services
- Local installations
🚀 Installation
Local Development Setup
-
Clone the repository
git clone <repository-url> cd sientia-dataops-laborious -
Create virtual environment
python3.11 -m venv venv source ./venv/bin/activate -
Install dependencies
-
Install github cli
bash sudo apt update sudo apt install gh -y -
Authenticate with github
bash gh auth login -
Run the install_dependencies.sh script
bash chmod +x install_dependencies.sh ./install_dependencies.sh -
Create environment configuration file
cp .env.example .env # Edit .env with your connection details -
Configure external dependencies
You'll need to set up port forwarding or connections to external services. For example:
# 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:
# 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:
# 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:
# 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/laborious/activities/
pytest tests/laborious/workflows/
Manual Application Execution
For manual execution without scripts:
# 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
Code Quality & Validation
Overview
Since Python is not compiled, we validate quality, security, and correctness before execution.
Validation Tools
- Ruff: Linting and formatting
- mypy: Static type checking
- Bandit: Security analysis
- pytest: Unit/integration testing with coverage
Tools Installation
pip install -r requirements-dev.txt
Complete Validation
Run each validation step individually:
ruff format --check laborious/ tests/
ruff check laborious/ tests/
mypy laborious/
bandit -r laborious/ -ll
pytest tests/ --cov=laborious --cov-report=term-missing
Automatic Fixes
ruff format laborious/ tests/
ruff check --fix laborious/ tests/
Configuration
All settings reside in pyproject.toml (Ruff, mypy, pytest, Bandit).
CI/CD Integration
The workflow at .github/workflows/quality-gate.yml executes validations on each push/PR.
Best Practices
- Run all validation steps before committing
- Use
ruff check --watchfor continuous feedback - Add type hints and tests for new code
🧪 Testing
Test Structure
tests/
├── conftest.py # Global fixtures and env setup
├── laborious/
│ ├── activities/ # Activity implementation tests
│ │ ├── test_activities.py # Activities aggregator tests
│ │ ├── test_gates.py # Data quality gates and formatting tests
│ │ ├── test_mlflow.py # MLFlow operations and reference data tests
│ │ ├── test_storage.py # Storage and MinIO offload tests
│ │ ├── test_model_metrics.py # Drift and simple metrics tests
│ │ ├── test_opc.py # OPC operations tests
│ │ └── test_api.py # PI Web API operations tests
│ ├── workflows/ # Workflow orchestration tests
│ │ ├── test_predictions_batch.py
│ │ ├── test_minimal_retrain.py
│ │ ├── test_drift.py
│ │ ├── test_simple_metrics.py
│ │ └── subworkflows/
│ │ ├── test_prediction_process.py
│ │ └── test_format_and_export_prediction.py
│ └── utils/ # Utility function tests
│ ├── test_connectors_config.py
│ ├── models/
│ │ └── test_minio_dataframe_payload.py
│ ├── filters/
│ │ ├── test_conditional_filters.py
│ │ └── test_mlflow_filters.py
│ └── repository/
│ ├── test_model_repository.py
│ └── test_opc_repository.py
Test Coverage
The test suite provides comprehensive coverage for:
- Data Quality Gates: Input, response, and content validation filters
- Data Formatting: Prediction, transformed data, and retrain report formatting
- MLFlow Operations: Transform, predict, retrain, and reference data retrieval
- Workflow Orchestration: Complete workflow execution paths and error handling
- Metrics Recording: Performance monitoring and OPC export metrics
Test Execution
# 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/laborious/activities/test_gates.py
pytest tests/laborious/workflows/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
- Labels:
Prediction Operation Metrics
laborious_predictions_written_count: Counter for successful prediction exports- Labels:
pod_id,model_name,workflow_name
- Labels:
laborious_prediction_confidence_monitor: Gauge for current prediction confidence levels- Labels:
pod_id,model_name,workflow_name
- Labels:
laborious_prediction_response_time_monitor: Histogram for prediction response times- Labels:
pod_id,model_name,workflow_name - Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
- Labels:
OPC Export Metrics
laborious_prediction_opc_writing_count: Counter for OPC server write operations- Labels:
pod_id,model_name,workflow_name,opc_server_id
- Labels:
laborious_prediction_opc_writing_response_time_monitor: Histogram for OPC write response times- Labels:
pod_id,model_name,workflow_name,opc_server_id - Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
- Labels:
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 |
PI_WEB_API_BASE_URL |
PI Web API server base URL | None |
No |
PI_WEB_API_AUTH_TYPE |
PI Web API authentication type (basic/bearer) | None |
No |
PI_WEB_API_AUTH_TOKEN |
PI Web API authentication token | None |
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 |
MINIO_ENDPOINT_URL |
MinIO endpoint URL | http://localhost:9000 |
Yes |
MINIO_ACCESS_KEY |
MinIO access key | minioadmin |
Yes |
MINIO_SECRET_KEY |
MinIO secret key | minioadmin |
Yes |
MINIO_REGION_NAME |
MinIO region name | us-east-1 |
No |
MINIO_DEFAULT_BUCKET |
Default MinIO bucket | laborious |
No |
MINIO_RETENTION_HOURS |
Retention window (hours) for offloaded MinIO objects | 24 |
No |
SIENTIA_MINIO_OFFLOAD_THRESHOLD_MEGABYTES |
Offload threshold for DataFrame-derived payloads | int(1.5 * 1024 * 1024) |
No |
MinIO Payload Offload & Retention
Laborious uses MinIO to prevent Temporal workflow history from carrying very large in-memory payloads (pandas DataFrame-derived dicts).
Whenever a payload exceeds a configurable size threshold, it is stored as a parquet file in MinIO and the workflow history only keeps a lightweight reference.
Notes:
SIENTIA_MINIO_OFFLOAD_THRESHOLD_MEGABYTESsupports:- Integer bytes (e.g.
"1572864") - Float MiB (e.g.
"1.5"), converted to bytes asMiB * 1024 * 1024
- Integer bytes (e.g.
- Fallback behavior uses
1.5 MiBwhen the env var is missing or invalid.
Wire Contract: MinioDataFramePayload
The payload is implemented in laborious/utils/models/minio_dataframe_payload.py.
The dataclass does not store a pandas DataFrame field.
Instead, the DataFrame is only used at build time by:
MinioDataFramePayload.from_dataframe(...)MinioDataFramePayload.from_dataframe_to_dict(...)
After evaluation, the payload is serialized for Temporal as a flat dict:
- Inline path:
datacontainsdf.to_dict(), and MinIO keys (object_key,bucket, ...) are absent /None. - MinIO path: the dict contains:
bucketobject_key(full MinIO object name returned byMinioRepository.upload_file)object_prefix(directory prefix used for cleanup listing; relative to the repository namespace)uri(best-efforts3://<bucket>/<...>string)datais omitted / set toNone.
When an activity needs pandas operations, it resolves references using:
MinioDataFramePayload.retrieve(minio_repo)— downloads from MinIO or returns inline data as a DataFrame
MinIO Object Naming (Retention Parsing)
MinIO object basename (required convention):
{model_name}-{operation}-{timestamp}.parquet
Where:
model_name: model identifier used by the pipelineoperation:initial(SQL/query load before transform) ortransform(after MLFlow transform)timestamp:DATETIME_FORMAT_FILENAMEfromsientia_do.temporal.constants
The relative object key (under the repository namespace) is always shaped as:
training_datasets/{model_name}/{basename}
Retention cleanup parses timestamps from the basename using the -initial- / -transform- anchors.
model_name may contain hyphens; parsing is resilient to it.
Workflows / Activities Integration
Predictions batch uses MinIO offload as follows:
predictions_batchcallsActivities.load_query_with_minio_offload- On success, it puts the serialized
MinioDataFramePayloaddict intoprediction_input["data"].
- On success, it puts the serialized
sub_workflows/prediction_process- Tracks which MinIO prefixes were referenced for offloaded payloads.
- Runs
Activities.cleanup_minio_objects_expiredin afinallyblock (only when MinIO offload happened).
laborious/activities/gates.pyandlaborious/activities/mlflow.py- Resolve offloaded payloads transparently before constructing pandas
DataFrameobjects.
- Resolve offloaded payloads transparently before constructing pandas
Legacy: query_to_minio (Minimal Retrain)
Storage.query_to_minio is intentionally kept with its legacy behavior for minimal_retrain.
It always uploads parquet and returns {success, object_key, uri}.
It is not used by predictions batch MinIO offload, and its objects are not part of the retention parser described above.
Legacy MinIO object layout (relative key):
training_datasets/{model_name}/{object_prefix}_{timestamp}.parquet where object_prefix is sanitized
(slashes replaced by underscores) to keep a stable model-level directory.
OPC Configuration
For multiple OPC servers, use the OPC_CONFIG environment variable:
{
"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_URLOPC_NAMEOPC_SERVER_URIOPC_CERT_PATHOPC_PRIVATE_KEY_PATHOPC_SERVER_CERT_PATHOPC_RECONNECTION_INTERVAL
PI Web API Configuration
PI Web API configuration is built from environment variables using the build_api_config function from sientia_do.connectors_config. The configuration includes:
PI_WEB_API_BASE_URL: Base URL of the PI Web API serverPI_WEB_API_AUTH_TYPE: Authentication type ('basic' or 'bearer')PI_WEB_API_AUTH_TOKEN: Authentication token for API access
The PI Web API export is optional and can be configured per workflow through the pi_web_api_output_config parameter:
{
"pi_web_api_output_config": {
"endpoint": "https://pi-server.com/piwebapi",
"prediction_tags": {
"tag1": "web_id_1",
"tag2": "web_id_2"
},
"confidence_tags": {
"tag3": "web_id_3",
"tag4": "web_id_4"
}
}
}
Where:
endpoint: PI Web API endpoint URLprediction_tags: Dictionary mapping tag names to web IDs for prediction valuesconfidence_tags: Dictionary mapping tag names to web IDs for confidence values
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.
{
"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.
{
"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 aggregator
│ ├── gates.py # Data quality gates and filtering
│ ├── mlflow.py # MLFlow model operations
│ ├── storage.py # PostgreSQL queries and MinIO offload
│ ├── model_metrics.py # Drift and regression metrics
│ ├── opc.py # OPC server operations
│ └── api.py # PI Web API operations
├── workflows/ # Temporal workflow definitions
│ ├── predictions_batch.py # Main batch prediction workflow
│ ├── minimal_retrain.py # Model retraining workflow
│ ├── drift.py # Data drift detection workflow
│ ├── simple_metrics.py # Regression metrics 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
│ └── prepare_worker.py # Worker factory with autoscaling config
├── utils/ # Utility functions
│ ├── connectors_config.py # Environment-driven config builders
│ ├── models/ # Data models
│ │ └── minio_dataframe_payload.py # MinIO-offloaded DataFrame payload
│ ├── 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
│ └── minio_manager.py # MinIO object storage operations
├── metrics.py # Prometheus metrics definitions
└── __init__.py
Adding New Features
- Follow Temporal patterns for new workflows and activities
- Add comprehensive docstrings for all public methods
- Include Prometheus metrics for monitoring
- Add unit tests for new functionality
- Update this README with new features and configuration
🐛 Troubleshooting
Common Issues
-
Temporal Connection Failures
- Verify Temporal server is running and accessible
- Check namespace configuration and permissions
- Review server logs for connection issues
-
MLFlow Connection Issues
- Verify MLFlow server is running and accessible
- Check authentication credentials and permissions
- Ensure model names and versions exist
-
Database Connection Issues
- Verify PostgreSQL service is running
- Check connection credentials and network access
- Ensure proper connection pool configuration
-
OPC Connection Failures
- Verify OPC server is accessible
- Check certificate and key file paths
- Review OPC server logs for connection issues
-
PI Web API Connection Failures
- Verify PI Web API server is accessible
- Check authentication credentials and token validity
- Verify web IDs exist and have write permissions
- Review PI Web API server logs for connection issues
- Check notification system for error details
-
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:
export LOG_LEVEL=DEBUG
⚡ Performance Tuning
Key Parameters
- Worker Concurrency: Adjust
max_concurrent_workflow_tasksandmax_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
- Fork the repository
- Create a feature branch
- Make your changes with comprehensive testing
- Update documentation and docstrings
- 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.