如何用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.Popento 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
argparseto 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:
- Configure the
loggingmodule to write to both a file and the console. - Capture each child script's stdout/stderr, prefix lines with the script name, and feed them into the logger.
- 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
Child Script Requirements: Each of your scripts (python1_a.py, etc.) needs to use
argparseto 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 logicLogging Behavior:
- Logs are saved to
pipeline_execution.logand 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.
- Logs are saved to
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
相关产品推荐
相关产品推荐

