Virtual-Threads-Lab

23 Java-Dateien. 9 zentrale Dateien werden direkt mit echtem Quellcode und ihrem Zusammenspiel erklärt.

Zurück zu Code-Labs

Kapitel 16 · Virtual Threads und Performance

Was dieses Lab zeigt

Zeigt, warum Virtual Threads knappe Downstreams nicht unbegrenzt skalieren. Bulkheads, Connection-Pool-Grenzen, Rate Limits, Timeouts und Metriken bilden ein realistisches Schutzsystem.

Lernziele

  • Virtual Threads richtig einordnen
  • Downstreams mit Bulkheads schützen
  • Timeouts und Limits messen

Technik und Schwerpunkte

Java 2123 Java-Dateien1 Tests/RunnerVirtual ThreadsBulkheadRate Limiting
Echter Quellcode aus diesem Lab

Geführter Codepfad

Der Performance-Pfad kombiniert parallele I/O-Aufrufe mit Bulkhead, Timeout, Rate Limit und Connection Pool. Ein Advisor trennt I/O-gebundene von CPU-gebundenen Workloads und begründet die Ausführungsstrategie.

Run8BDemoOrderViewServiceBulkheadTimeoutPolicyWindowRateLimiterSimulatedConnectionPoolWorkloadClassifierPerformanceAdvisorRun8BTestRunner
Lesereihenfolge der zentralen Klassen. Die Pfeile zeigen den didaktischen Weg durch den realen Quellcode, nicht zwingend jeden Laufzeitaufruf.
1. Run8BDemoStartet verschiedene Lastprofile und zeigt erfolgreiche, gedrosselte und abgebrochene Aufrufe.
2. OrderViewServiceLädt Order-, Payment- und Inventory-Daten parallel und setzt daraus die Sicht zusammen.
3. BulkheadBegrenzt gleichzeitig laufende Aufrufe, damit eine Abhängigkeit nicht alle Ausführungskapazität bindet.
4. TimeoutPolicyBeendet Arbeit, die das definierte Zeitbudget überschreitet.
5. WindowRateLimiterBegrenzt die Anzahl akzeptierter Aufrufe pro Zeitfenster.
6. SimulatedConnectionPoolMacht sichtbar, dass virtuelle Threads externe Poolgrenzen nicht aufheben.
7. WorkloadClassifierUnterscheidet I/O- und CPU-Last als Grundlage der Laufzeitentscheidung.
8. PerformanceAdvisorLeitet aus Workload und Ressourcenengpässen eine konkrete PerformanceDecision ab.
9. Run8BTestRunnerPrüft Parallelität, Schutzmechanismen, Poolgrenzen und Beratungsergebnisse.

1. Run8BDemo

src/main/java/com/example/run8b/Run8BDemo.java
Java-Datei öffnen
Rolle im Ablauf

Startet verschiedene Lastprofile und zeigt erfolgreiche, gedrosselte und abgebrochene Aufrufe.

Im Lesepfad folgt OrderViewService: Lädt Order-, Payment- und Inventory-Daten parallel und setzt daraus die Sicht zusammen.

Typ
class Run8BDemo
Verwendet
OrderViewService, PerformanceAdvisor
Verwendet von
Einstiege
main(String[] args)
package com.example.run8b;
public final class Run8BDemo {
    public static void main(String[] args) {
        MetricsRegistry metrics = new MetricsRegistry();
        OrderView view = OrderViewService.demoService(metrics).build(new OrderId("ORD-8B-1"));
        PerformanceDecision decision = new PerformanceAdvisor().decide(true, true, false);
        System.out.println("RUN8B_DEMO_OK " + view.order().id().value() + " " + decision.classifier());
        System.out.println(metrics.snapshot());
    }
}

2. OrderViewService

src/main/java/com/example/run8b/OrderViewService.java
Java-Datei öffnen
Rolle im Ablauf

Lädt Order-, Payment- und Inventory-Daten parallel und setzt daraus die Sicht zusammen.

Im Lesepfad folgt Bulkhead: Begrenzt gleichzeitig laufende Aufrufe, damit eine Abhängigkeit nicht alle Ausführungskapazität bindet.

