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

如何用Python subprocess构建并行串行结合的任务流水线?含传参与日志方案

Great question! Let's walk through building this pipeline step by step, covering parallel execution, argument passing, and robust logging that handles those progress print statements effectively.

Solution Overview

First, let's outline the core components we need:

  • Parallel Execution: Use subprocess.Popen to launch multiple scripts at once, then wait for all to complete successfully before moving on.
  • Serial Execution: Run subsequent scripts one after another, only proceeding if the previous script exits without errors.
  • Argument Passing: Use argparse to define and propagate parameters to all child scripts.
  • Logging: Capture and organize output from all scripts with timestamps, script identifiers, and both console + file persistence.
Best Practices for Logging

Since your scripts have progress print statements, you want to:

  • Keep output ordered and traceable (mark which script each log line comes from)
  • See progress in real-time (not just after the script finishes)
  • Save logs for later debugging
  • Avoid mixing stdout/stderr in a confusing way

The best approach here is to:

  1. Configure the logging module to write to both a file and the console.
  2. Capture each child script's stdout/stderr, prefix lines with the script name, and feed them into the logger.
  3. Use line buffering so progress prints show up immediately.
Example Code

Here's a complete implementation that ties everything together:

import subprocess
import argparse
import logging
import sys
import os

def setup_logging():
    """Configure logging to output to both file and console with timestamps."""
    logger = logging.getLogger("PipelineRunner")
    logger.setLevel(logging.INFO)
    
    # Format for log lines: timestamp - logger name - level - message
    formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
    
    # File handler to save logs
    file_handler = logging.FileHandler('pipeline_execution.log')
    file_handler.setFormatter(formatter)
    
    # Console handler for real-time output
    console_handler = logging.StreamHandler(sys.stdout)
    console_handler.setFormatter(formatter)
    
    logger.addHandler(file_handler)
    logger.addHandler(console_handler)
    return logger

def parse_args():
    """Parse command-line arguments to pass to all child scripts."""
    parser = argparse.ArgumentParser(
        description="Run a pipeline with parallel initial steps followed by serial execution."
    )
    # Example parameter - adjust based on your actual needs
    parser.add_argument(
        "--data-path",
        required=True,
        help="Path to input data directory, passed to all scripts"
    )
    parser.add_argument(
        "--script-dir",
        default=".",
        help="Directory where your Python scripts are located (default: current directory)"
    )
    return parser.parse_args()

def run_parallel(scripts, args, logger):
    """Launch multiple scripts in parallel, capture output, and check success."""
    processes = []
    for script_name in scripts:
        script_path = os.path.join(args.script_dir, script_name)
        # Build the command to run the script with arguments
        cmd = [sys.executable, script_path, "--data-path", args.data_path]
        
        logger.info(f"Starting parallel script: {script_name}")
        # Launch process with line-buffered output capture
        proc = subprocess.Popen(
            cmd,
            stdout=subprocess.PIPE,
            stderr=subprocess.STDOUT,  # Merge stderr into stdout for unified logging
            text=True,
            bufsize=1,  # Line buffering for real-time output
            universal_newlines=True
        )
        processes.append((proc, script_name))
    
    # Wait for all processes and handle their output
    all_successful = True
    for proc, script_name in processes:
        logger.info(f"Waiting for {script_name} to finish...")
        # Read output line by line as it's generated
        for line in proc.stdout:
            logger.info(f"[{script_name}] {line.strip()}")
        
        # Check exit code
        return_code = proc.wait()
        if return_code != 0:
            logger.error(f"{script_name} failed with exit code {return_code}")
            all_successful = False
        else:
            logger.info(f"{script_name} completed successfully")
    
    return all_successful

def run_serial(script_name, args, logger):
    """Run a single script serially, capture output, and return success status."""
    script_path = os.path.join(args.script_dir, script_name)
    cmd = [sys.executable, script_path, "--data-path", args.data_path]
    
    logger.info(f"Starting serial script: {script_name}")
    proc = subprocess.Popen(
        cmd,
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
        text=True,
        bufsize=1,
        universal_newlines=True
    )
    
    # Capture and log output
    for line in proc.stdout:
        logger.info(f"[{script_name}] {line.strip()}")
    
    return_code = proc.wait()
    if return_code != 0:
        logger.error(f"{script_name} failed with exit code {return_code}")
        return False
    logger.info(f"{script_name} completed successfully")
    return True

def main():
    args = parse_args()
    logger = setup_logging()
    
    # Step 1: Parallel execution of first three scripts
    parallel_scripts = ["python1_a.py", "python1_b.py", "python1_c.py"]
    logger.info("=== Starting Parallel Step ===")
    if not run_parallel(parallel_scripts, args, logger):
        logger.error("Parallel step failed. Aborting pipeline.")
        sys.exit(1)
    
    # Step 2: Run python2.py if parallel step succeeded
    logger.info("\n=== Starting Serial Step: python2.py ===")
    if not run_serial("python2.py", args, logger):
        logger.error("python2.py failed. Aborting pipeline.")
        sys.exit(1)
    
    # Step 3: Run python3.py if python2.py succeeded
    logger.info("\n=== Starting Serial Step: python3.py ===")
    if not run_serial("python3.py", args, logger):
        logger.error("python3.py failed.")
        sys.exit(1)
    
    logger.info("\n=== All Pipeline Steps Completed Successfully! ===")

if __name__ == "__main__":
    main()
Key Details to Note
  1. Child Script Requirements: Each of your scripts (python1_a.py, etc.) needs to use argparse to accept the parameters you're passing (like --data-path). For example, in python1_a.py:

    import argparse
    parser = argparse.ArgumentParser()
    parser.add_argument("--data-path", required=True)
    args = parser.parse_args()
    # Rest of your script logic
    
  2. Logging Behavior:

    • Logs are saved to pipeline_execution.log and printed to the console.
    • Each line from a child script is prefixed with [script_name] so you know which script generated it.
    • Progress prints from the child scripts show up in real-time thanks to line buffering.
  3. Error Handling:

    • If any parallel script fails, the pipeline stops immediately.
    • If python2.py fails, python3.py never runs.
    • All failures are logged with exit codes for debugging.

内容的提问来源于stack exchange,提问作者marosh maroshka

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:38:54