You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何监控函数执行流水线的运行状态?是否可借助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 raise line 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 like pydantic to 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 08:02:23