Typ
class OrderViewService
Verwendet
Bulkhead, TimeoutPolicy, WindowRateLimiter, SimulatedConnectionPool
Verwendet von
Run8BDemo, Run8BTestRunner
Einstiege
build(OrderId id), demoService(MetricsRegistry metrics)
Application Service
package com.example.run8b;
import java.time.Duration;
import java.util.concurrent.*;
// Pattern: Application Service - koordiniert Use Case, Ports und technische Schutzgrenzen.
public final class OrderViewService {
    private final FakeOrderRepository orders;
    private final FakePaymentClient payments;
    private final FakeInventoryClient inventory;
    private final Bulkhead paymentBulkhead;
    private final Bulkhead inventoryBulkhead;
    private final TimeoutPolicy timeout;
    public OrderViewService(FakeOrderRepository orders, FakePaymentClient payments, FakeInventoryClient inventory,
                            Bulkhead paymentBulkhead, Bulkhead inventoryBulkhead, TimeoutPolicy timeout) {
        this.orders = orders; this.payments = payments; this.inventory = inventory; this.paymentBulkhead = paymentBulkhead; this.inventoryBulkhead = inventoryBulkhead; this.timeout = timeout;
    }
    public OrderView build(OrderId id) {
        try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
            Future<OrderData> order = executor.submit(() -> orders.load(id));
            Future<PaymentData> payment = executor.submit(() -> paymentBulkhead.call(() -> payments.load(id)));
            Future<InventoryData> stock = executor.submit(() -> inventoryBulkhead.call(() -> inventory.load(id)));
            return new OrderView(timeout.get(order, "order"), timeout.get(payment, "payment"), timeout.get(stock, "inventory"), false);
        }
    }
    public static OrderViewService demoService(MetricsRegistry metrics) {
        SimulatedConnectionPool pool = new SimulatedConnectionPool(2, metrics);
        return new OrderViewService(new FakeOrderRepository(pool), new FakePaymentClient(new WindowRateLimiter(10, Duration.ofSeconds(1).toMillis())), new FakeInventoryClient(), new Bulkhead("payment", 4, metrics), new Bulkhead("inventory", 4, metrics), new TimeoutPolicy(Duration.ofMillis(500)));
    }
}

3. Bulkhead

src/main/java/com/example/run8b/Bulkhead.java
Java-Datei öffnen
Rolle im Ablauf

Begrenzt gleichzeitig laufende Aufrufe, damit eine Abhängigkeit nicht alle Ausführungskapazität bindet.

Im Lesepfad folgt TimeoutPolicy: Beendet Arbeit, die das definierte Zeitbudget überschreitet.

Typ
class Bulkhead
Verwendet
Verwendet von
OrderViewService, Run8BTestRunner
Einstiege
call(Supplier<T> work)
Bulkhead
package com.example.run8b;
import java.util.concurrent.Semaphore;
import java.util.function.Supplier;
// Pattern: Bulkhead - schützt knappe Downstream-Ressourcen vor Überlast.
public final class Bulkhead {
    private final String name;
    private final Semaphore permits;
    private final MetricsRegistry metrics;
    public Bulkhead(String name, int maxConcurrentCalls, MetricsRegistry metrics) {
        this.name = name;
        this.permits = new Semaphore(maxConcurrentCalls);
        this.metrics = metrics;
    }
    public <T> T call(Supplier<T> work) {
        if (!permits.tryAcquire()) {
            metrics.increment(name + ".rejected");
            throw new BulkheadRejectedException("bulkhead rejected: " + name);
        }
        metrics.increment(name + ".accepted");
        try { return work.get(); }
        finally { permits.release(); }
    }
}

4. TimeoutPolicy

src/main/java/com/example/run8b/TimeoutPolicy.java
Java-Datei öffnen
Rolle im Ablauf

Beendet Arbeit, die das definierte Zeitbudget überschreitet.

Im Lesepfad folgt WindowRateLimiter: Begrenzt die Anzahl akzeptierter Aufrufe pro Zeitfenster.

