Python ProcessPoolExecutor:カスタムプロセスプール
Ostorlabでは毎日数百のモバイルアプリケーションをスキャンしています。各スキャンは非常に多くのリソースを消費するため、当初から、コードの速度を最適化し、クラウドリソースの利用を最大化する必要がありました。
Ostorlabでは毎日数百のモバイルアプリケーションをスキャンしています。各スキャンは非常に多くのリソースを消費しますが、同時に、当初からコードの速度を最適化し、クラウドリソースの利用を最大化する必要がありました。
Pythonにおける並行処理と並列処理についてここで改めて論じることはしません。このテーマについては優れた資料が数多くあり、私のお気に入りはLuciano Ramalhoの「Fluent Python」です。
ただし、簡単にいくつか補足しておきます。
CPythonでは、GILによって複数のネイティブスレッドが同時にPythonバイトコードを実行することが妨げられます。つまり、スレッディングはCPUバウンドの並列処理には適していません。このテーマについて詳しくは、「Understanding the Python GIL」と「Removing Python's GIL: The Gilectomy」の講演をお勧めします。 プロセスのほうがCPUバウンドの並列処理には適していますが、プロセスの生成には大きなオーバーヘッドが伴います。複数のタスクを並行かつ並列に実行することが目的であれば、各タスクを個別のプロセスで実行するのは非常に非効率です。 この問題を解決するため、concurrent.futuresライブラリがPython 3.2で導入され、Python 2.5にバックポートされました。このライブラリでは、複数のワーカー(デフォルトはコンピューターのコア数)を生成し、そこにタスクを渡して処理させることができます。タスクは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によるシリアライゼーションの制約のために、すぐに多くの制限に直面しました。主な問題は次の2つです。
一部のタスクがパラメーター付きのデコレーターで修飾されており、Error pickling <function>のようなエラーが発生しました。また、渡す引数が非常に大きくなることがあり、recursion limitエラーが発生しました。
2つ目の問題については、キューとシリアライゼーションの処理を経由せずにタスクに引数を渡す方法を考える必要がありました。
ProcessPoolExecutorを使う前は、argsパラメーターで引数を受け取るシンプルなProcessを使用していました。
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))
同じオブジェクトのシリアライゼーションに起因する余分な処理をすべて取り除いたことで、効果はさらに大きくなりました。
multiprocessingの活用、プロセス生成のオーバーヘッドの固定化、余分なシリアライゼーションの全面的な回避により、同じコードの実行時間は20分以上から1分未満へと短縮されました。
今後は、さらなるパフォーマンス向上が得られるかどうかを確かめるため、GoとUnikernelを使った実験を予定しています。