foundations · Stage 1

プロセス・スレッド・並行性の不変条件

隔離と共有の境界を選び、再現可能な競合を不変条件で修正する。

学習時間
270分
難易度
foundation
更新日
2026-07-30
到達証拠
成果物・説明・判断根拠・転用

到達目標

  1. プロセスの隔離資源とスレッドの共有資源を、障害伝播と通信経路を含めて図示できる

    • 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
    • 隔離と共有の失敗半径を比較する5分説明
  2. 100回以上の反復で競合を再現し、失敗時のinterleavingと破られた不変条件を記録できる

    • 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録
    • 競合、deadlock、可視性の仮説を観測で切り分ける回答
  3. 同期方法を適用した後に正しさ、停止性、性能を再測定し、別の共有境界で選択を再比較できる

    • 競合、deadlock、可視性の仮説を観測で切り分ける回答
    • 未知のワーカー処理でプロセス分離とスレッド共有を再評価した設計

能力の進行

  1. recognize

    アドレス空間、ファイル記述子、実行状態、共有データを隔離または共有の観点で分類できる

    証拠: 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録

  2. explain

    操作のatomicity、可視性、順序性、deadlockを別の正しさ条件として説明できる

    証拠: 隔離と共有の失敗半径を比較する5分説明

  3. apply

    失敗頻度を測れる競合実験を作り、明示した不変条件に対して同期方法を適用できる

    証拠: 再現条件、interleaving、不変条件、修正、回帰測定を含む実験記録

  4. diagnose

    イベント列、ロック取得順、進捗カウンタからdata race、deadlock、starvationを切り分けられる

    証拠: 競合、deadlock、可視性の仮説を観測で切り分ける回答

  5. lead

    信頼境界、障害半径、通信量、デバッグ容易性から隔離単位のレビューを主導できる

    証拠: 未知のワーカー処理でプロセス分離とスレッド共有を再評価した設計

なぜ重要か

並行処理では、同じコードが正しい順序で動くという暗黙の期待が壊れる。二つのworkerが在庫を一つずつ減らすだけでも、読取りと書込みが分離していれば更新を失う。低頻度の競合を「再現しないから直った」と扱うと、負荷が高い本番で不変条件を破る。

プロセスは独立したアドレス空間を持つ実行環境、スレッドはプロセス内の資源を共有する実行単位である。ただし、常にプロセスが重くスレッドが軽いという一軸では選べない。生成方式、OS、通信量、権限境界、クラッシュ時の巻込みが判断を変える。

メンタルモデル

並行性の正しさを四つに分ける。atomicityは操作が途中状態を見せないこと、可視性は一方の書込みを他方が観測できること、順序性は必要な前後関係が成立すること、停止性は処理がいつか進むことである。一つを満たしても他を自動では満たさない。Java SE 21では、あるvolatile書込みは同じfieldへの後続volatile読取りにhappens-beforeするが、stock--全体はread、compute、writeの複合操作である。

「data raceがない」と「業務上のrace conditionがない」も別である。個々のアクセスがatomicでも、在庫確認と確定の間に別操作が入れば業務不変条件を破り得る。守る範囲はコード行ではなく状態遷移で決める。

共有境界から不変条件を守る診断経路

注記

図を読む際の補足情報です。

  1. この注記は旧図の読み順を保持する補助です。
  2. 隔離単位: process A、process B、または同一process内thread AとBを置く。
  3. 共有対象: memory、file、socket、database row、queueを列挙する。
  4. 操作分解: read、compute、write、publishをイベントへ分ける。
  5. 不変条件: 在庫は0以上、合計減算数と最終値が一致、同じ注文を二度確定しない。
  6. 同期関係: mutex、atomic operation、message passingのどれが前後関係を作るか示す。
  7. 停止性: lock順、待機資源、timeout、cancel経路を観測する。
  8. 回帰: 同じstress条件で違反頻度と所要時間を再測定する。

