エラー処理から構造化並行性へ
原文は Nelson Elhage により に公開されました。 このブログを購読する
並行プログラムにおけるエラー処理を、どのように考えるべきだろうか。
シングルスレッドのプログラムでは、実装や具体的なやり方には多様なバリエーションがあるものの、ほぼ標準的なパターンに収束している。エラーが発生すると、それを処理できるスタックフレームが見つかるまでスタックを遡って伝播させる。その過程でスタックフレームを順番にアンワインドし、各フレームにリソースを適切にクリーンアップしたり破棄したりする機会を与えるのだ。
このパターンは、多くの現代的な言語(C++、Python、Java)における明示的な例外処理機構をそのまま説明している。これらの言語はいずれも、フレームのアンワインド時にクリーンアップを行う仕組み(RAII、finallyブロック、Pythonのコンテキストマネージャ)を備えている。だがそれだけでなく、Rustにおける標準的なパターン(Resultを返す、?演算子、アンワインド時にdropを呼び出す)や、Goにおける定番のパターン(おなじみのif err != nil { return err }とクリーンアップのためのdefer)、さらには現代のCコードの多くさえも、goto errorパターンといった形でこの説明に当てはまる(Linuxカーネルにおける例を参照)。
今日の視点から見ると、この説明はあまりに一般的で中身がないように思えるかもしれないが、常にそうだったわけではない。他の(ほとんど廃れた)エラー処理アプローチとしては、Lispの「restarts」機構、Cのlongjmp1、UNIXシグナルのような「トラップ」機構、そして悪名高いVisual Basicのon error節などがある。「アンワインド」自体も、構造化されたコールスタックが存在することを前提としており、この概念自体も発明され普及する必要があったものだ。
今回は問いたい。単一のスタックが存在しない並行プログラムのために、このパターンをどうアップデートすべきだろうか。複数の並行タスクが存在する中で2、エラー状態を処理するためにコードをどう構成すべきなのだろうか。
未処理のエラー
エラー処理について考えるうえで、おそらく最も単純なケースは、エラーが発生したにもかかわらず、それを明示的に処理するコードがどこにもない場合だ。シングルスレッドのプログラムでは、エラーがエントリーポイントまで「バブルアップ」し、有用なエラーメッセージやスタックトレースを残してプログラムを終了することを期待する。
複数のタスクを持つ並行プログラムでは、何が起こるべきかはそれほど自明ではない。エラーを上へバブルさせ、エラーを送出したタスクを終了させることはできるが、その後はどうなるのだろうか。具体的に考えるために、次のようなトイプログラムを考えてみよう3:
import threading
import time
def background_thread():
# This was supposed to be running some background work, but it
# encountered an error!
raise ValueError("oops")
def main():
threading.Thread(target=background_thread).start()
# do the main work
time.sleep(5)
print("All done, exiting!")
main()このプログラムは、二つの並行なスレッドを動かそうとする。一つは「メイン」の処理を行う関数、もう一つは「バックグラウンド」で何らかの処理を行うものだ。ところが、バックグラウンドスレッドは即座に例外を送出してしまう。何が起こるべきだろうか。
このプログラムをさまざまな言語実装で実行してみると、大きく分けて二つの一般的なアプローチがあることがわかる。私に言わせれば、これらはほとんど自明な二つの選択肢だ。
- エラーを出力し、スレッドを終了させ、そのまま他のすべてのスレッド(あるいはバリエーションとして、メインスレッドだけ、あるいはすべての「非デーモン」スレッド)が終了するまで実行を続ける。(Java、Python)
- エラーを出力し、即座にプログラム全体を終了させる。(Go、Rust、C++)
Pythonでは、プログラムは例外を即座にログ出力するが、その後も実行を続ける。time.sleepが終わるのを待ち、成功の終了コードで終了してしまうのだ。
$ time python exception_thread.py
Exception in thread Thread-1 (background_thread):
Traceback (most recent call last):
File "/Users/nelhage/.pyenv/versions/3.11.0/lib/python3.11/threading.py", line 1038, in _bootstrap_inner
self.run()
File "/Users/nelhage/.pyenv/versions/3.11.0/lib/python3.11/threading.py", line 975, in run
self._target(*self._args, **self._kwargs)
File "/Users/nelhage/Sync/code/structured-concurrency/exception_thread.py", line 7, in background_thread
raise ValueError("oops")
ValueError: oops
All done, exiting!
python exception_thread.py 0.03s user 0.02s system 0% cpu 5.124 total
$ echo $?
0どちらの選択肢も、あまり満足できるものではない。
プログラム全体を強制終了させるのは、あまりに強力すぎる手段だ。私自身、バックグラウンドで動いていた監視用のgoroutineで処理されなかったpanicが、クリティカルなデーモンをクラッシュさせたことが直接の原因となった、重大度の高いインシデントの対応を手伝ったことがある。そのケースでは、「本来の」処理は中断させずにそのまま続けさせたかったはずだ。
逆に、プログラムをそのまま動かし続けると、ほぼ確実にテストも想定もされていない状態に陥り、他のタスクが死んだタスクの進行や特定のアクションを期待している場合、デッドロックあるいはそれ以上の事態に陥るリスクが高くなる。
プログラムを動かし続けると、開発中にも苦労することになる。開発中は、多くのエラーが「しょうもない」開発者のミス――タイプミスや単純なロジックエラー、その他局所的でささいな間違い――である。そうしたエラーに遭遇したときは、通常できるだけ早く修正してプログラムを再起動したい。プログラムがまだ動き続けていると、手動で再起動しなければならず、他のタスクが出力を生成していると、例外がスクロールバックの中に埋もれてしまうこともある。
エラーはどこに届けるべきか
ある意味で私たちが求めているのは、エラーを「転送」するためのより良い場所だ。シングルスレッドのプログラムでは、その場所は「呼び出し元」である。並行処理がある場合、タスクには最終的にリターンする呼び出し元が存在しないため、代わりにどうすべきだろうか。
ここでPythonのasyncioフレームワークから少しヒントを得よう。asyncioは、上記のどちらとも異なる「第三の道」を取っている。asyncioでは、タスクはTaskオブジェクトとして表される。これはイベントやロック、ソケットなどと同様に待機(wait)できるオブジェクトだ。Taskを待機すると、そのタスクが完了するまでブロックする。タスクが例外を送出した場合、その例外はTaskを待機していた者に再配送される。
このアプローチは、タスクから漏れ出た例外を誰が処理すべきかについて意見を押し付けない。代わりに、プログラマ自身がその判断を下すための手段を提供するのだ。
しかし、これには大きな欠点がある。誰もTaskを待機しなかった場合に、そのタスクが例外を送出すると、プログラム終了時まで例外は実質的に飲み込まれ、最終的に警告とともに表示されるだけになる。もし誰かがそのタスクがキューなどを介して出力を生成するのを待っていたとすれば、プログラムはただ永遠にハングし、静かに、不可解に、謎めいたままになるのだ。
実際、状況はかなり深刻で、私はときどき「デフォルトでは、asyncioはメイン以外のタスクの例外を完全に飲み込む」と要約することさえある。文字通りに正しいわけではないが、私の経験では出発点としては十分に適切なメンタルモデルであり、多くの開発者のasyncioに関する経験をよく説明している。
もし常にタスクを待つようにしたらどうなるか
asyncioのアプローチには、すべてのTaskを忘れずに待機する限り、推奨できる点が多くある。この不変条件をルールとして強制するにはどうすればよいだろうか。最も単純なルールはこうかもしれない。タスクをスポーンしたら、そのタスクを待機する責任も負う、というものだ。このパターンを、asyncio.create_task(coro)を非同期コンテキストマネージャにして、ブロックを抜ける前にタスクを待機させることで実装できるだろう。
# (n.b. this is not real API in any version of Python)
async with asyncio.create_task(background_task()) as task:
# …
# `task` will be waited for on exit from this block, and any exception raisedもちろん、多くのタスクを――あるいは動的に決まる数のタスクを――生成したいことも非常に多い。そこでcontextlib.ExitStackのイディオムを借りて、任意の数のタスクを生成できる一つのコンテキストマネージャオブジェクトを用意するという手がある。例えば次のようなものだ。
async with TaskLauncher() as tasks:
tasks.create_task(background_task())
# do the main work via a second task
tasks.create_task(asyncio.sleep(5))
# All tasks will be waited for on exit from the regionTaskLauncher経由でのみ新しいタスクをスポーンする限り、すべてのタスクが明確な親に「所属」し、親が子から漏れ出た例外を待機する責任を負うという、明確な親子関係が得られる。いずれかのタスクが処理されない例外を送出すれば、その例外はこの階層を伝って上へバブルする。誰もそれを捕捉しなければ、最終的にはルートタスクとasyncio.runの呼び出しに行き着く。並行な例外処理の問題を、馴染みのあるシングルスレッド版の問題へと大きく変換できたわけだ。
二つの問題
残念ながら、これで終わりにはほど遠い。上記のスケッチには、深刻で関連し合った二つの課題があり、どちらも自明な修正方法はない。
デッドロック
第一に、デッドロックだ。タスクの親は「いずれ」それを待機すると先ほど述べたが、それは親が実際にコンテキストマネージャを抜けた場合にのみ真実となる。だが、もし親が、決して完了しない処理を待機しているとしたらどうなるだろうか。子タスクの一つがエラーに遭遇したために、その処理を担うはずだった子が動かなくなった場合だ。
問題の一端を示す短い例がこちらだ。
from task_launcher import TaskLauncher
import asyncio
async def do_work(job_id, done_event):
if job_id == 1:
raise ValueError("Oops, job 1 failed!")
done_event.set()
async def main():
async with TaskLauncher() as tasks:
events = []
for i in range(4):
done = asyncio.Event()
tasks.create_task(do_work(i, done))
events.append(done)
for ev in events:
await ev.wait()
print("All done!")
if __name__ == '__main__':
asyncio.run(main())この例は、よくあるパターンを様式化したものだ。多くの「ファンアウト」や「ファンアウト/ファンイン」型の並行パターンは、基本的に「いくつかのタスクを起動し、それらのタスクが何らかの処理を行い、親がその処理を待つ」という流れを持つ。
この特定のバグに対する単純な修正方法はたくさんある4。しかし、私たちが求めているのは、この特定のバグを自動的に、あるいは何らかの一般的な方法で修正するパターン、あるいは少なくとも、そうした安易な罠が仕掛けられたままにならないようにするパターンだ。私の経験では、並行システムを初めて開発する際に、こうした種類のデッドロックに遭遇するのは驚くほど簡単だ。
リソースリーク
シングルスレッドのプログラムでは、エラーからの復旧の課題は、制御フローをアンワインドすることだけでなく、失敗した操作に関連して進行中だった処理やそれに紐づくリソースを「元に戻す」あるいは「クリーンアップ」することを確実にすることにある。メモリを解放したり、ファイルハンドルを閉じたり、データ構造を一貫した状態に戻したりする必要があるかもしれない。
並行プログラムでは、失敗した操作に関連するリソースが、複数の異なる実行タスクにまたがっている可能性がある。TaskLauncherに関連付けられた一つのタスクが失敗した場合、スポーンされたすべてのタスクが何らかの形で停止するか、少なくとも放棄された状態をクリーンアップする機会を得られるようにする必要がある。
ある意味で最も単純な修正は、TaskLauncherが終了する前にスポーンされたすべてのタスクを待機するようにすることだ。しかし、実際には、その変更だけではデッドロックの問題が劇的に悪化してしまう。
キャンセルはしたくないのだが……
もし私たちが以下の両方を望むなら、
- 追加のエラーがないか把握し、関連するリソースをクリーンアップするために、リターンする前にすべての子タスクを待機すること、そして
- 潜在的に無制限の時間待つことなく、いずれかの子タスクのエラーに迅速に対応すること、
そのためには、別のタスクで起きたイベントに応じて、任意のタスクに早期かつ迅速な終了を要求する方法が基本的に必要だと私は考える。言い換えれば、キャンセル機構が必要なのだ。
私たちはここで展開してきた特定のパラダイムに照らしてこの結論に至ったが、これはもっと広範に当てはまり、振り返ってみればかなり直感的でもある。どんな並行性パラダイムにおいても、「協調して動作する複数の並行タスク」というものは必ず存在し、それは「そのうちの一つが予期せず死んだらどうなるか」という問いに答える必要があることを意味する。そして、一般的な答えとして「他のタスクにキャンセルして早期に終了するよう求める」以外を想像するのは難しい。
もちろん、特定の並行プログラムやパターンであれば、汎用的なキャンセル機構なしでも、場当たり的な仕組みの組み合わせや、慎重な推論と構築によって実装することは十分可能だ。しかし、汎用的で合成可能な並行パラダイムを、それなしで思い描くことは私には難しい。
この結論は私にとって愉快なものではない。キャンセルを実装しサポートすることは難しいからだ。キャンセルは事実上すべてのコードに新たなエラーパスを導入し、それは本質的に非同期で、推論やテストが困難なものとなる。Cのpthread_cancelや、JavaのThread.stop、RubyのThread.terminateといった、歴史的なキャンセル機構の試みは、せいぜい極めて繊細でエラーが起こりやすく、最悪の場合は根本的に使い物にならないものだった。
並行処理という文脈では、少なくともいくつかの利点はある。すでに並行コードを書いているのであれば、キャンセル機構は「非同期に起こりうること」をもう一つ増やすことになるが、少なくとも私たちはすでにその種の問題を抱えている。asyncioのような協調的並行システムでは、キャンセルをawaitポイントでのみ発生するように制限でき、潜在的な混乱の範囲を狭めることができる。一方Goでは、キャンセルはContextオブジェクトに組み込まれ、コードが明示的にキャンセルをチェックすることを求めるという、異なるトレードオフを取っている。
より一般的には、私たちはこの数十年で多くのことを学んできており、これらの新しいシステムの中には、実際に多かれ少なかれ実用的な汎用キャンセル機構を備えているように見えるものもある。この記事はキャンセル機構の課題や設計空間について深掘りすることを意図したものではないので、ここでは何らかのキャンセルAPIが存在するものと仮定して5、先に進むことにする。
キャンセルを伴うタスクのツリー
もしタスクをキャンセルする能力があるなら、それを「タスクのツリー」というアイデアと組み合わせて、並行なエラー処理に対するかなり汎用的な解決策を作り出すことができる。
TaskLauncherによって起動されたいずれかのタスクが例外を送出した場合(コンテキストマネージャを実行している親タスク自身も含む)、他のすべてのタスク(子タスクも親タスク自身も)をキャンセルする。- 子タスクには、キャンセルを検知し自身のリソースをクリーンアップする機構を与える。通常これは、何らかの形で通常のエラー処理機構を再利用することを意味する。例えば、キャンセルによって
CancelledError例外が送出され、タスクはそれを捕捉して再送出したり、finallyブロックやコンテキストマネージャを使ったりできる。 TaskLauncherコンテキストを抜ける際に、すべての子タスクが終了するのを待つ。正常終了であれ、未処理の例外による終了であれ、キャンセルに応じた終了であれ、である。- そして、いずれかの子タスクがエラーを送出していた場合、それを親タスクに再送出する。
この機構により、並行なエラーはシングルスレッドの場合とかなり似た振る舞いをするようになる。処理されなければ捕捉され、上へ伝播する。通常のエラー処理機構を使って捕捉し処理することもできる(子タスクに由来するエラーを含めて)。通常のシングルスレッドの場合と同じようにエラー後のクリーンアップを行うコードを書いておけば、複数のタスクが存在する場合でも(少なくとも大半は)適切なクリーンアップが得られるはずだ。
その代わり、並行コードに追加の構造を課すことになる。タスクを親子の階層にネストし、それらのライフタイムが適切に入れ子になるようにしなければならないのだ。
構造化並行性
ここで白状すると、これらのアイデアはいずれも新しいものではなく、私の発明でもない(この形でアイデアにアプローチした解説を他に見たことはないが)。ネストしたライフタイムを持つタスクのツリーと、ツリー全体にわたる自動的なキャンセルというこのパラダイムは、近年「構造化並行性(Structured concurrency)」という名前のもとで、ゆっくりと着実に人気と採用を広げてきたものだ。
多くのアイデアや構成要素には長い系譜があるが、この考え方は私の知る限り2016年に初めて命名され、おそらくtrioフレームワークと、trioの作者であり主任メンテナであるnjsによるこの考え方を探求した詳細なエッセイによって最も広く知られるようになった。
Python 3.11以降、PythonのasyncioにはTaskGroupクラスが含まれており、これは本質的に上でスケッチしたTaskLauncherの本番対応版である。trioでは同様のものを「nursery」と呼んでいる。Goでは、errgroupパッケージが本質的に同じセマンティクスを提供しており、キャンセルをサポートするcontextパッケージの上に構築されている。
構造化並行性には多くの利点があり、私自身や他の多くの人々は、このスタイルでプログラムを書くことで、正確で安全な並行コードを書くことがはるかに容易になることを実感している(もちろん他の課題が依然として残るにせよ)。このパラダイムとその利点についてより徹底的に探求したものとして、先ほどもリンクしたnjsの古典的なエッセイを強く推薦する。
コーダ:なぜエラー処理なのか
最後に、エラー処理について振り返って締めくくりたい。そもそもなぜこの切り口やこの記事について考え始めたのか、ということだ。
プログラマがエラー処理について考えるとき、それはしばしば「堅牢性」に関することや、「本番」あるいは「ちゃんとしたソフトウェア」のための関心事として分類されがちだと思う――「大規模に」なったときや、何かが「信頼性」を求められるとき、あるいは多くの異なる環境で動作しネットワークからの予期せぬ入力を処理する必要があるときなどに、気にかけなければならないトピック、といった具合だ。
そしてそれらはすべて正しく、そのようなシステムにとっては、何がうまくいかなくなる可能性があるか、そしてそれをどう慎重に処理するかをよく考えることが確かに重要だ。
とはいえ、私がこの思考を始めたのは「成熟した、堅牢なプログラム」という方向からではなく、まったく逆の端、すなわち新しいコードをゼロから書く際の開発体験について考えることからだった。先ほど簡単に触れたように、新しいプログラム――並行であるか否かにかかわらず――を書くときには、非常に頻繁に「しょうもないバグ」が大量にある初期段階があり、それらをできるだけ素早く片付けていく必要がある。
構造化並行性フレームワークの外で並行プログラムを書いた私の経験では、「プログラムを実行する、しょうもないバグを見つける、修正する」という基本的な開発ループを回すこと自体が、非常に苛立たしいほど困難になることが多い。正確に言えば、シングルスレッドのプログラムならきれいなスタックトレースを出力して終了するようなしょうもないバグが、デッドロックになったり、飲み込まれたり、あるいはさらに厄介なことになったりする悪癖があるからだ。そして、場当たり的にエラー処理を追加しようとすると、事態がさらに悪化することさえあると私は感じている。例えば、大きな並行処理の最後ですべてのエラーを集めて一箇所でログ出力できるように、パイプラインを通じてエラーを「転送」するのが「自然な」アプローチだと気づくことがときどきあった。このアプローチは機能するが、プログラム全体が完了するまでいかなるエラーにも気づけないことを意味する場合もあり、開発中には本当に苛立たしい。
だからこそ、構造化並行性のアプローチを採用すること、あるいは少なくともそれを基本的な考え方やパラダイムとして取り入れること――たとえ環境に「真の」構造化並行性ライブラリがなくても――が、実際には並行プログラムをそもそも書いたりデバッグしたりすることを劇的に容易にすることに私は気づいた。投げ捨て用のプロトタイプでさえ、最終的に、あるいは本番環境で、といった先の話ではなく、ほぼ即座に利益をもたらすのだ。
longjmpは例外やアンワインドを実装するためのプリミティブとして使うことができる。しかしそれ自体ははるかに低レベルなプリミティブであり、多種多様な別のパターンを可能にする。↩私はハードウェアの並列性よりも論理的な並行性に関心があるため、「タスク」という言葉を、他の任意の数の同様のシーケンスと時間的にインターリーブされうる、別個の線形な実行シーケンスを指すために使う。↩
この議論の多くは多くの言語や並行フレームワークに広く当てはまることを意図しているが、具体例としてPythonの例を使うことにする。他の利点に加え、Pythonはスレッドと協調的非同期の両方をサポートしており、複数のパラダイムを探求できるからだ。↩
おそらく最も単純な修正は、
asyncio.Eventを完全に取り除き、代わりにTaskオブジェクト自体をasyncio.gatherすることだ。ここではそれは簡単だが、「処理の単位」と「子タスク」の間に1対1の対応がない場合には、常にそれほど単純とは限らない。↩実際、Pythonの
asyncioには遍在的なキャンセル機構があることを付記しておく。↩
記事をランダムに読む
コメント
ログインしてコメントする