隙間を見つけて数え、塞ぐ — スレッド・asyncio・SQLite
目標
counter += 1がなぜアトミックでないかをバイトコードで確認し、スレッドの更新の喪失・座席の二重予約・asyncioのawaitの間の過剰な出金・SQLiteの2つの接続での更新の損失を、自分で作って数え、直します。イベントループを塞ぐ呼び出しのコストを「ループ遅延」として測り、GILの下で、スレッドとプロセスがCPU作業をどれだけ速く終えるかを測ります。
なぜ重要なのか
競合状態は、読んだ値と書く値の間の隙間で生じます。スレッドではバイトコード命令の間、asyncioではawait、データベースではSELECTとUPDATEの間が、その隙間です。隙間の一部だけをロックで包んだ修正は、エラーなしで動きながら、数字だけが間違うため、レビューではなかなか引っかかりません。そのためこのラボの採点ツールは、書いた数字だけを見るのではなく、あなたの関数を、用意された素材ではない入力で再実行して、基準の実装と突き合わせます。時間が関わる判定は、PodがCPUを共有しているので、広い範囲でしか見ません。
用意するもの
/opt/fixtures/svccs/concurrency/の下にあります。読み取り専用で使ってください。
params.json threads·loops(2·3단계) · sqlite_rounds(7단계) · lag_handlers·lag_block_s·lag_tick_s(6단계)
cpu_n·io_tasks·io_sleep_s(8단계)
bookings.json 좌석 예약 요청 목록 [{"user": "u001", "seat": "A01"}, …] — 같은 좌석 요청이 몰려 온다
accounts.json {"balances": {계좌: 잔액}, "withdrawals": [{"account": 계좌, "amount": 금액}, …]}
このコードブロックの韓国語の説明は、順に、params.jsonにはthreadsとloops(ステップ2・3)、sqlite_rounds(ステップ7)、lag_handlers・lag_block_s・lag_tick_s(ステップ6)、cpu_n・io_tasks・io_sleep_s(ステップ8)が入っていること、bookings.jsonは座席予約リクエストのリストで同じ座席へのリクエストが集中して来ること、accounts.jsonは口座ごとの残高のマップと、口座と金額を持つ出金のリストであること、を述べています。
ステップ
/root/svccs/concurrency/bytecode.jsonに、python(例: 「3.12」のようにメジャーとマイナーのバージョン)、ops、atomicを書きます。opsは、def bump(): global counter; counter += 1をdis.get_instructionsで展開した命令名のリストから、RESUMEとRETURNで始まるものを除いたものであり、atomicは、この増加がスレッド間でアトミックかどうか(true/false)です。/root/svccs/concurrency/conc.pyにunsafe_increment(threads, loops)を作成します。threads個のスレッドが、共有の値をloops回ずつ増やしますが、1回増やすたびに、読み取り →time.sleep(0)→ 読んだ値+1の書き込みの順序で行います。{"expected": threads×loops, "actual": 최종값, "lost": expected − actual}を返します(韓国語で「最終値」を意味する語です)。params.jsonのthreadsとloopsで実行した結果を、/root/svccs/concurrency/lost.jsonに、threads・loops・expected・actual・lostとして書きます。- 同じファイルに
locked_increment(threads, loops)を作成します。ステップ2と同じ、読み取り →time.sleep(0)→ 書き込みを行いますが、threading.Lockで防いで、lostが0になるようにします。譲る処理(time.sleep(0))は消しません。採点ツールが呼び出された回数を数えます。 - 同じファイルに
book(requests, safe)を作成します。リクエストごとに、スレッドが1つ、座席が空いているかの確認 →time.sleep(0)→ 予約の記録を行います。safeが真ならロックで防ぎます。{"requests", "seats"(서로 다른 좌석 수), "confirmed"(예약 성공 응답 수), "double_booked"(confirmed − 성공 응답에 나온 서로 다른 좌석 수)}を返します(コード内の韓国語は、順に「異なる座席の数」「予約成功の応答の数」「confirmed − 成功の応答に出てきた異なる座席の数」を意味する語です)。bookings.jsonで2回(ロックなし・ロックあり)実行して、/root/svccs/concurrency/seats.jsonに、requests・seats・unsafe・safeとして書きます(unsafeとsafeは、それぞれconfirmedとdouble_bookedを持つオブジェクトです)。 - 同じファイルに
run_withdrawals(balances, withdrawals, mode)を作成します。出金ごとにコルーチンを1つ、asyncio.gatherでリクエストのリストの順序どおりに起動し、各コルーチンは、残高 ≥ 金額かの確認 →await asyncio.sleep(0)→ 差し引きを行います(足りなければ拒否)。modeが"locked"なら、口座ごとのasyncio.Lockで防ぎます。{"approved", "rejected", "overdrawn_accounts"(최종 잔액이 음수인 계좌 수), "final"(계좌별 최종 잔액)}を返します(コード内の韓国語は、順に「最終残高が負の口座の数」「口座ごとの最終残高」を意味する語です)。accounts.jsonで2つのモードを実行して、/root/svccs/concurrency/overdraft.jsonに{"unsafe": …, "locked": …}の形式で書きます。 - 同じファイルに
async def handler(offload, block_s=0.2)とmeasure_lag(offload, handlers=3, tick=0.01)を作成します。handlerは、offloadが偽ならtime.sleep(block_s)をそのまま、真ならasyncio.to_threadへ移して呼び出します。measure_lagは、下の「ループ遅延のルール」のとおりに測ります。2つの場合を/root/svccs/concurrency/lag.jsonに、handlers・blocking・offloadedとして書きます(blockingとoffloadedは、measure_lagが返したオブジェクトそのままです)。 - 同じファイルに
sqlite_lost(path, rounds, mode)を作成します。下の「SQLiteのルール」のとおりに、2つの接続を1つのスレッドで交互に使います。params.jsonのsqlite_roundsで3つのモードを実行して、/root/svccs/concurrency/sqlite.jsonに、rounds・read_modify_write・atomic・deferred_txnとして書きます。 - 同じファイルに
cpu_bound(n)(0からn−1までのi*i % 7の合計)、gil_compare(n)、io_compare(tasks, sleep_s)を作成し、下の「GILのルール」のとおりに、/root/svccs/concurrency/gil.jsonを書きます。
ループ遅延のルール
ticker 작업: 반복해서 due = loop.time() + tick → await asyncio.sleep(tick) → loop.time() − due 를 기록
ticker 를 띄우고 tick×3 만큼 기다린 뒤, handler(offload) 를 handlers 개 gather 한다(걸린 시간 = elapsed_s)
gather 가 끝나면 ticker 를 멈추고 기다린다
돌려줄 것: max_lag_ms(기록 최댓값 × 1000, 소수 첫째 자리) · ticks(기록 개수) · elapsed_s(소수 셋째 자리)
このコードブロックの韓国語の説明は、順に、ticker作業は、繰り返し、due = loop.time() + tickを求め、await asyncio.sleep(tick)をして、loop.time() − dueを記録すること、tickerを起動して、tick×3だけ待ってから、handler(offload)をhandlers個gatherし(かかった時間がelapsed_s)、gatherが終わったらtickerを止めて待つこと、返す値はmax_lag_ms(記録の最大値×1000、小数第1位)・ticks(記録の個数)・elapsed_s(小数第3位)であること、を述べています。
SQLiteのルール
path 에 새 DB: CREATE TABLE counter (id INTEGER PRIMARY KEY, n INTEGER NOT NULL), 행 (1, 0)
연결 a, b 는 sqlite3.connect(path, timeout=0.05). 라운드마다
read_modify_write a 가 n 을 읽고, b 가 n 을 읽고, a 가 (읽은 값+1) 을 쓰고 commit, b 도 같게
atomic a 가 UPDATE counter SET n = n + 1 후 commit, b 도 같게
deferred_txn a·b 의 isolation_level = None. a: BEGIN·읽기, b: BEGIN·읽기, a: 읽은 값+1 쓰기,
b: 읽은 값+1 쓰기 후 COMMIT — OperationalError 면 busy +1 하고 ROLLBACK, 끝으로 a: COMMIT
돌려줄 것: final(최종 n) · expected(2×rounds) · busy · lost(expected − final − busy)
このコードブロックの韓国語の説明は、順に、pathに新しいDBを作り、counterテーブルと行(1, 0)を用意すること、接続aとbはtimeoutが0.05のsqlite3.connectで、ラウンドごとに、read_modify_writeではaがnを読み、bがnを読み、aが読んだ値+1を書いてcommitし、bも同じようにすること、atomicではaがUPDATE counter SET n = n + 1をしてcommitし、bも同じようにすること、deferred_txnではaとbのisolation_levelをNoneにして、a: BEGINと読み取り、b: BEGINと読み取り、a: 読んだ値+1を書き込み、b: 読んだ値+1を書き込んでCOMMITし、OperationalErrorならbusyを1増やしてROLLBACKし、最後にa: COMMITすること、返す値はfinal(最終的なn)・expected(2×rounds)・busy・lost(expected − final − busy)であること、を述べています。
GILのルール
gil_compare(n): cpu_bound(n) 을 한 번 돌려 데운 뒤, cpu_bound(n) 두 번을 순서대로 / 스레드 2개로 /
ProcessPoolExecutor(max_workers=2) 로 돌린다. 셋 다 벽시계(perf_counter)와 CPU 시간을 함께 잰다
→ cpu_seq_s · cpu_threads_s · cpu_procs_s (벽시계 초)
→ seq_cpu_s · threads_cpu_s (time.process_time 차이 — 이 프로세스 모든 스레드의 CPU 시간)
→ procs_cpu_s (os.times() 의 children_user + children_system 차이 — 끝난 자식 프로세스의 CPU 시간)
(여섯 값 모두 소수 셋째 자리)
io_compare(tasks, sleep_s): time.sleep(sleep_s) tasks 번을 순서대로 / 스레드 tasks 개로 → io_seq_s · io_threads_s
gil.json: gil_disabled_build(sysconfig.get_config_var("Py_GIL_DISABLED") 를 bool 로), 위 여덟 값,
thread_speedup = cpu_seq_s ÷ cpu_threads_s, proc_speedup = cpu_seq_s ÷ cpu_procs_s,
io_thread_speedup = io_seq_s ÷ io_threads_s,
thread_cores = threads_cpu_s ÷ cpu_threads_s, proc_cores = procs_cpu_s ÷ cpu_procs_s
(다섯 비율 모두 소수 둘째 자리 — cores 는 '그동안 평균 몇 개의 코어가 일했나' 입니다)
lost_updates_with_gil = lost.json 의 lost
このコードブロックの韓国語の説明は、順に、gil_compare(n)は、cpu_bound(n)を1回実行してウォームアップしたあと、cpu_bound(n)を、順番に2回、スレッド2つで、ProcessPoolExecutor(max_workers=2)で実行し、3つとも壁時計(perf_counter)とCPU時間を一緒に測ること、壁時計の秒はcpu_seq_s・cpu_threads_s・cpu_procs_sであり、seq_cpu_sとthreads_cpu_sはtime.process_timeの差(このプロセスのすべてのスレッドのCPU時間)、procs_cpu_sはos.times()のchildren_userとchildren_systemの差(終了した子プロセスのCPU時間)で、6つの値はすべて小数第3位であること、io_compare(tasks, sleep_s)は、time.sleep(sleep_s)をtasks回、順番に、およびスレッドtasks個で実行して、io_seq_sとio_threads_sを得ること、gil.jsonには、gil_disabled_build(sysconfig.get_config_varのPy_GIL_DISABLEDをboolにしたもの)、上の8つの値、thread_speedup = cpu_seq_s ÷ cpu_threads_s、proc_speedup = cpu_seq_s ÷ cpu_procs_s、io_thread_speedup = io_seq_s ÷ io_threads_s、thread_cores = threads_cpu_s ÷ cpu_threads_s、proc_cores = procs_cpu_s ÷ cpu_procs_sを入れ、5つの比率はすべて小数第2位であり、coresは、その間に平均で何個のコアが働いたかを表すこと、そしてlost_updates_with_gilはlost.jsonのlostであること、を述べています。
参考
- 採点ツールは
conc.pyを読み込みます。結果のJSONを作るコードは、if __name__ == "__main__":の下か、別のスクリプトに置いてください。 - スレッドの競合による損失の数は、OSが順序を決めるため、毎回異なります。採点ツールは、直したほうがちょうど0か、直していないほうが0より大きいかを確認します。asyncioとSQLiteの側の数字は、毎回同じです。
- よくある間違い: 書き込みだけをロックで包むこと、確認だけをロックで包んで行動を外で行うこと、コルーチンの中でthreading.Lockを握ったままawaitすること(ループが止まります)、to_threadに移したと書いてtime.sleepはそのままにしておくこと。
- 出力物はセッションが終わると消えます。必要なら別に保管してください。
counter += 1は命令がいくつか
counter += 1をdisで展開して、命令名のリストとアトミック性の判断を、/root/svccs/concurrency/bytecode.jsonにpython・ops・atomicとして書いてください。
dis.get_instructions(関数)は命令を1つずつ返し、各命令のopnameが名前です。読み取り(LOAD_)と書き込み(STORE_)が別の命令なら、GILが命令1つずつしか保証しないという事実と合わせて、何を意味するのかを考えてみてください。
スレッドで更新を失う
conc.pyにunsafe_increment(threads, loops)を作成し、params.jsonのthreadsとloopsで実行して、/root/svccs/concurrency/lost.jsonにthreads・loops・expected・actual・lostを書いてください。採点ツールは、別のスレッド数・反復数でも関数を呼び出します。
共有の値は、グローバル変数かdict 1つに置き、スレッドごとに、v = 値 → time.sleep(0) → 値 = v + 1を繰り返します。sleep(0)はGILを手放すので、別のスレッドが同じ値を読める隙間ができます。expectedはthreads × loopsです。
ロックで隙間全体を防ぐ
conc.pyにlocked_increment(threads, loops)を作成してください。読み取り → time.sleep(0) → 書き込みはそのままにして、threading.Lockでlostを0にします。採点ツールは、複数のスレッド数で呼び出し、time.sleepが呼ばれた回数も数えます。
ロックが何を包むべきかが、このステップのすべてです。書き込みの1行だけを包むと、ロックを待つ間に、すでに読んでおいた値が古くなり、順番に上書きされます。読み取りから書き込みまでを、ひとまとまりとして見てください。
座席の二重予約: 確認してから行動する
conc.pyにbook(requests, safe)を作成し、bookings.jsonで、ロックなし・ロックありで実行して、/root/svccs/concurrency/seats.jsonにrequests・seats・unsafe・safeを書いてください。採点ツールは、別のリクエストのリストでも呼び出します。
「空いているか」と「予約」の間に、別のスレッドが同じ座席を確認すると、両方が成功の応答を受け取ります。確認だけをロックの中に入れて、記録を外に置くと、隙間はそのままです。成功の応答のリストに座席を集めておけば、double_bookedを数えられます(list.appendは、FAQがアトミックだと記している操作です)。
awaitの間の過剰な出金
conc.pyにrun_withdrawals(balances, withdrawals, mode)を作成し、accounts.jsonでunsafeとlockedを実行して、/root/svccs/concurrency/overdraft.jsonに書いてください。採点ツールは、別の口座・出金のリストで2つのモードを呼び出し、基準の実装と正確に突き合わせます。
asyncioはスレッドが1つですが、awaitで別のコルーチンに順番が移ります。確認と差し引きの間にawaitがあると、その間に、同じ口座の別の出金が、同じ残高を見ます。ロックは、asyncio.Lockをasync withで、口座ごとに1つずつ用意し、確認から差し引きまでを包んでください。threading.Lockをコルーチンで握ったままawaitすると、ループが止まります。
塞ぐ呼び出し1つがすべてを止める
conc.pyにhandler(offload, block_s=0.2)とmeasure_lag(offload, handlers=3, tick=0.01)を作成し、2つの場合を/root/svccs/concurrency/lag.jsonにhandlers・blocking・offloadedとして書いてください。採点ツールは、あなたのhandlerを、自分の測定器で測り直します。
time.sleepが動いている間、ループのスレッドは、ほかのどのタスクも起こせません。handler 3つをgatherすると、3つが同じループの順番の中で続けて塞ぐので、遅延が合算されます。to_threadに移した場合も、作業自体は行う必要があるので、elapsed_sは0になりません。
SQLiteの2つの接続: 静かな損失と騒がしい失敗
conc.pyに、SQLiteのルールどおりsqlite_lost(path, rounds, mode)を作成し、sqlite_roundsで3つのモードを実行して、/root/svccs/concurrency/sqlite.jsonにrounds・read_modify_write・atomic・deferred_txnを書いてください。採点ツールは、別のラウンド数で3つのモードを呼び出します。
Pythonのsqlite3は、デフォルトの設定では、UPDATEの前でだけトランザクションを開きます。SELECTで読んだ値は保護されません。UPDATE ... SET n = n + 1は、読み取りと書き込みが1つの文です。自分でBEGINした2つの接続が、どちらも読んだあとで書こうとすると、片方は書き込みに昇格できず、「database is locked」を受け取ります。
GIL: 速くならないのに、安全でもない
conc.pyにcpu_bound・gil_compare・io_compareを作成し、GILのルールどおり/root/svccs/concurrency/gil.jsonを書いてください。lost_updates_with_gilは、ステップ2のlost.jsonのlostです。採点ツールは、io_compareを直接呼び出し、cpu_boundの値を突き合わせます。
CPUだけを使う関数は、GILをほとんど手放さないので、2つのスレッドが順番に動きます。プロセスはインタープリターが別なので、GILも別です。壁時計の時間は、隣のPodが混んでいると揺らぎますが、「CPU時間 ÷ 壁時計の時間」は、その間に平均で何個のコアが働いたかなので、GILの痕跡がよりはっきり出ます。time.sleepは、待っている間GILを手放すので、スレッドで重なります。