Typ
class TimeoutPolicy
Verwendet
Verwendet von
OrderViewService
Einstiege
get(Future<T> future, String label)
Policy Object
package com.example.run8b;
import java.time.Duration;
import java.util.concurrent.*;
// Pattern: Policy Object - kapselt Timeout-Entscheidung je Use Case.
public final class TimeoutPolicy {
    private final Duration timeout;
    public TimeoutPolicy(Duration timeout) { this.timeout = timeout; }
    public <T> T get(Future<T> future, String label) {
        try { return future.get(timeout.toMillis(), TimeUnit.MILLISECONDS); }
        catch (TimeoutException e) { future.cancel(true); throw new TimeoutExceededException("timeout in " + label); }
        catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new TimeoutExceededException("interrupted in " + label); }
        catch (ExecutionException e) { throw new RuntimeException(e.getCause()); }
    }
}

5. WindowRateLimiter

src/main/java/com/example/run8b/WindowRateLimiter.java
Java-Datei öffnen
Rolle im Ablauf

Begrenzt die Anzahl akzeptierter Aufrufe pro Zeitfenster.

Im Lesepfad folgt SimulatedConnectionPool: Macht sichtbar, dass virtuelle Threads externe Poolgrenzen nicht aufheben.

Typ
class WindowRateLimiter
Verwendet
Verwendet von
OrderViewService, Run8BTestRunner
Einstiege
tryAcquire()
Rate Limiter
package com.example.run8b;
// Pattern: Rate Limiter - begrenzt Aufrufe pro Zeitfenster.
public final class WindowRateLimiter {
    private final int maxPerWindow;
    private final long windowMillis;
    private long windowStart = System.currentTimeMillis();
    private int used;
    public WindowRateLimiter(int maxPerWindow, long windowMillis) { this.maxPerWindow = maxPerWindow; this.windowMillis = windowMillis; }
    public synchronized boolean tryAcquire() {
        long now = System.currentTimeMillis();
        if (now - windowStart >= windowMillis) { windowStart = now; used = 0; }
        if (used >= maxPerWindow) return false;
        used++;
        return true;
    }
}

6. SimulatedConnectionPool

src/main/java/com/example/run8b/SimulatedConnectionPool.java
Java-Datei öffnen
Rolle im Ablauf

Macht sichtbar, dass virtuelle Threads externe Poolgrenzen nicht aufheben.

Im Lesepfad folgt WorkloadClassifier: Unterscheidet I/O- und CPU-Last als Grundlage der Laufzeitentscheidung.

Typ
class SimulatedConnectionPool
Verwendet
Verwendet von
OrderViewService
Einstiege
withConnection(Supplier<T> query)
Resource Pool
package com.example.run8b;
import java.util.concurrent.Semaphore;
import java.util.function.Supplier;
// Pattern: Resource Pool - modelliert begrenzte Datenbankverbindungen.
public final class SimulatedConnectionPool {
    private final Semaphore connections;
    private final MetricsRegistry metrics;
    public SimulatedConnectionPool(int size, MetricsRegistry metrics) { this.connections = new Semaphore(size); this.metrics = metrics; }
    public <T> T withConnection(Supplier<T> query) {
        if (!connections.tryAcquire()) {
            metrics.increment("db.pool.exhausted");
            throw new PoolExhaustedException("no database connection available");
        }
        metrics.increment("db.pool.used");
        try { return query.get(); }
        finally { connections.release(); }
    }
}

7. WorkloadClassifier

src/main/java/com/example/run8b/WorkloadClassifier.java
Java-Datei öffnen
Rolle im Ablauf

Unterscheidet I/O- und CPU-Last als Grundlage der Laufzeitentscheidung.

Im Lesepfad folgt PerformanceAdvisor: Leitet aus Workload und Ressourcenengpässen eine konkrete PerformanceDecision ab.

Typ
enum WorkloadClassifier
Verwendet
Verwendet von
PerformanceAdvisor, Run8BTestRunner
Einstiege
package com.example.run8b;
public enum WorkloadClassifier { IO_BOUND, CPU_BOUND, RESOURCE_BOUND }

8. PerformanceAdvisor

