コース進捗 コース目次 24レッスン中 24件を公開中
Pythonの基礎
データとコレクション
信頼できるプログラムを作る
オブジェクトでモデル化する
プロフェッショナルなPython
上級Python
重なって進むことと同時実行は別物
並行性は、同じ時間帯に複数の仕事が進行中であることです。並列性は、通常は別々のCPUコアで、仕事が本当に同じ瞬間に動くことです。
- 逐次 一皿を終えて次へ
- 並行 待つ間に別の皿へ
- 並列 二人が今、同時に作る
ファイルや通信を待つ仕事は、多くの場合I/O待ち中心です。長いPure Python計算はCPU中心です。スレッドはブロッキング待ちを重ねるのに向くことがあり、プロセスは独立したCPU仕事を別コアで動かせます。予定調整と通信にもコストがあるため、高速化は測ってから判断します。
スレッドは同じ部屋を共有する
一つのプロセス内のスレッドは、同じアドレス空間、グローバル変数、ヒープ上のオブジェクトを共有します。呼び出しスタックは別ですが、二つの参照が同じ変更可能オブジェクトを指すことがあります。
標準的なCPython 3.12ビルドでは、GILにより、一つのプロセスでPythonバイトコードを実行するスレッドは同時に一つです。CPythonは多くのブロッキングI/O待機中にGILを解放し、一部のネイティブ拡張も解放します。そのためスレッドはI/O待機を重ねられますが、スレッドを増やすだけでは通常、Pure PythonのCPU処理を複数コアで並列実行できません。別のインタープリターや特殊ビルドでは異なる場合があります。
- 1プロセス 共有メモリの部屋
- スレッド 部屋の物を共有
- GIL Pythonバイトコードは1つ
- I/O待機 別スレッドが進める
GILは複数段階の不変条件を守りません。確認と変更の全体をLockで囲みます。
from concurrent.futures import ThreadPoolExecutor
from threading import Lock
balance = 3
balance_lock = Lock()
def spend_one():
global balance
with balance_lock:
if balance > 0:
balance -= 1
with ThreadPoolExecutor(max_workers=4) as pool:
list(pool.map(lambda _: spend_one(), range(10)))
print(balance)
0
ロックが「残高を0未満にしない」という確認と変更をまとめて守ります。クリティカルセクションは小さくします。共有が不要なら、ワーカーが値を返し、1つの所有者が結合する方が簡単です。
Executorが働き手を、Futureが結果を持つ
ThreadPoolExecutorは上限付きのワーカースレッドを再利用します。map()は完了順が違っても入力順に結果を返します。submit()はFutureを返します。
from concurrent.futures import ThreadPoolExecutor
def parse_minutes(text):
minutes = int(text)
if minutes < 0:
raise ValueError("negative minutes")
return minutes
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(parse_minutes, "-2")
try:
future.result()
except ValueError as error:
print(f"ワーカー失敗:{error}")
ワーカー失敗:negative minutes
result()は値を返すか、ワーカーの例外を調整役スレッドで再送出します。Futureは必ず観察します。as_completed()が見せる完了順は意図的に不安定です。識別子を結び付け、公開出力の前に並べ替えます。
プロセスは別々の部屋を使う
ワーカープロセスは、別のインタープリターとメモリを持ちます。標準CPythonでは独立したPython CPU処理を別コアで動かせます。入力、呼び出し可能オブジェクト、結果は通常pickleによる直列化で境界を越えます。
- プロセスA 専用メモリ
- 直列化データ 部屋の間を移動
- プロセスB 別コアを使える
プロセスタスクには、インポート可能なトップレベル関数を使います。引数と結果は適度な大きさで直列化可能にします。開いたファイル、ロック、接続、ラムダ、入れ子関数は不向きです。実行制御は移植可能な起動ガードの後ろへ置きます。
def main():
run_process_pool()
if __name__ == "__main__":
main()
ガードがなければ、spawnされた子がモジュールを読み込み、さらにプールを作る再帰が起きることがあります。
上限と停止の意味を正直に理解する
ワーカーを増やすと、ディスク競合、メモリ、サービス拒否、直列化コストが増える場合があります。max_workersは実際の容量から決めます。小さなCPU仕事はまとめ、終わらない入力には有限バッチや上限付き生成キューを使います。ワーカー数だけでは、投入済みFuture数を制限できません。
停止APIの意味は限定的です。
future.result(timeout=1)が制限するのは呼び出し側の待ち時間で、仕事は続くかもしれないfuture.cancel()が成功するのは、呼び出し開始前だけshutdown(wait=True)は新規受付を止め、実行中と保留中の仕事を待つshutdown(wait=False)は早く戻るが、呼び出しは続き、Python終了前には待たれるcancel_futures=Trueが取り消すのは保留中だけで、実行中は止めない
I/O API自身のタイムアウトと、協調的な停止合図を使います。任意の実行中スレッドをPythonから安全に強制巻き戻しすることはできません。
ミニ制作:ローカル読書ダイジェスト
reading_digest.pyを作ります。オフラインで標準ライブラリだけを使います。上限付きスレッドプールが一時ファイルを読み、上限付きプロセスプールが独立したテキスト分析をします。並べ替えにより出力は決定的です。
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
from pathlib import Path
from tempfile import TemporaryDirectory
def load_article(path):
return path.name, path.read_text(encoding="utf-8")
def analyze_article(record):
filename, text = record
words = text.split()
mentions = sum(
word.strip("。、,,.!?!?()() ").casefold() == "python"
for word in words
)
return filename, len(words), mentions
def build_digest(paths):
paths = list(paths)
if not paths:
return []
with ThreadPoolExecutor(max_workers=min(2, len(paths))) as pool:
loaded = list(pool.map(load_article, paths))
with ProcessPoolExecutor(max_workers=min(2, len(loaded))) as pool:
analyzed = list(pool.map(analyze_article, loaded))
return sorted(analyzed)
def main():
with TemporaryDirectory() as temporary:
directory = Path(temporary)
samples = {
"a.txt": "Python waits. Python reads files.",
"b.txt": "Processes analyze independent Python text.",
}
paths = []
for name, text in samples.items():
path = directory / name
path.write_text(text, encoding="utf-8")
paths.append(path)
items = build_digest(reversed(paths))
print("読書ダイジェスト")
for filename, words, mentions in items:
print(f"- {filename}:{words}語、python={mentions}")
with ThreadPoolExecutor(max_workers=1) as pool:
missing = pool.submit(load_article, directory / "missing.txt")
try:
missing.result()
except FileNotFoundError:
print("欠けたファイルを検出:True")
if __name__ == "__main__":
main()
ファイルとして保存し、python3 reading_digest.pyで実行します。
読書ダイジェスト
- a.txt:5語、python=2
- b.txt:5語、python=1
欠けたファイルを検出:True
このサンプルは小さすぎるため、高速化の証明にはなりません。プロセス起動、直列化、例外伝達、ワーカー上限、安定出力を確認する教材です。小さな入力なら逐次分析の方が速い場合があります。
3つのミッションと第23章への準備
- 共有をなくす。 ロック付き残高を、ワーカーがローカル件数を返し、中央で結合する形へ変えます。
- 順序を保つ。
submit()とas_completed()を使い、入力インデックスを保存してから、その順に表示します。 - 正直に測る。 大きめのローカル文書で逐次版とプロセス版を比べ、ワーカー数とデータ量も記録します。
次を説明、実行できれば準備完了です。
- 並行性と並列性を分けられる
- スレッドが共有するものと標準CPython 3.12のGIL境界を説明できる
- 複合的な不変条件をロックで守れる
- すべてのFutureを観察し、例外を表面化できる
- トップレベルの直列化可能なプロセスタスクをメインガードの後ろで使える
- リソース上限からワーカー数と投入数を制限できる
- timeout、cancel、shutdownが実行中処理を強制終了しないと分かる
- ダイジェストを実行し、並べ替え済み出力を再現できる
次章では、1つのイベントループがコルーチン、タスク、期限、協調的キャンセルで多くのI/O待機を調整します。