説明用のx=10で、Thread A/Bのどのinterleavingがlost updateを起こし、mutex後に何が変わるか。

  1. 境界と説明用fixture

    同一processのThread A/Bが共有整数xを1ずつ減算する説明用scenario。初期値x = 10は普遍値ではない。

    1. 隔離単位・shared context

      process A/Bまたは同一processのThread A/Bを置く。この説明用traceは同一processの二threadを使う。

      順序: 0

      lane: shared-context

    2. 共有対象・x

      memory上の共有整数x = 10を説明用fixtureとし、実systemではfile、socket、database row、queueも列挙する。

      順序: 1

      lane: shared-context

  2. 同期なしのlost update

    両threadが同じx = 10を読んで9を書き、二回減算の期待値x = 8を破る最小trace。

    1. 同期なし開始

      説明用の値は初期値x = 10。Thread AとThread Bがmutexなしで1ずつ減算する。

      順序: 2

      lane: shared-context

    2. Thread A read

      Thread Aが共有値x = 10を読む。

      順序: 3

      lane: thread-a

    3. Thread B read

      Thread BもAのwrite前に同じ共有値x = 10を読む。

      順序: 4

      lane: thread-b

    4. Thread A compute

      Thread Aはlocalに10 - 1 = 9を計算する。

      順序: 5

      lane: thread-a

    5. Thread B compute

      Thread Bもlocalに10 - 1 = 9を計算する。

      順序: 6

      lane: thread-b

    6. Thread A write

      Thread Aが共有値へx = 9を書き込む。

      順序: 7

      lane: thread-a

    7. Thread B write

      Thread Bが同じx = 9を上書きし、Thread Aの減算を失わせる。

      順序: 8

      lane: thread-b

    8. lost update違反点

      二回減算後の期待値 x = 8に対しactual x = 9。Thread B writeの完了時点でlost updateが観測可能になる。

      順序: 9

      lane: shared-context

  3. mutexによる比較trace

    同じ説明用fixtureをmutexで直列化し、二回目のreadがx = 9を見る比較。

    1. Thread A mutex取得

      同じ説明用fixtureをx = 10へ戻し、Thread Aがmutexを取得する。

      順序: 10

      lane: thread-a

    2. Thread A read/compute/write

      mutex内でThread Aがx = 10を読み、9を計算してx = 9を書く。

      順序: 11

      lane: thread-a

    3. Thread A mutex解放

      Thread Aのwrite後にmutexを解放し、happens-beforeを作る。

      順序: 12

      lane: thread-a

    4. Thread B mutex取得

      Thread BはAの解放後にmutexを取得する。

      順序: 13

      lane: thread-b

    5. Thread B read/compute/write

      Thread Bは同期後のx = 9を読み、8を計算してx = 8を書く。

      順序: 14

      lane: thread-b

    6. 同期後の不変条件

      同期後の x = 8は二回減算後の期待値 x = 8と一致し、lost updateはない。

      順序: 15

      lane: shared-context

  4. 停止性と回帰

    安全性だけでなくlock順、timeout、同条件stressを再確認する。

    1. 停止性

      mutexのlock順、待機資源、timeout、cancel経路を観測する。

      順序: 16

      lane: shared-context

    2. 回帰

      同期なしと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、同期前後のコード差、停止性、性能分布を含める。

lost updateとmutex回復の完全な状態遷移

注記

図を読む際の補足情報です。

  1. 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。
状態遷移
イベント開始終了判定理由
nextA readB readallowed
nextB readA computeallowed
nextA computeB computeallowed
nextB computeA writeallowed
nextA writeB write / lost updateallowed
nextB write / lost updateA lockallowed
nextA lockA locked read/compute/writeallowed
nextA locked read/compute/writeA unlockallowed
nextA unlockB lockallowed
nextB lockB locked read/compute/writeallowed
nextB locked read/compute/writeinvariant completeallowed
nextinvariant completeB write / lost updaterejectedread/compute/write全体を同じmutex境界で保護する。

