如何用Python的concurrent.futures并行化共享数组的访问操作
解决思路与代码优化
你的问题核心是把原本在主线程串行执行的result[locs] += mask[locs]也放到多线程里跑,但首先要注意:多线程直接修改共享numpy数组会有数据竞争问题——+=操作不是原子的,多个线程同时改同一个元素会导致结果错误。下面给你两种可行方案,还有一个能直接秒杀问题的数学优化。
方案1:局部增量数组合并(推荐)
这个思路是让每个线程生成自己的局部结果数组,线程之间互不干扰,最后主线程把所有局部数组加起来得到最终结果。完全避免锁的开销,安全又高效。
修改后的代码:
import numpy as np import time import concurrent.futures MAX = 100 SIZE = 500 mask = np.random.randint(0, MAX, (SIZE, SIZE)) def process_image(i): start = time.time() locs = np.where(mask > i) # 每个线程单独创建局部数组,只在目标位置赋值 local_result = np.zeros((SIZE, SIZE)) local_result[locs] = mask[locs] print(f" process_image({i}) took {round(time.time() - start, 2)} secs.") return local_result if __name__ == '__main__': result = np.zeros((SIZE, SIZE)) with concurrent.futures.ThreadPoolExecutor(max_workers=32) as executor: # 批量提交任务,获取所有局部结果 for local_res in executor.map(process_image, range(MAX)): # 累加局部数组到最终结果 result += local_res print(result)
为什么选这个?
- 没有共享资源竞争,线程安全无需担心。
- numpy的数组累加是底层优化过的,比串行更新效率高很多。
方案2:锁保护共享数组更新
如果必须直接在线程里修改共享数组,就得用锁来保证同一时间只有一个线程操作数组。但锁会把更新操作串行化,要是更新本身特别耗时,并行的优势会大打折扣,只适合内存紧张(局部数组占空间太大)的场景。
代码示例:
import numpy as np import time import concurrent.futures import threading MAX = 100 SIZE = 500 mask = np.random.randint(0, MAX, (SIZE, SIZE)) # 创建全局锁 update_lock = threading.Lock() def process_image(i): start = time.time() locs = np.where(mask > i) # 加锁后再执行更新,避免竞争 with update_lock: result[locs] += mask[locs] print(f" process_image({i}) took {round(time.time() - start, 2)} secs.") if __name__ == '__main__': result = np.zeros((SIZE, SIZE)) with concurrent.futures.ThreadPoolExecutor(max_workers=32) as executor: executor.map(process_image, range(MAX)) print(result)
额外惊喜:数学等价的极致优化
仔细看你的逻辑:每个元素mask[x,y] = v,会在i从0到v-1的循环里被累加v次(每次加v),所以最终结果其实是v * v!直接用numpy的向量化操作就能一步搞定,比任何多线程都快:
result = np.square(mask) # 或者等价写法:result = mask * mask
这是最效率的方案,因为numpy的底层是C实现的向量化运算,完全绕开了Python的线程开销。
内容的提问来源于stack exchange,提问作者nickponline
相关产品推荐
相关产品推荐

