引言
在 python 的 multiprocessing 模块中,信号量(semaphore) 是一种强大的同步原语,用于控制对共享资源的并发访问数量。它与互斥锁(lock)不同,互斥锁一次只允许一个进程访问,而信号量可以允许指定数量的进程同时访问。
你可以把信号量想象成一个有固定数量钥匙的柜子:一个进程要执行任务,必须先acquire()(取走)一把钥匙,任务完成后必须release()(归还)钥匙。当所有钥匙都被取走时,其他进程就必须等待。这个“钥匙”的数量,就是信号量允许的最大并发数。
核心api
semaphore(value=1): 创建一个信号量对象。value是初始化时“钥匙”的数量,即允许的最大并发进程数。acquire(blocking=true, timeout=none): 请求一把“钥匙”。- 若计数器 > 0,则计数器减1,进程继续执行。
- 若计数器 = 0,进程会阻塞,直到有可用的“钥匙”。
release(): 归还一把“钥匙”,计数器加1。- 上下文管理器: 推荐使用
with semaphore:语句,它会自动调用acquire()和release(),能有效避免因异常导致锁未释放的问题。
关键注意事项
在多进程编程中,信号量的创建位置至关重要。
- 正确做法:必须在主进程中创建信号量,然后作为参数传递给子进程。
- 错误做法:将信号量创建为全局变量。在 windows 或 macos(使用
'spawn'方式创建进程)上,每个子进程会重新导入模块,创建出各自的信号量副本,导致同步失效。
完整示例:模拟限流打印任务
下面这个例子模拟了一个最多允许3个进程同时执行任务的场景。
import multiprocessing
import time
import random
def worker(process_id, semaphore):
"""
模拟一个工作进程。
在进入临界区前,必须获取信号量。
"""
# 使用 with 语句管理信号量,更安全
with semaphore:
# 进入临界区
start_time = time.strftime("%h:%m:%s")
print(f"[{start_time}] 进程 {process_id} 开始执行任务...")
# 模拟一个耗时任务,耗时 1~3 秒
work_duration = random.randint(1, 3)
time.sleep(work_duration)
end_time = time.strftime("%h:%m:%s")
print(f"[{end_time}] 进程 {process_id} 任务完成 (耗时 {work_duration}秒)。")
if __name__ == "__main__":
# 1. 在主进程中创建信号量,指定最大并发数为 3
# 这就像创建了一个有 3 把钥匙的锁
max_concurrent = 3
semaphore = multiprocessing.semaphore(max_concurrent)
# 2. 创建并启动 10 个工作进程
processes = []
for i in range(10):
# 重要:将信号量对象作为参数传递给子进程
p = multiprocessing.process(target=worker, args=(i, semaphore))
processes.append(p)
p.start()
# 3. 等待所有子进程结束
for p in processes:
p.join()
print("所有任务执行完毕。")代码详解
- 创建信号量:
semaphore = multiprocessing.semaphore(3)创建了一个允许 3 个进程同时进入临界区的信号量。 - 传递信号量:在创建
process时,将semaphore作为参数传入worker函数。这是确保所有进程共享同一个信号量的关键。 - 控制并发:在
worker函数中,with semaphore:语句块内的代码就是临界区。任何时候,最多只有3个进程能同时执行其中的代码。 - 模拟任务:
time.sleep()模拟了进程的实际工作。
可能的输出结果(部分)
[14:23:01] 进程 0 开始执行任务... [14:23:01] 进程 1 开始执行任务... [14:23:01] 进程 2 开始执行任务... [14:23:03] 进程 1 任务完成 (耗时 2秒)。 [14:23:03] 进程 3 开始执行任务... # 进程1释放信号量,进程3立刻开始 [14:23:04] 进程 2 任务完成 (耗时 3秒)。 [14:23:04] 进程 4 开始执行任务... ...
从输出中可以看到,开始执行任务 的消息不会连续出现超过3次。每当一个进程完成并释放信号量后,下一个等待的进程便会立即获取并开始执行。
信号量 vs 进程池 (multiprocessing.pool)
这两种方式都可以控制并发,但适用场景不同:
semaphore:更灵活,可以精确控制代码中任意一个代码块的并发数,而不仅仅是整个任务函数。pool:更高层,主要用于管理一组同质任务的并发执行。
如果你的需求是“同时最多运行n个任务”,pool 是更简单的选择。如果你的需求是“在任务执行过程中,某个特定资源的访问并发数不能超过n”,那么 semaphore 是更合适的工具。
总结
- 核心思想:信号量通过一个计数器来控制对共享资源的并发访问数量。
- 关键用法:在
if __name__ == "__main__":块中创建multiprocessing.semaphore(n),并作为参数传递给子进程。 - 最佳实践:使用
with semaphore:语句管理信号量的获取和释放,避免死锁。 - 适用场景:限制对有限资源(如数据库连接、网络带宽、特定硬件)的并发访问。
以上就是python通过信号量控制多进程并发的完整示例的详细内容,更多关于python信号量控制多进程并发的资料请关注代码网其它相关文章!
发表评论