read-old-valueからlost-updateを経て、mutexで二回の減算を直列化したlocked-completeまでを再現する。

  1. A read: read/compute/write trace 1: Aがx=10をread。; 条件 常時; node read-old-value; edge なし
  2. B read: trace 2: Bもx=10をread。; 条件 常時; node b-read-old-value; edge step-01
  3. A compute: trace 3: Aがlocal 9をcompute。; 条件 常時; node a-compute; edge step-02
  4. B compute: trace 4: Bもlocal 9をcompute。; 条件 常時; node b-compute; edge step-03
  5. A write: trace 5: Aがx=9をwrite。; 条件 常時; node a-write; edge step-04
  6. B write: trace 6: Bが9を上書き。期待値x=8、actual x=9でinvariant違反。; 条件 常時; node lost-update; edge step-05
  7. A lock: trace 7: fixtureを10へ戻しAがlock。; 条件 常時; node lock-acquired; edge step-06
  8. A locked update: trace 8: mutex内でAがread/compute/writeし10→9。; 条件 常時; node a-locked-write; edge step-07
  9. A unlock: trace 9: Aがwrite後にunlock。; 条件 常時; node unlock; edge step-08
  10. B lock: trace 10: Bがlockし、同期後の9を観測。; 条件 常時; node b-lock-acquired; edge step-09
  11. B locked update: trace 11: Bがread/compute/writeし9→8。; 条件 常時; node b-locked-write; edge step-10
  12. invariant complete: trace 12: 期待値x=8、actual x=8。lock/unlockを含む完全trace。; 条件 常時; node locked-complete; edge step-11
完全な遷移
イベント開始終了条件
nextread-old-valueb-read-old-value常時
timerread-old-valueb-read-old-value常時
nextb-read-old-valuea-compute常時
timerb-read-old-valuea-compute常時
previousb-read-old-valueread-old-value常時
resetb-read-old-valueread-old-value常時
nexta-computeb-compute常時
timera-computeb-compute常時
previousa-computeb-read-old-value常時
reseta-computeread-old-value常時
nextb-computea-write常時
timerb-computea-write常時
previousb-computea-compute常時
resetb-computeread-old-value常時
nexta-writelost-update常時
timera-writelost-update常時
previousa-writeb-compute常時
reseta-writeread-old-value常時
nextlost-updatelock-acquired常時
timerlost-updatelock-acquired常時
previouslost-updatea-write常時
resetlost-updateread-old-value常時
nextlock-acquireda-locked-write常時
timerlock-acquireda-locked-write常時
previouslock-acquiredlost-update常時
resetlock-acquiredread-old-value常時
nexta-locked-writeunlock常時
timera-locked-writeunlock常時
previousa-locked-writelock-acquired常時
reseta-locked-writeread-old-value常時
nextunlockb-lock-acquired常時
timerunlockb-lock-acquired常時
previousunlocka-locked-write常時
resetunlockread-old-value常時
nextb-lock-acquiredb-locked-write常時
timerb-lock-acquiredb-locked-write常時
previousb-lock-acquiredunlock常時
resetb-lock-acquiredread-old-value常時
nextb-locked-writelocked-complete常時
timerb-locked-writelocked-complete常時
previousb-locked-writeb-lock-acquired常時
resetb-locked-writeread-old-value常時
previouslocked-completeb-locked-write常時
resetlocked-completeread-old-value常時
観測結果
結果状態
同期なしtraceでは期待値8とactual 9の差を示す。lost-update
mutex traceでは最終値8と停止性の両方を検証する。locked-complete

現在の状態: A read — read/compute/write trace 1: Aがx=10をread。

このモデルは例示的かつ決定的であり、実システムの完全な再現ではありません。

トレードオフと失敗モード

隔離と共有を選ぶdecision table
主要制約 候補 利点 主要リスク
大量の小さい共有状態 同一processのthreads 低い通信変換費用 race、deadlock、クラッシュ伝播
信頼できない処理または権限分離 別process address spaceと権限の隔離 IPC契約、copy、部分失敗
所有権を一箇所へ集約可能 message passing 共有mutable stateの縮小 mailbox滞留、順序、backpressure

