如何为asyncio.create_subprocess_exec设置最大并发数
Great question! Limiting concurrent subprocesses when using asyncio.create_subprocess_exec is totally doable, and it’s a smart move to keep your machine from getting overwhelmed. The best tool for this job is an asyncio.Semaphore—it’s straightforward, effective, and lets you precisely control how many subprocesses run at once.
Using asyncio.Semaphore (Recommended Approach)
A Semaphore acts like a gatekeeper: it lets a fixed number of tasks pass through at a time. For your use case, you’ll initialize it with your desired concurrency limit, then wrap each subprocess call in a context manager that automatically acquires and releases the semaphore as processes start and finish.
Step-by-Step Example
Here’s a complete, runnable example that runs 500 subprocesses with a maximum of 5 concurrent processes at a time (adjust the limit based on your machine’s resources):
import asyncio async def run_single_process(command, semaphore): # Acquire the semaphore before launching the subprocess async with semaphore: proc = await asyncio.create_subprocess_exec( *command, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) # Wait for the process to finish and capture output stdout, stderr = await proc.communicate() # Optional: Handle output or errors as needed if stdout: print(f"Task {command[-1]} output: {stdout.decode().strip()}") if stderr: print(f"Task {command[-1]} error: {stderr.decode().strip()}") return proc.returncode async def main(): # Set your desired concurrency limit max_concurrent_processes = 5 semaphore = asyncio.Semaphore(max_concurrent_processes) # Generate 500 unique commands (replace this with your actual input-driven commands) commands = [["echo", f"job_{i}"] for i in range(500)] # Create all async tasks tasks = [run_single_process(cmd, semaphore) for cmd in commands] # Wait for all tasks to complete and collect results results = await asyncio.gather(*tasks) # Optional: Process results (e.g., check for failed processes) failed_tasks = [i for i, code in enumerate(results) if code != 0] print(f"All {len(results)} tasks completed. {len(failed_tasks)} tasks failed.") if __name__ == "__main__": asyncio.run(main())
Key Details to Keep in Mind
- Semaphore Behavior: The
async with semaphoreblock ensures that onlymax_concurrent_processessubprocesses are running at any given time. When a process finishes, the semaphore is released automatically, allowing the next waiting task to start. - Why Not
set_max_tasks?: TheAbstractEventLoop.set_max_tasksmethod limits the total number of async tasks in the event loop—not just subprocess-related ones. This is too broad for your use case, as it would include any other async tasks you might have running. A Semaphore gives you targeted control over subprocess concurrency. - Tune the Concurrency Limit: Adjust
max_concurrent_processesbased on your machine’s capabilities:- For CPU-bound subprocesses: Match the number to your CPU core count (e.g., 4-8 for a typical desktop) to avoid overloading the processor.
- For IO-bound subprocesses: You can safely set a higher limit (e.g., 10-20) since they spend most of their time waiting on input/output.
Bonus: Dynamic Adjustments (Optional)
If you need to tweak the concurrency limit on the fly (e.g., based on real-time system load), you can modify the semaphore’s internal count using semaphore._value (note: this is an internal attribute and not officially documented) or implement a custom semaphore class. For most use cases, though, a fixed limit will work perfectly.
内容的提问来源于stack exchange,提问作者Max Heiber