src/main/java/com/example/run8b/PerformanceAdvisor.java
Java-Datei öffnen
Rolle im Ablauf

Leitet aus Workload und Ressourcenengpässen eine konkrete PerformanceDecision ab.

Im Lesepfad folgt Run8BTestRunner: Prüft Parallelität, Schutzmechanismen, Poolgrenzen und Beratungsergebnisse.

Typ
class PerformanceAdvisor
Verwendet
WorkloadClassifier
Verwendet von
Run8BDemo, Run8BTestRunner
Einstiege
decide(boolean waitsForNetwork, boolean usesDbPool, boolean heavyCpu)
Strategy/Decision Service
package com.example.run8b;
// Pattern: Strategy/Decision Service - ordnet Workload-Profile einer Architekturentscheidung zu.
public final class PerformanceAdvisor {
    public PerformanceDecision decide(boolean waitsForNetwork, boolean usesDbPool, boolean heavyCpu) {
        if (heavyCpu) return new PerformanceDecision(WorkloadClassifier.CPU_BOUND, "limit to CPU cores and profile allocation");
        if (usesDbPool) return new PerformanceDecision(WorkloadClassifier.RESOURCE_BOUND, "use virtual threads but protect DB pool with bulkhead and timeout");
        if (waitsForNetwork) return new PerformanceDecision(WorkloadClassifier.IO_BOUND, "virtual threads are a good fit with downstream limits");
        return new PerformanceDecision(WorkloadClassifier.CPU_BOUND, "measure first; do not add concurrency blindly");
    }
}

9. Run8BTestRunner

src/test/java/com/example/run8b/Run8BTestRunner.java
Java-Datei öffnen
Rolle im Ablauf

Prüft Parallelität, Schutzmechanismen, Poolgrenzen und Beratungsergebnisse.

Damit ist der zentrale Pfad abgeschlossen; der Test-/Runner-Code und die vollständige Dateiliste darunter zeigen die übrigen Varianten.

Typ
class Run8BTestRunner
Verwendet
OrderViewService, Bulkhead, WindowRateLimiter, WorkloadClassifier, PerformanceAdvisor
Verwendet von
Einstiege
main(String[] args)
package com.example.run8b;
import java.time.Duration;
public final class Run8BTestRunner {
    public static void main(String[] args) {
        buildsOrderViewWithVirtualThreads();
        rejectsWhenBulkheadIsFull();
        detectsRateLimit();
        classifiesResourceBoundWorkload();
        cpuChecksumIsStable();
        System.out.println("RUN8B_TESTS_OK");
    }
    static void buildsOrderViewWithVirtualThreads() {
        MetricsRegistry metrics = new MetricsRegistry();
        OrderView view = OrderViewService.demoService(metrics).build(new OrderId("ORD-1"));
        require(view.payment().authorized(), "payment should be authorized");
        require(metrics.get("payment.accepted") == 1, "payment metric missing");
    }
    static void rejectsWhenBulkheadIsFull() {
        MetricsRegistry metrics = new MetricsRegistry();
        Bulkhead bulkhead = new Bulkhead("tiny", 1, metrics);
        bulkhead.call(() -> "ok");
        require(metrics.get("tiny.accepted") == 1, "accepted metric missing");
    }
    static void detectsRateLimit() {
        WindowRateLimiter limiter = new WindowRateLimiter(1, Duration.ofSeconds(60).toMillis());
        require(limiter.tryAcquire(), "first call allowed");
        require(!limiter.tryAcquire(), "second call rejected");
    }
    static void classifiesResourceBoundWorkload() {
        PerformanceDecision d = new PerformanceAdvisor().decide(true, true, false);
        require(d.classifier() == WorkloadClassifier.RESOURCE_BOUND, "db pool workload must be resource bound");
    }
    static void cpuChecksumIsStable() {
        long a = new CpuWorkload().checksum(1000);
        long b = new CpuWorkload().checksum(1000);
        require(a == b, "checksum must be deterministic");
    }
    static void require(boolean condition, String message) { if (!condition) throw new AssertionError(message); }
}
Alle Projektdateien öffnen (27 Einträge)
⌂ Cockpit