もっともらしい誤診と反証

  1. 誤診: 「二workerが止まったのでdeadlock」。反証: lock取得と待機の循環を記録する。I/O待ちなら外部完了で進み、starvationなら他workerだけ進捗する。停止だけでは循環待ちを証明しない。
  2. 誤診: 「Java volatileを付けて100回通ったのでrace修正済み」。反証: read-modify-writeをeventへ分解し、同じinterleavingがJava Language Specification上可能か確認する。volatileがvisibilityとorderingのhappens-beforeを作っても複合incrementのatomicityを保証しない。

lock範囲を広げると不変条件は守りやすいが並列性が下がり、複数lockは並列性を上げる一方で順序規約が必要になる。timeoutでlockを諦める場合も、途中状態のrollbackまたは再試行契約が要る。

知識チェック

  1. 整数代入がatomicならcounter増加もatomicか。通常は同じではない。増加はread、compute、writeの複合操作になり得る。
  2. process間ならraceは起きないか。共有memory、file、databaseを介せば起きる。address space隔離は全資源の隔離ではない。
  3. 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を再現して不変条件で直す

提出成果物: 競合を再現し不変条件で修正した実験記録

  1. Java 21で教材のRaceLab.javaをコンパイルし、volatile intへの非atomicなread-modify-writeを200試行する
  2. 最終値は0であり負数にならないという不変条件を先に明文化し、違反回数を記録する
  3. 失敗試行の読取り、計算、書込みinterleavingを最小のイベント列として再構成する
  4. Java synchronizedでread-modify-write全体を修正し、同じ200試行で不変条件と停止性を確認する
  5. 直列、同期ありスレッド、プロセス間メッセージの所要時間と障害半径を比較する

説明して理解を確かめる

5分で、プロセスとスレッドの共有境界、lost updateのinterleaving、Java volatileがvisibilityとorderingを扱っても複合incrementのatomicityを得られない理由を説明する。

アセスメント

  1. 問い: 二つのworkerが停止した。『deadlockだ』という診断を確定するために必要な証拠と、別のもっともらしい原因を示す。

    期待する証拠: 待機資源とロック取得順の循環、進捗観測、I/O待ちまたはstarvationとの反証

  2. 問い: 共有変数へvolatileを付けたら再現しなくなった。修正完了と判断できない理由は何か。

    期待する証拠: Java SE 21のvolatile happens-before、複合操作のatomicity、確率的な一回の非再現を分けた説明

別問題へ転用する

プロセス分離とスレッド共有の再比較

復習スケジュール

  1. 1日後

    atomicity、可視性、順序性のうち、今回の不変条件に必要なのはどれか

  2. 7日後

    deadlock説を反証する最小のイベント証拠は何か

  3. 30日後

    隔離を強めると通信と復旧の設計はどう変わるか

  4. 90日後

    atomicity、可視性、順序性のうち、今回の不変条件に必要なのはどれか

評価ルーブリック

4段階の評価基準
観点未達発展途上熟達卓越
technical-correctnessvolatile、barrier、mutexを同義に扱うか、lost updateの操作列を説明できない競合は説明するが、atomicityと可視性またはdeadlockとstarvationを混同する不変条件、interleaving、happens-before、同期範囲を一貫して説明する停止性とメモリ順序の反例を含め、実装とOS境界の双方を正確に扱う
judgmentプロセスは重い、スレッドは軽いという一軸だけで選択する性能と隔離を比較するが、通信、復旧、信頼境界が抜ける障害半径、共有量、通信費、観測性を制約に合わせて比較する権限分離と段階的な並列度を設計し、負荷変化時の撤回条件を示す
evidence一度再現しなかったことを修正の証拠にする反復試験はあるが、失敗時のイベント列または停止性の観測がない100試行以上の失敗頻度、最小interleaving、修正後の同条件回帰を示すスケジューリングを変えたstress試験と性能分布で正しさと費用を検証する
communicationraceがあったという結論だけで、不変条件と再現方法がないコード変更は示すが、なぜ修正になるかをイベント列で説明しない再現条件、不変条件、原因、同期範囲、回帰結果をレビュー可能に結ぶ運用者へ検知方法を、実装者へ同期規約を同じ証拠から説明できる

出典

以下の外部資料は利用者が選択したときだけ開きます。