如何监控函数执行流水线的运行状态?是否可借助RxPy实现?
Alright, let's break down how to add robust monitoring to your function pipeline while keeping things flexible for different storage targets (files, databases, MongoDB). Here's a practical, modular approach that plays nicely with your existing reduce-based pipeline setup:
1. Build a Reusable Monitoring Decorator
First, we'll create a decorator that wraps each function in your pipeline to track execution details. This way, you don't have to modify your original f1/f2/etc. functions at all.
import time import traceback from typing import Callable, Dict, Any def monitor_function(metadata_handler: Callable[[Dict[str, Any]], None]): def decorator(func: Callable) -> Callable: def wrapper(*args, **kwargs): start_time = time.time() # Initialize monitoring data with basic info monitoring_data = { "function_name": func.__name__, "call_time": start_time, "finish_time": None, "success": True, "return_value": None, "exception": None, "traceback": None } try: # Execute the original function and capture its return value result = func(*args, **kwargs) monitoring_data["return_value"] = result except Exception as e: # Log failure details if an exception occurs monitoring_data["success"] = False monitoring_data["exception"] = str(e) monitoring_data["traceback"] = traceback.format_exc() # Optional: Uncomment below to stop the pipeline on first failure # raise finally: # Calculate duration and send data to your storage handler monitoring_data["finish_time"] = time.time() monitoring_data["execution_duration"] = monitoring_data["finish_time"] - start_time metadata_handler(monitoring_data) return monitoring_data["return_value"] return wrapper return decorator
2. Create Storage Handlers for Your Targets
Next, we'll build modular handlers to send monitoring data to files, SQLite, or MongoDB. Each handler follows the same interface, so you can swap them out easily.
File Storage Handler
Logs data to a JSON lines file for simple, human-readable records:
import json from datetime import datetime class FileMetadataHandler: def __init__(self, file_path: str = "pipeline_monitoring.log"): self.file_path = file_path def __call__(self, data: Dict[str, Any]): # Convert timestamps to readable ISO strings data["call_time_str"] = datetime.fromtimestamp(data["call_time"]).isoformat() data["finish_time_str"] = datetime.fromtimestamp(data["finish_time"]).isoformat() with open(self.file_path, "a") as f: f.write(json.dumps(data) + "\n")
MongoDB Storage Handler
Stores data directly in a MongoDB collection (requires pymongo):
from pymongo import MongoClient class MongoMetadataHandler: def __init__(self, db_name: str = "pipeline_monitor", collection_name: str = "function_executions"): # Adjust the connection string here for remote MongoDB instances self.client = MongoClient() self.collection = self.client[db_name][collection_name] def __call__(self, data: Dict[str, Any]): # MongoDB natively supports float timestamps, so we can insert directly self.collection.insert_one(data)
SQLite Storage Handler
Stores data in a local SQLite database for structured querying:
import sqlite3 class SQLiteMetadataHandler: def __init__(self, db_path: str = "pipeline_monitor.db"): self.db_path = db_path self._initialize_database() def _initialize_database(self): # Create the table if it doesn't exist with sqlite3.connect(self.db_path) as conn: cursor = conn.cursor() cursor.execute(""" CREATE TABLE IF NOT EXISTS function_executions ( id INTEGER PRIMARY KEY AUTOINCREMENT, function_name TEXT NOT NULL, call_time REAL NOT NULL, finish_time REAL NOT NULL, execution_duration REAL NOT NULL, success BOOLEAN NOT NULL, return_value TEXT, exception TEXT, traceback TEXT ) """) conn.commit() def __call__(self, data: Dict[str, Any]): with sqlite3.connect(self.db_path) as conn: cursor = conn.cursor() cursor.execute(""" INSERT INTO function_executions ( function_name, call_time, finish_time, execution_duration, success, return_value, exception, traceback ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) """, ( data["function_name"], data["call_time"], data["finish_time"], data["execution_duration"], data["success"], str(data["return_value"]), # Convert objects to strings (adjust if needed) data["exception"], data["traceback"] )) conn.commit()
3. Wrap Your Pipeline & Execute
Now we'll modify your existing run function to wrap each pipeline function with our monitor, then execute as usual:
import functools def run_monitored_pipeline(pipeline: list[Callable], metadata_handler: Callable[[Dict[str, Any]], None]) -> Callable: # Wrap every function in the pipeline with the monitoring decorator monitored_pipeline = [monitor_function(metadata_handler)(func) for func in pipeline] # Reuse your original reduce logic to compose the monitored functions return functools.reduce(lambda f, g: lambda x: f(g(x)), monitored_pipeline, lambda x: x) # Example test functions (replace with your actual functions) def f1(x): return x + 1 def f2(x): return x * 2 def f3(x): # Simulate an error for testing if x > 5: raise ValueError("Value too big!") return x - 3 def f4(x): return {"result": x} # Usage Example if __name__ == "__main__": # Pick your storage handler (swap with Mongo/SQLite as needed) handler = FileMetadataHandler("pipeline_logs.json") # handler = MongoMetadataHandler() # handler = SQLiteMetadataHandler() # Your original pipeline pipeline = [f1, f2, f3, f4] # Create the monitored pipeline runner run_pipeline = run_monitored_pipeline(pipeline, handler) # Test with a valid input print("Running pipeline with x=2:") try: final_result = run_pipeline(2) print(f"Final output: {final_result}") except Exception as e: print(f"Pipeline failed: {e}") # Test with an input that triggers an error in f3 print("\nRunning pipeline with x=3:") try: final_result = run_pipeline(3) print(f"Final output: {final_result}") except Exception as e: print(f"Pipeline failed: {e}")
Key Considerations
- Exception Behavior: By default, the decorator swallows exceptions to let the pipeline continue, but you can uncomment the
raiseline in the decorator to make the pipeline fail immediately on the first error. - Return Value Serialization: For storage targets like files or SQLite, we convert Python objects to strings. If you need to preserve object state, consider using
pickle(note security risks for untrusted data) or a library likepydanticto serialize objects to JSON-compatible structures. - Scalability: The handler pattern keeps your code modular—you can easily add new storage targets (like PostgreSQL, Elasticsearch) by creating a new class that implements the
__call__method.
内容的提问来源于stack exchange,提问作者nanounanue

