foundations · Stage 1
プロセス・スレッド・並行性の不変条件
隔離と共有の境界を選び、再現可能な競合を不変条件で修正する。
到達目標
プロセスの隔離資源とスレッドの共有資源を、障害伝播と通信経路を含めて図示できる
- 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
- 隔離と共有の失敗半径を比較する5分説明
100回以上の反復で競合を再現し、失敗時のinterleavingと破られた不変条件を記録できる
- 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
- 競合、deadlock、可視性の仮説を観測で切り分ける回答
同期方法を適用した後に正しさ、停止性、性能を再測定し、別の共有境界で選択を再比較できる
- 競合、deadlock、可視性の仮説を観測で切り分ける回答
- 未知のワーカー処理でプロセス分離とスレッド共有を再評価した設計
能力の進行
recognize
アドレス空間、ファイル記述子、実行状態、共有データを隔離または共有の観点で分類できる
証拠: 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
explain
操作のatomicity、可視性、順序性、deadlockを別の正しさ条件として説明できる
証拠: 隔離と共有の失敗半径を比較する5分説明
apply
失敗頻度を測れる競合実験を作り、明示した不変条件に対して同期方法を適用できる
証拠: 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
diagnose
イベント列、ロック取得順、進捗カウンタからdata race、deadlock、starvationを切り分けられる
証拠: 競合、deadlock、可視性の仮説を観測で切り分ける回答
lead
信頼境界、障害半径、通信量、デバッグ容易性から隔離単位のレビューを主導できる
証拠: 未知のワーカー処理でプロセス分離とスレッド共有を再評価した設計
なぜ重要か
並行処理では、同じコードが正しい順序で動くという暗黙の期待が壊れる。二つのworkerが在庫を一つずつ減らすだけでも、読取りと書込みが分離していれば更新を失う。低頻度の競合を「再現しないから直った」と扱うと、負荷が高い本番で不変条件を破る。
プロセスは独立したアドレス空間を持つ実行環境、スレッドはプロセス内の資源を共有する実行単位である。ただし、常にプロセスが重くスレッドが軽いという一軸では選べない。生成方式、OS、通信量、権限境界、クラッシュ時の巻込みが判断を変える。
メンタルモデル
並行性の正しさを四つに分ける。atomicityは操作が途中状態を見せないこと、可視性は一方の書込みを他方が観測できること、順序性は必要な前後関係が成立すること、停止性は処理がいつか進むことである。一つを満たしても他を自動では満たさない。Java SE 21では、あるvolatile書込みは同じfieldへの後続volatile読取りにhappens-beforeするが、stock--全体はread、compute、writeの複合操作である。
「data raceがない」と「業務上のrace conditionがない」も別である。個々のアクセスがatomicでも、在庫確認と確定の間に別操作が入れば業務不変条件を破り得る。守る範囲はコード行ではなく状態遷移で決める。
注記
図を読む際の補足情報です。
- この注記は旧図の読み順を保持する補助です。
- 隔離単位: process A、process B、または同一process内thread AとBを置く。
- 共有対象: memory、file、socket、database row、queueを列挙する。
- 操作分解: read、compute、write、publishをイベントへ分ける。
- 不変条件: 在庫は0以上、合計減算数と最終値が一致、同じ注文を二度確定しない。
- 同期関係: mutex、atomic operation、message passingのどれが前後関係を作るか示す。
- 停止性: lock順、待機資源、timeout、cancel経路を観測する。
- 回帰: 同じstress条件で違反頻度と所要時間を再測定する。
説明用のx=10で、Thread A/Bのどのinterleavingがlost updateを起こし、mutex後に何が変わるか。
- 境界と説明用fixture
同一processのThread A/Bが共有整数xを1ずつ減算する説明用scenario。初期値x = 10は普遍値ではない。
- 隔離単位・shared context
process A/Bまたは同一processのThread A/Bを置く。この説明用traceは同一processの二threadを使う。
順序: 0
lane: shared-context
- 共有対象・x
memory上の共有整数x = 10を説明用fixtureとし、実systemではfile、socket、database row、queueも列挙する。
順序: 1
lane: shared-context
- 隔離単位・shared context
- 同期なしのlost update
両threadが同じx = 10を読んで9を書き、二回減算の期待値x = 8を破る最小trace。
- 同期なし開始
説明用の値は初期値x = 10。Thread AとThread Bがmutexなしで1ずつ減算する。
順序: 2
lane: shared-context
- Thread A read
Thread Aが共有値x = 10を読む。
順序: 3
lane: thread-a
- Thread B read
Thread BもAのwrite前に同じ共有値x = 10を読む。
順序: 4
lane: thread-b
- Thread A compute
Thread Aはlocalに10 - 1 = 9を計算する。
順序: 5
lane: thread-a
- Thread B compute
Thread Bもlocalに10 - 1 = 9を計算する。
順序: 6
lane: thread-b
- Thread A write
Thread Aが共有値へx = 9を書き込む。
順序: 7
lane: thread-a
- Thread B write
Thread Bが同じx = 9を上書きし、Thread Aの減算を失わせる。
順序: 8
lane: thread-b
- lost update違反点
二回減算後の期待値 x = 8に対しactual x = 9。Thread B writeの完了時点でlost updateが観測可能になる。
順序: 9
lane: shared-context
- 同期なし開始
- mutexによる比較trace
同じ説明用fixtureをmutexで直列化し、二回目のreadがx = 9を見る比較。
- Thread A mutex取得
同じ説明用fixtureをx = 10へ戻し、Thread Aがmutexを取得する。
順序: 10
lane: thread-a
- Thread A read/compute/write
mutex内でThread Aがx = 10を読み、9を計算してx = 9を書く。
順序: 11
lane: thread-a
- Thread A mutex解放
Thread Aのwrite後にmutexを解放し、happens-beforeを作る。
順序: 12
lane: thread-a
- Thread B mutex取得
Thread BはAの解放後にmutexを取得する。
順序: 13
lane: thread-b
- Thread B read/compute/write
Thread Bは同期後のx = 9を読み、8を計算してx = 8を書く。
順序: 14
lane: thread-b
- 同期後の不変条件
同期後の x = 8は二回減算後の期待値 x = 8と一致し、lost updateはない。
順序: 15
lane: shared-context
- Thread A mutex取得
- 停止性と回帰
安全性だけでなくlock順、timeout、同条件stressを再確認する。
- 停止性
mutexのlock順、待機資源、timeout、cancel経路を観測する。
順序: 16
lane: shared-context
- 回帰
同期なしとmutexありを同じstress条件で反復し、違反頻度と所要時間を再測定する。
順序: 17
lane: shared-context
- 停止性
両threadがx=10を読む最小trace、期待値x=8に対する違反点、mutex後のx=8を順に説明できる。
動く例で考える
lost updateをイベント列で再現する
次をRaceLab.javaとして保存する。racyとsynchronizedは同じJava 21、初期値、worker数、反復数、出力契約を使う。
import java.io.BufferedReader;
import java.io.BufferedWriter;
import java.io.InputStreamReader;
import java.io.OutputStreamWriter;
import java.io.PrintWriter;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.util.Arrays;
import java.util.HashMap;
import java.util.Locale;
import java.util.Map;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
public final class RaceLab {
private static final int INITIAL = 1000;
private static final int PER_WORKER = 500;
private static final int TRIALS = 200;
private static final long TIMEOUT_MILLIS = 2000;
private static final long FAILURE_TIMEOUT_MILLIS = 100;
private static final Object LOCK = new Object();
private static volatile int stock;
private record Trial(
int finalStock,
int minimumStock,
long elapsedNs,
boolean stalled) {}
private record OwnerResult(
int finalStock,
int minimumStock,
int dedupePrevented,
int childExit,
long startupNs,
long ipcNs,
long bytesSent,
long bytesReceived) {}
private record FailureProbe(
boolean timedOut,
boolean forcedCleanup,
int childExit) {}
private static Trial runTrial(String mode) throws InterruptedException {
stock = INITIAL;
AtomicInteger minimumStock = new AtomicInteger(INITIAL);
long start = System.nanoTime();
if (mode.equals("serial")) {
stock -= PER_WORKER;
minimumStock.accumulateAndGet(stock, Math::min);
stock -= PER_WORKER;
minimumStock.accumulateAndGet(stock, Math::min);
return new Trial(
stock,
minimumStock.get(),
System.nanoTime() - start,
false);
}
CyclicBarrier readsCompleted = new CyclicBarrier(2);
AtomicReference<Throwable> failure = new AtomicReference<>();
Runnable operation = mode.equals("racy")
? () -> {
try {
int observed = stock;
readsCompleted.await();
stock = observed - PER_WORKER;
minimumStock.accumulateAndGet(stock, Math::min);
} catch (Throwable error) {
failure.compareAndSet(null, error);
}
}
: () -> {
synchronized (LOCK) {
int observed = stock;
stock = observed - PER_WORKER;
minimumStock.accumulateAndGet(stock, Math::min);
}
};
Thread first = new Thread(operation, "worker-A");
Thread second = new Thread(operation, "worker-B");
first.start();
second.start();
first.join(TIMEOUT_MILLIS);
second.join(TIMEOUT_MILLIS);
boolean stalled = first.isAlive() || second.isAlive();
if (stalled) {
first.interrupt();
second.interrupt();
first.join(TIMEOUT_MILLIS);
second.join(TIMEOUT_MILLIS);
}
if (failure.get() != null) {
throw new IllegalStateException("worker failed", failure.get());
}
return new Trial(
stock,
minimumStock.get(),
System.nanoTime() - start,
stalled);
}
private static long median(long[] sorted) {
int middle = sorted.length / 2;
return sorted.length % 2 == 1
? sorted[middle]
: sorted[middle - 1] / 2 + sorted[middle] / 2
+ (sorted[middle - 1] % 2 + sorted[middle] % 2) / 2;
}
private static void runThreadTrials(String mode)
throws InterruptedException {
int violations = 0;
int minimumStock = Integer.MAX_VALUE;
int finalStock = Integer.MIN_VALUE;
int stalled = 0;
long[] elapsedNsSamples = new long[TRIALS];
for (int trial = 0; trial < TRIALS; trial++) {
Trial result = runTrial(mode);
finalStock = result.finalStock();
if (finalStock != 0) {
violations++;
}
minimumStock = Math.min(
minimumStock,
result.minimumStock());
if (result.minimumStock() < 0) {
throw new AssertionError(
"negative intermediate stock in " + mode);
}
elapsedNsSamples[trial] = result.elapsedNs();
if (result.stalled()) {
stalled++;
}
}
long[] sortedElapsed = elapsedNsSamples.clone();
Arrays.sort(sortedElapsed);
long totalElapsedNs = Arrays.stream(elapsedNsSamples).sum();
long elapsedNsMedian = median(sortedElapsed);
long elapsedNsMin = sortedElapsed[0];
long elapsedNsMax = sortedElapsed[sortedElapsed.length - 1];
int expectedViolations = mode.equals("racy") ? TRIALS : 0;
if (violations != expectedViolations || stalled != 0) {
throw new AssertionError("unexpected trial result");
}
System.out.printf(
"mode=%s trials=%d expected=%d final=%d min=%d violations=%d stalled=%d elapsed_ns=%d%n",
mode, TRIALS, 0, finalStock, minimumStock,
violations, stalled, totalElapsedNs);
if (mode.equals("racy")) {
System.out.println(
"event=A:read(1000),B:read(1000),A:write(500),B:write(500)");
}
System.out.printf(
Locale.ROOT,
"{\"schema\":\"race-lab-v3\",\"mode\":\"%s\","
+ "\"trials\":%d,\"expected_final\":0,"
+ "\"final_stock\":%d,\"minimum_stock\":%d,"
+ "\"violations\":%d,\"stalled\":%d,"
+ "\"elapsed_ns_median\":%d,"
+ "\"elapsed_ns\":{\"samples\":%s,\"min\":%d,"
+ "\"median\":%d,\"max\":%d,\"total\":%d}}%n",
mode,
TRIALS,
finalStock,
minimumStock,
violations,
stalled,
elapsedNsMedian,
Arrays.toString(elapsedNsSamples),
elapsedNsMin,
elapsedNsMedian,
elapsedNsMax,
totalElapsedNs);
}
private static void runOwner() throws Exception {
int ownerStock = INITIAL;
Map<String, int[]> applied = new HashMap<>();
BufferedReader input = new BufferedReader(
new InputStreamReader(
System.in,
StandardCharsets.US_ASCII));
PrintWriter output = new PrintWriter(
new OutputStreamWriter(
System.out,
StandardCharsets.US_ASCII),
true);
output.println("READY");
String line;
while ((line = input.readLine()) != null) {
if (line.equals("GET")) {
output.printf("STATE,%d%n", ownerStock);
continue;
}
if (line.equals("STOP")) {
output.printf("STOPPED,%d%n", ownerStock);
return;
}
String[] fields = line.split(",", -1);
if (fields.length != 2
|| !fields[0].matches(
"[A-Za-z0-9][A-Za-z0-9._-]{0,63}")
|| !fields[1].matches("[0-9]+")) {
throw new IllegalArgumentException("invalid owner request");
}
int decrement = Integer.parseInt(fields[1]);
if (decrement < 0 || decrement > INITIAL) {
throw new IllegalArgumentException("invalid decrement");
}
int[] previous = applied.get(fields[0]);
if (previous != null) {
if (previous[0] != decrement) {
throw new IllegalArgumentException(
"request ID reused with another decrement");
}
output.printf(
"RESULT,%s,%d,true%n",
fields[0],
previous[1]);
continue;
}
if (ownerStock - decrement < 0) {
throw new IllegalStateException("negative stock");
}
ownerStock -= decrement;
applied.put(
fields[0],
new int[] {decrement, ownerStock});
output.printf(
"RESULT,%s,%d,false%n",
fields[0],
ownerStock);
}
throw new IllegalStateException("owner input ended before STOP");
}
private static void runStalledOwner() throws InterruptedException {
System.out.println("READY");
System.out.flush();
while (true) {
Thread.sleep(TIMEOUT_MILLIS);
}
}
private static ProcessBuilder ownerProcess(String mode) {
String javaBinary = Path.of(
System.getProperty("java.home"),
"bin",
"java").toString();
return new ProcessBuilder(
javaBinary,
"-cp",
System.getProperty("java.class.path"),
RaceLab.class.getName(),
mode);
}
private static String readLineWithTimeout(
BufferedReader reader,
ExecutorService readerExecutor) throws Exception {
Future<String> future = readerExecutor.submit(reader::readLine);
try {
String line = future.get(
TIMEOUT_MILLIS,
TimeUnit.MILLISECONDS);
if (line == null) {
throw new IllegalStateException("owner output ended");
}
return line;
} catch (TimeoutException error) {
future.cancel(true);
throw new IllegalStateException("owner response timeout", error);
}
}
private static String exchange(
BufferedWriter input,
BufferedReader output,
ExecutorService readerExecutor,
String request) throws Exception {
input.write(request);
input.newLine();
input.flush();
return readLineWithTimeout(output, readerExecutor);
}
private static OwnerResult exerciseOwner() throws Exception {
Process child = null;
ExecutorService readerExecutor =
Executors.newSingleThreadExecutor();
long startupStart = System.nanoTime();
try {
child = ownerProcess("owner").start();
BufferedWriter input = new BufferedWriter(
new OutputStreamWriter(
child.getOutputStream(),
StandardCharsets.US_ASCII));
BufferedReader output = new BufferedReader(
new InputStreamReader(
child.getInputStream(),
StandardCharsets.US_ASCII));
String ready = readLineWithTimeout(output, readerExecutor);
long startupNs = System.nanoTime() - startupStart;
if (!ready.equals("READY")) {
throw new IllegalStateException("owner did not become ready");
}
String[] requests = {"A,500", "A,500", "B,500", "GET", "STOP"};
String[] expected = {
"RESULT,A,500,false",
"RESULT,A,500,true",
"RESULT,B,0,false",
"STATE,0",
"STOPPED,0",
};
long bytesSent = 0;
long bytesReceived = ready.getBytes(
StandardCharsets.US_ASCII).length + 1;
int minimumStock = INITIAL;
int dedupePrevented = 0;
long ipcStart = System.nanoTime();
for (int index = 0; index < requests.length; index++) {
String response = exchange(
input,
output,
readerExecutor,
requests[index]);
if (!response.equals(expected[index])) {
throw new IllegalStateException(
"unexpected owner response: " + response);
}
bytesSent += requests[index].getBytes(
StandardCharsets.US_ASCII).length + 1;
bytesReceived += response.getBytes(
StandardCharsets.US_ASCII).length + 1;
if (response.endsWith(",true")) {
dedupePrevented++;
}
String[] fields = response.split(",");
if (fields.length >= 2
&& fields[fields.length - 1].matches("[0-9]+")) {
minimumStock = Math.min(
minimumStock,
Integer.parseInt(fields[fields.length - 1]));
} else if (fields.length >= 3
&& fields[fields.length - 2].matches("[0-9]+")) {
minimumStock = Math.min(
minimumStock,
Integer.parseInt(fields[fields.length - 2]));
}
}
long ipcNs = System.nanoTime() - ipcStart;
if (!child.waitFor(TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) {
throw new IllegalStateException("owner exit timeout");
}
int childExit = child.exitValue();
if (childExit != 0) {
throw new IllegalStateException(
"owner exited " + childExit);
}
return new OwnerResult(
0,
minimumStock,
dedupePrevented,
childExit,
startupNs,
ipcNs,
bytesSent,
bytesReceived);
} finally {
if (child != null && child.isAlive()) {
child.destroyForcibly();
if (!child.waitFor(
TIMEOUT_MILLIS,
TimeUnit.MILLISECONDS)) {
throw new IllegalStateException(
"owner force cleanup failed");
}
}
readerExecutor.shutdownNow();
if (!readerExecutor.awaitTermination(
TIMEOUT_MILLIS,
TimeUnit.MILLISECONDS)) {
throw new IllegalStateException(
"owner reader cleanup failed");
}
}
}
private static FailureProbe exerciseFailureRadius()
throws Exception {
Process child = null;
BufferedReader output = null;
ExecutorService readerExecutor =
Executors.newSingleThreadExecutor();
try {
child = ownerProcess("owner-stall").start();
output = new BufferedReader(
new InputStreamReader(
child.getInputStream(),
StandardCharsets.US_ASCII));
if (!readLineWithTimeout(output, readerExecutor).equals("READY")) {
throw new IllegalStateException(
"stalled owner did not become ready");
}
boolean timedOut = !child.waitFor(
FAILURE_TIMEOUT_MILLIS,
TimeUnit.MILLISECONDS);
boolean forcedCleanup = false;
if (timedOut) {
forcedCleanup = true;
child.destroyForcibly();
}
if (!child.waitFor(TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) {
throw new IllegalStateException(
"stalled owner force cleanup failed");
}
return new FailureProbe(
timedOut,
forcedCleanup,
child.exitValue());
} finally {
if (output != null) {
output.close();
}
if (child != null && child.isAlive()) {
child.destroyForcibly();
child.waitFor(TIMEOUT_MILLIS, TimeUnit.MILLISECONDS);
}
readerExecutor.shutdownNow();
if (!readerExecutor.awaitTermination(
TIMEOUT_MILLIS,
TimeUnit.MILLISECONDS)) {
throw new IllegalStateException(
"stalled owner reader cleanup failed");
}
}
}
private static void runProcessMessage() throws Exception {
OwnerResult owner = exerciseOwner();
FailureProbe failure = exerciseFailureRadius();
if (owner.minimumStock() < 0
|| owner.finalStock() != 0
|| owner.dedupePrevented() != 1
|| !failure.timedOut()
|| !failure.forcedCleanup()) {
throw new AssertionError("unexpected process-message result");
}
long elapsedNs = owner.startupNs() + owner.ipcNs();
System.out.printf(
"mode=process-message trials=1 expected=%d final=%d min=%d violations=%d stalled=%d elapsed_ns=%d%n",
0,
owner.finalStock(),
owner.minimumStock(),
0,
0,
elapsedNs);
System.out.printf(
Locale.ROOT,
"{\"schema\":\"race-lab-v3\","
+ "\"mode\":\"process-message\","
+ "\"trial_scope\":"
+ "\"single-owner-session-not-200-thread-trials\","
+ "\"protocol\":[\"A,500\",\"A,500\","
+ "\"B,500\",\"GET\",\"STOP\"],"
+ "\"final_stock\":%d,\"minimum_stock\":%d,"
+ "\"dedupe_prevented\":%d,\"child_exit\":%d,"
+ "\"startup_ns\":%d,\"ipc_ns\":%d,"
+ "\"bytes_sent\":%d,\"bytes_received\":%d,"
+ "\"failure_probe_timed_out\":%b,"
+ "\"forced_cleanup\":%b,"
+ "\"failure_probe_exit\":%d,"
+ "\"parent_continued\":true,"
+ "\"failure_radius\":\"owner-child-process\"}%n",
owner.finalStock(),
owner.minimumStock(),
owner.dedupePrevented(),
owner.childExit(),
owner.startupNs(),
owner.ipcNs(),
owner.bytesSent(),
owner.bytesReceived(),
failure.timedOut(),
failure.forcedCleanup(),
failure.childExit());
}
public static void main(String[] args) throws Exception {
if (args.length != 1) {
throw new IllegalArgumentException(
"usage: java RaceLab "
+ "racy|synchronized|serial|process-message");
}
if (args[0].equals("owner")) {
runOwner();
} else if (args[0].equals("owner-stall")) {
runStalledOwner();
} else if (args[0].equals("process-message")) {
runProcessMessage();
} else if (args[0].equals("racy")
|| args[0].equals("synchronized")
|| args[0].equals("serial")) {
runThreadTrials(args[0]);
} else {
throw new IllegalArgumentException(
"usage: java RaceLab "
+ "racy|synchronized|serial|process-message");
}
}
}
java --version
javac RaceLab.java
java RaceLab racy | tee race-before.txt
java RaceLab synchronized | tee race-after.txt
java RaceLab serial | tee serial.txt
java RaceLab process-message | tee process-message.txt
/usr/bin/time -p java RaceLab synchronized
/usr/bin/time -p java RaceLab serial
/usr/bin/time -p java RaceLab process-message
thread三方式の出力契約はmode=... trials=200 expected=0 final=... min=... violations=... stalled=... elapsed_ns=...と、その直後のJSONである。JSONは全200件のelapsed_ns.samples、min、median、max、totalを保持する。各TrialのminimumStockは書込み直後に更新するため最終値の別名ではなく、racy、synchronized、serialの全方式で試行中の負数がないことをassertする。racy版はCyclicBarrierで両workerのreadを先に完了させるため、schedulerやyieldの偶然ではなく全200試行で最小lost updateを構成する。synchronized版とserial版は同じ1000件相当の固定fixtureで200試行0違反を要求し、join timeoutとisAliveで停止性を検査する。なお、一般の確率的raceでは一回の非再現を修正証拠にしない。
- 前提
- Java 21、在庫1000、worker AとBが各500件分を減らす。不変条件はexpected final 0、負数なし、stalled 0。
- 入力
- 同じ1000件の減算fixtureをracy、synchronized、serialへ渡し、各200試行する。
- 操作
- racy版では両read後にbarrierを解放して二つのwriteを行い、同期版ではread-subtract-write全体を同じmonitorで囲む。
- 観測
- 最小event列はA:read(1000)、B:read(1000)、A:write(500)、B:write(500)で、expected 0に対しfinal 500となる。thread三方式ではfinal、試行中minimum、violations、stalled、200個のelapsed分布を比較する。
- 結論
- synchronizedでread-modify-write全体を囲む。修正後200試行0違反と、monitor unlockから後続lockへのhappens-beforeを組にして修正根拠にする。
性能値はSystem.nanoTimeと/usr/bin/time -pの実測だけを報告し、例示値を捏造しない。process-messageは共有memory版の200試行へ混ぜず、single-owner-session-not-200-thread-trialsという別の測定単位にする。RaceLab自身がowner子プロセスを起動し、標準入力protocolをrequestId,decrement、GET、STOPに固定する。A,500、同じA,500、B,500、GET、STOPを送り、重複Aを二重減算しないこと、final/minimum 0、送受信byte、正常child exit、process起動時間とIPC時間を機械可読JSONへ分離して記録する。さらに別の停止用childを100msでtimeoutさせ、destroyForcibly、exit確認後も親が結果を出せた事実を、共有heapを巻き込まない障害半径の実測証拠にする。
ラボ成果物: 「競合を再現し不変条件で修正した実験記録」には、失敗頻度、最小interleaving、同期前後のコード差、停止性、性能分布を含める。
注記
図を読む際の補足情報です。
- x=10から二回減算する説明用fixtureであり、schedulerの普遍的な実行順を表さない。
二つのthreadが同じ古い値を読んだ後、どこで不変条件を破り、mutexでどう回復するか。
- 初期状態: A read: primary event a-read: Thread Aがx=10を読む。
- B read: primary event b-read: BもA write前のx=10を読む。
- A compute: primary event a-compute: localで10-1=9。
- B compute: primary event b-compute: localで10-1=9。
- A write: primary event a-write: x=9を書く。
- B write / lost update: primary events b-write/lost-update-violation: 9を上書きし、期待値x=8に対してactual x=9。
- A lock: primary event a-lock: fixtureを10へ戻してmutex取得。
- A locked read/compute/write: primary event a-locked-update: mutex内で10→9。
- A unlock: primary event a-unlock: write後にmutex解放。
- B lock: primary event b-lock: A unlock後にmutex取得。
- B locked read/compute/write: primary event b-locked-update: 同期後の9を読み9→8。
- invariant complete: primary event synchronized-invariant: 期待値x=8とactual x=8。lock/unlock単位の完全trace。
| イベント | 開始 | 終了 | 判定 | 理由 |
|---|---|---|---|---|
| next | A read | B read | allowed | — |
| next | B read | A compute | allowed | — |
| next | A compute | B compute | allowed | — |
| next | B compute | A write | allowed | — |
| next | A write | B write / lost update | allowed | — |
| next | B write / lost update | A lock | allowed | — |
| next | A lock | A locked read/compute/write | allowed | — |
| next | A locked read/compute/write | A unlock | allowed | — |
| next | A unlock | B lock | allowed | — |
| next | B lock | B locked read/compute/write | allowed | — |
| next | B locked read/compute/write | invariant complete | allowed | — |
| next | invariant complete | B write / lost update | rejected | read/compute/write全体を同じmutex境界で保護する。 |
read-old-valueからlost-updateを経て、mutexで二回の減算を直列化したlocked-completeまでを再現する。
- A read: read/compute/write trace 1: Aがx=10をread。; 条件 常時; node
read-old-value; edge なし - B read: trace 2: Bもx=10をread。; 条件 常時; node
b-read-old-value; edgestep-01 - A compute: trace 3: Aがlocal 9をcompute。; 条件 常時; node
a-compute; edgestep-02 - B compute: trace 4: Bもlocal 9をcompute。; 条件 常時; node
b-compute; edgestep-03 - A write: trace 5: Aがx=9をwrite。; 条件 常時; node
a-write; edgestep-04 - B write: trace 6: Bが9を上書き。期待値x=8、actual x=9でinvariant違反。; 条件 常時; node
lost-update; edgestep-05 - A lock: trace 7: fixtureを10へ戻しAがlock。; 条件 常時; node
lock-acquired; edgestep-06 - A locked update: trace 8: mutex内でAがread/compute/writeし10→9。; 条件 常時; node
a-locked-write; edgestep-07 - A unlock: trace 9: Aがwrite後にunlock。; 条件 常時; node
unlock; edgestep-08 - B lock: trace 10: Bがlockし、同期後の9を観測。; 条件 常時; node
b-lock-acquired; edgestep-09 - B locked update: trace 11: Bがread/compute/writeし9→8。; 条件 常時; node
b-locked-write; edgestep-10 - invariant complete: trace 12: 期待値x=8、actual x=8。lock/unlockを含む完全trace。; 条件 常時; node
locked-complete; edgestep-11
| イベント | 開始 | 終了 | 条件 |
|---|---|---|---|
| next | read-old-value | b-read-old-value | 常時 |
| timer | read-old-value | b-read-old-value | 常時 |
| next | b-read-old-value | a-compute | 常時 |
| timer | b-read-old-value | a-compute | 常時 |
| previous | b-read-old-value | read-old-value | 常時 |
| reset | b-read-old-value | read-old-value | 常時 |
| next | a-compute | b-compute | 常時 |
| timer | a-compute | b-compute | 常時 |
| previous | a-compute | b-read-old-value | 常時 |
| reset | a-compute | read-old-value | 常時 |
| next | b-compute | a-write | 常時 |
| timer | b-compute | a-write | 常時 |
| previous | b-compute | a-compute | 常時 |
| reset | b-compute | read-old-value | 常時 |
| next | a-write | lost-update | 常時 |
| timer | a-write | lost-update | 常時 |
| previous | a-write | b-compute | 常時 |
| reset | a-write | read-old-value | 常時 |
| next | lost-update | lock-acquired | 常時 |
| timer | lost-update | lock-acquired | 常時 |
| previous | lost-update | a-write | 常時 |
| reset | lost-update | read-old-value | 常時 |
| next | lock-acquired | a-locked-write | 常時 |
| timer | lock-acquired | a-locked-write | 常時 |
| previous | lock-acquired | lost-update | 常時 |
| reset | lock-acquired | read-old-value | 常時 |
| next | a-locked-write | unlock | 常時 |
| timer | a-locked-write | unlock | 常時 |
| previous | a-locked-write | lock-acquired | 常時 |
| reset | a-locked-write | read-old-value | 常時 |
| next | unlock | b-lock-acquired | 常時 |
| timer | unlock | b-lock-acquired | 常時 |
| previous | unlock | a-locked-write | 常時 |
| reset | unlock | read-old-value | 常時 |
| next | b-lock-acquired | b-locked-write | 常時 |
| timer | b-lock-acquired | b-locked-write | 常時 |
| previous | b-lock-acquired | unlock | 常時 |
| reset | b-lock-acquired | read-old-value | 常時 |
| next | b-locked-write | locked-complete | 常時 |
| timer | b-locked-write | locked-complete | 常時 |
| previous | b-locked-write | b-lock-acquired | 常時 |
| reset | b-locked-write | read-old-value | 常時 |
| previous | locked-complete | b-locked-write | 常時 |
| reset | locked-complete | read-old-value | 常時 |
| 結果 | 状態 |
|---|---|
| 同期なしtraceでは期待値8とactual 9の差を示す。 | lost-update |
| mutex traceでは最終値8と停止性の両方を検証する。 | locked-complete |
現在の状態: A read — read/compute/write trace 1: Aがx=10をread。
このモデルは例示的かつ決定的であり、実システムの完全な再現ではありません。
トレードオフと失敗モード
| 主要制約 | 候補 | 利点 | 主要リスク |
|---|---|---|---|
| 大量の小さい共有状態 | 同一processのthreads | 低い通信変換費用 | race、deadlock、クラッシュ伝播 |
| 信頼できない処理または権限分離 | 別process | address spaceと権限の隔離 | IPC契約、copy、部分失敗 |
| 所有権を一箇所へ集約可能 | message passing | 共有mutable stateの縮小 | mailbox滞留、順序、backpressure |
もっともらしい誤診と反証
- 誤診: 「二workerが止まったのでdeadlock」。反証: lock取得と待機の循環を記録する。I/O待ちなら外部完了で進み、starvationなら他workerだけ進捗する。停止だけでは循環待ちを証明しない。
- 誤診: 「Java volatileを付けて100回通ったのでrace修正済み」。反証: read-modify-writeをeventへ分解し、同じinterleavingがJava Language Specification上可能か確認する。volatileがvisibilityとorderingのhappens-beforeを作っても複合incrementのatomicityを保証しない。
lock範囲を広げると不変条件は守りやすいが並列性が下がり、複数lockは並列性を上げる一方で順序規約が必要になる。timeoutでlockを諦める場合も、途中状態のrollbackまたは再試行契約が要る。
知識チェック
- 整数代入がatomicならcounter増加もatomicか。通常は同じではない。増加はread、compute、writeの複合操作になり得る。
- process間ならraceは起きないか。共有memory、file、databaseを介せば起きる。address space隔離は全資源の隔離ではない。
- lockを一つにすればdeadlockは消えるか。そのlock同士の循環は避けやすいが、callback、I/O、再入、別資源待ちを含む停止性は別に確認する。
5分teach-back
lost updateの四イベントを描き、不変条件、atomicity、可視性、停止性を区別する。mutex案とmessage passing案の障害半径を一つずつ述べる。
未知へのtransfer
プロセス分離とスレッド共有の再比較を、画像変換pluginまたは機密データ集計へ適用する。CPU時間だけでなく、権限、入力の信頼度、IPC量、worker crash後の再実行を評価し直す。
出典と次の学習
Java Language Specification SE 21 Chapter 17はthread、monitor、volatile、happens-beforeの正準、POSIX.1-2024のpthread_createとforkは生成後に共有または複製される実行文脈、Linux kernel memory barriers文書はkernel実装上の順序保証、MicrosoftのProcesses and ThreadsはWindows上の資源境界を確認するために使う。異なる言語・OSの契約を混ぜず、正確なURLはmetadataのsourcesに置いた。
次はネットワークの部分失敗を扱い、process外の待機と期限を時系列で診断する。1日後はinterleaving、7日後はdeadlock反証、30日後は実コードの不変条件、90日後は隔離単位の判断を再レビューする。
実践ラボ
在庫更新のlost updateを再現して不変条件で直す
提出成果物: 競合を再現し不変条件で修正した実験記録
- Java 21で教材のRaceLab.javaをコンパイルし、volatile intへの非atomicなread-modify-writeを200試行する
- 最終値は0であり負数にならないという不変条件を先に明文化し、違反回数を記録する
- 失敗試行の読取り、計算、書込みinterleavingを最小のイベント列として再構成する
- Java synchronizedでread-modify-write全体を修正し、同じ200試行で不変条件と停止性を確認する
- 直列、同期ありスレッド、プロセス間メッセージの所要時間と障害半径を比較する
説明して理解を確かめる
5分で、プロセスとスレッドの共有境界、lost updateのinterleaving、Java volatileがvisibilityとorderingを扱っても複合incrementのatomicityを得られない理由を説明する。
アセスメント
問い: 二つのworkerが停止した。『deadlockだ』という診断を確定するために必要な証拠と、別のもっともらしい原因を示す。
期待する証拠: 待機資源とロック取得順の循環、進捗観測、I/O待ちまたはstarvationとの反証
問い: 共有変数へvolatileを付けたら再現しなくなった。修正完了と判断できない理由は何か。
期待する証拠: Java SE 21のvolatile happens-before、複合操作のatomicity、確率的な一回の非再現を分けた説明
別問題へ転用する
プロセス分離とスレッド共有の再比較
復習スケジュール
- 1日後
atomicity、可視性、順序性のうち、今回の不変条件に必要なのはどれか
- 7日後
deadlock説を反証する最小のイベント証拠は何か
- 30日後
隔離を強めると通信と復旧の設計はどう変わるか
- 90日後
atomicity、可視性、順序性のうち、今回の不変条件に必要なのはどれか
評価ルーブリック
| 観点 | 未達 | 発展途上 | 熟達 | 卓越 |
|---|---|---|---|---|
| technical-correctness | volatile、barrier、mutexを同義に扱うか、lost updateの操作列を説明できない | 競合は説明するが、atomicityと可視性またはdeadlockとstarvationを混同する | 不変条件、interleaving、happens-before、同期範囲を一貫して説明する | 停止性とメモリ順序の反例を含め、実装とOS境界の双方を正確に扱う |
| judgment | プロセスは重い、スレッドは軽いという一軸だけで選択する | 性能と隔離を比較するが、通信、復旧、信頼境界が抜ける | 障害半径、共有量、通信費、観測性を制約に合わせて比較する | 権限分離と段階的な並列度を設計し、負荷変化時の撤回条件を示す |
| evidence | 一度再現しなかったことを修正の証拠にする | 反復試験はあるが、失敗時のイベント列または停止性の観測がない | 100試行以上の失敗頻度、最小interleaving、修正後の同条件回帰を示す | スケジューリングを変えたstress試験と性能分布で正しさと費用を検証する |
| communication | raceがあったという結論だけで、不変条件と再現方法がない | コード変更は示すが、なぜ修正になるかをイベント列で説明しない | 再現条件、不変条件、原因、同期範囲、回帰結果をレビュー可能に結ぶ | 運用者へ検知方法を、実装者へ同期規約を同じ証拠から説明できる |
出典
以下の外部資料は利用者が選択したときだけ開きます。
- Chapter 17. Threads and Locks (standard)
- pthread_create — thread creation (standard)
- fork — create a new process (standard)
- Linux kernel memory barriers (primary)
- About Processes and Threads (primary)