Python ProcessPoolExecutor:一个自定义的进程池
在 Ostorlab,我们每天扫描数百款移动应用,每次扫描都极其消耗资源,但与此同时,从一开始我们就不得不为速度优化代码,并最大化地利用云资源。
在 Ostorlab,我们每天扫描数百款移动应用,每次扫描都极其消耗资源,但与此同时, 从一开始我们就不得不为速度优化代码,并最大化地利用云资源。
这里不再重新讨论 Python 中的并发与并行,关于这一主题有大量优秀的资料, 我最喜欢的是 Luciano Ramalho 所著的《Fluent Python》。
不过,这里有一些简要的笔记:
在 CPython 中,GIL 会阻止多个原生线程同时执行 Python 字节码。这意味着多线程(Threading) 并不适合 CPU 密集型的并行。关于这一主题,我推荐以下演讲以获取更多信息:《Understanding the Python GIL》和《Removing Python's GIL: The Gilectomy》。 多进程(Processing)更适合 CPU 密集型的并行,然而派生一个进程会带来可观的开销,如果 目标是以并发且并行的方式运行多个任务,那么在单独的进程中运行每个任务效率极低。 为了解决这个问题,concurrent.futures 库在 Python 3.2 中被引入,并被向后移植到 Python 2.5。该库 允许派生出若干个工作进程(worker,默认数量为您计算机上的核心数),您可以把任务传递给它们 进行处理。任务在使用 Pickle 序列化之后,会通过队列发送。
以下是来自 Fluent Python 公开 Github 的一个示例:
"""Download flags of top 20 countries by population
ProcessPoolExecutor version
Sample run::
$ python3 flags_threadpool.py BD retrieved. EG retrieved. CN retrieved. ... PH retrieved.
US retrieved. IR retrieved. 20 flags downloaded in 0.93s
"""
# BEGIN FLAGS_PROCESSPOOL
from concurrent import futures
from flags import save_flag, get_flag, show, main
MAX_WORKERS = 20
def download_one(cc):
image = get_flag(cc)
show(cc)
save_flag(image, cc.lower() + '.gif')
return cc
def download_many(cc_list):
with futures.ProcessPoolExecutor() as executor: # <1>
res = executor.map(download_one, sorted(cc_list))
return len(list(res))
if __name__ == '__main__':
main(download_many)
# END FLAGS_PROCESSPOOL
在实际使用中,我们很快就因为 Pickle 序列化的局限而遇到了许多限制,两个主要 问题是:
一些任务被带参数的装饰器所装饰,这会引发一个与此类似的错误
Error pickling <function>。所传递的参数有时相当大,会导致一个 recursion limit 错误。
对于第二个问题,我们不得不设想一种方法,在不经过队列和
序列化过程的情况下,向任务传递参数。
在使用 ProcessPoolExecutor 之前,我们使用的是简单的 Process,它们在 args 参数中接收参数:
p = multiprocessing.Process(
target=_process_worker,
args=(self._call_queue,
self._result_queue,
self._initial_args,
self._initial_kwargs))
以这种方式传递参数在大小上没有任何问题,而且当然无需序列化。Python 允许传递 队列、管道,甚至使用 multiprocessing.Value 和 multiprocessing.Array 类型来传递共享内存(参见 multiprocessing)。
我们代码的特殊之处在于,所传递的参数始终相同且体积可观,因此我们 决定修改 ProcessPoolExecutor,在初始化时传入这些参数,再由它传递给所派生的 工作进程。
任务将以这些参数作为初始参数来执行,以下是主要的改动,添加 initial_args 和 initial_kwargs 属性:
class ProcessPoolExecutor(_base.Executor):
def __init__(self, max_workers=None, *args, **kwargs):
"""Initializes a new ProcessPoolExecutor instance.
Args: max_workers: The maximum number of processes that can be used to
execute the given calls. If None or not given then as many
worker processes will be created as the machine has processors.
"""
_check_system_limits()
if max_workers is None:
self._max_workers = multiprocessing.cpu_count()
else:
self._max_workers = max_workers
# Make the call queue slightly larger than the number of processes to
# prevent the worker processes from idling. But don't make it too big
# because futures in the call queue cannot be cancelled.
self._call_queue = multiprocessing.Queue(self._max_workers + EXTRA_QUEUED_CALLS)
self._result_queue = multiprocessing.Queue()
self._work_ids = queue.Queue()
self._queue_management_thread = None
self._processes = set()
# Shutdown is a two-step process. self._shutdown_thread = False self._shutdown_lock = threading.Lock()
self._queue_count = 0
self._pending_work_items = {}
self._initial_args = args
self._initial_kwargs = kwargs
将这些属性传递给所派生的进程:
def _adjust_process_count(self):
for _ in range(len(self._processes), self._max_workers):
p = multiprocessing.Process(
target=_process_worker,
args=(self._call_queue,
self._result_queue,
self._initial_args,
self._initial_kwargs))
p.start()
self._processes.add(p)
使用这些参数来运行每个任务:
def _process_worker(call_queue, result_queue, args, kwargs):
"""Evaluates calls from call_queue and places the results in result_queue.
This worker is run in a separate process.
Args: call_queue: A multiprocessing.Queue of _CallItems that will be read and
evaluated by the worker.
result_queue: A multiprocessing.Queue of _ResultItems that will written
to by the worker.
shutdown: A multiprocessing.Event that will be set as a signal to the
worker that it should exit when call_queue is empty.
"""
while True:
call_item = call_queue.get(block=True)
if call_item is None:
# Wake up queue management thread
result_queue.put(None)
return
try:
args += call_item.args
kwargs.update(**call_item.kwargs)
r = call_item.fn(*args, **kwargs)
except BaseException:
e = sys.exc_info()[1]
result_queue.put(_ResultItem(call_item.work_id, exception=e))
else:
result_queue.put(_ResultItem(call_item.work_id, result=r))
由于我们去除了因对同一对象反复序列化而产生的所有额外处理,收益甚至更大。
得益于多进程的使用、固定的进程派生开销,以及绕过了所有额外的序列化,同一段代码 的执行时间从超过 20 分钟缩短到了不足一分钟。
更进一步,我们计划尝试使用 Go 和 Unikernel,看看能否获得更多的性能提升。