İçeriğe geç

CompletableFuture ve Structured Concurrency

Uzman 10 dk Sık karşılaşılır

Önce şunu oku: Virtual Threads vs Platform Threads

30 saniyede özet

Üç işi aynı anda başlattın ve biri patladı. allOf diğer ikisini durdurmaz, onlar boşuna çalışmaya devam eder. Structured concurrency tam olarak bu dağınıklığı toplamak için var.

Üç dükkândan aynı anda sipariş verdin ve biri hemen “yok” dedi. Öteki iki kurye ise hâlâ yolda.

Yazılımda da bir istek üç servise paralel gidiyor ve biri patlıyor. Asıl soru hatanın kendisi değil, geride kalanlar.

allOf ile üç çağrı paralel başladı ve biri 30 ms'de patladı. Diğer ikisi ne olur? Cevabı göster

Çoğu kişi “iptal olurlar” der. Olmazlar — çalışmaya, bağlantı tutmaya devam ederler. Hangi executorGörevleri hangi thread'lerin çalıştıracağına karar veren bileşen. Belirtmezsen çoğu asenkron API ortak ForkJoinPool'u kullanır — ve o havuz tüm uygulamayla paylaşılır.Sözlükte gör → üzerinde çalıştıkları da bunu değiştirmez.

Aynı hata, iki son: biri iş sızdırır, öteki toplar.
Adım adım oku
  1. Üç çağrı aynı anda başlar.
  2. Biri erken patlar ve isteğe hata cevabı döner.
  3. allOf ile öteki iki çağrı cevap döndükten sonra da çalışmaya ve bağlantı tutmaya devam eder.
  4. StructuredTaskScope ile hata anında öteki iki çağrı da iptal edilir. Sızan iş kalmaz.
  1. Bayt: Üç servisi allOf ile aynı anda çağırdım. Biri hemen patladı, ben de hatayı döndüm.

  2. Sen: Güzel, iş bitti o zaman.

  3. Bayt: Bitmedi! Öteki iki çağrı hâlâ çalışıyor, bağlantı tutuyor, kimse onları beklemiyor.

  4. Bayt: Structured concurrency hepsini tek çantaya koyar: çantayı bırakınca içindekiler de iptal olur.

Kendin gör

allOf kardeşleri iptal etmez — structured concurrency neden var

Tohum 8

CompletableFuture.allOf

  • fetchUserçalışıyor
  • fetchOrdersçalışıyor
  • fetchRecommendationsçalışıyor
cevap30 ms300 ms

Cevap süresi

30 ms

sırayla olsa 110 ms

Sızan iş

320 ms

cevaptan sonra çalışan

İşler ne zaman durdu

300 ms

iptal yok

Hız
Adım 0

Şu an ne oldu?

3 çağrı aynı anda uçuşta — 0 ms

Çağrılar birbirini beklemiyor; toplam süre en yavaş çağrı kadar olacak.

Aklında kalsın: thenCompose zincirlemek paralellik üretmez; paralellik için çağrıların birbirinden bağımsız başlatılması gerekir.

Olay günlüğü (0)

Henüz olay yok. Oynat veya adımla.

Üç kurulumu aynı hata senaryosunda karşılaştır:

  1. allOf + Erken hata. Cevap 30 ms’de döndü ama 320 ms iş sızdı — kırmızı bloklar cevaptan sonra çalışan kısım. İşler 300 ms’ye kadar sürdü.
  2. Sıraylaya geç. Sızan iş sıfır — çünkü sonraki çağrılar hiç başlamadı. Ama cevap 110 ms’de geldi.
  3. StructuredTaskScope. Cevap yine 30 ms, sızan iş sıfır, iki çağrı iptal edildi.
  4. Hatayı Patlamasın yap. Üç kurulum arasında allOf ile structured ayırt edilemiyor.

Dördüncü adım önemli: bu iki API’nin farkı mutlu yolda görünmez. Fark yalnızca hata anında ortaya çıkar.

Kafam karıştı, daha basit anlat

Üç arkadaşı aynı anda markete yolladın. Biri “market kapalı” diye döndü. Diğer ikisi hâlâ yolda ve kimse onlara “dönün” demedi.

Hızlı kontrolOrta

Her durumda hangi birleştirme metodunu kullanırsın?

Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

Sınıflandırılmamış

thenCompose / zincirleme

Sonraki adım öncekinin sonucuna ihtiyaç duyuyor

    thenCombine / allOf

    Çağrılar birbirinden bağımsız

      Üç yol, üç bedel

      GecikmeSızan işHatayı öğrenme
      SıraylaToplam (530 ms)✅ yokGeç
      allOfEn yavaş (300 ms)❌ 320 msHızlı
      StructuredTaskScopeEn yavaş (300 ms)✅ yokHızlı

      Sıralı çağrının tek avantajı, başlatılmamış işin sızamamasıdır. Structured concurrency bunu paralellikten vazgeçmeden verir — var oluş sebebi budur.

      Kafam karıştı, daha basit anlat

      Structured concurrency, ekibi tek bir lider altında yollamaktır. Biri başarısız olursa lider herkesi geri çağırır, kimse yolda unutulmaz.

      Hızlı kontrolİleri

      Üç bağımsız servis çağrısını paralel yapıp hepsini beklemek istiyorsun. En temiz yol?

      Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

      allOf neden iptal etmiyor

      var user = supplyAsync(() -> fetchUser(id));
      var orders = supplyAsync(() -> fetchOrders(id));
      var recs = supplyAsync(() -> fetchRecommendations(id));
      allOf(user, orders, recs).join(); // orders patlarsa burası hemen atar

      allOf dönen future, kaynaklardan biri patladığında istisnai olarak tamamlanır. Ama user ve recs bağımsız nesnelerdir; allOf onların sahibi değildir, yalnızca izleyicisidir.

      Sonuç: join() istisna atar, çağıran metottan çıkar, arka planda iki çağrı hâlâ çalışır — sonuçlarını kimse okumaz.

      Üstelik cancel(true) çağırmak da çözmez: CompletableFuture.cancel çalışan görevi kesmez, yalnızca future’ı iptal edilmiş olarak işaretler.

      Aşağıdaki örnek bir bankanın döviz masasından: aynı kur üç likidite sağlayıcısına aynı anda soruluyor ve en iyi teklif seçiliyor. İki yaklaşım yan yana.

      Derinleş · Üç likidite sağlayıcısından en iyi kur: iki yaklaşım 4 dosya · ~88 satır · ilk okumada atlayabilirsin
      Proje dosyaları

      src/main/java/fx/ LiquidityProvider.java Likidite sağlayıcısı sözleşmesi ve kur teklifi tipi.

      src/main/java/fx/LiquidityProvider.java
      // The bank's FX desk buys currency from several liquidity providers (other banks)
      // and gives the customer the best ask it can get.
      public interface LiquidityProvider {
      String name();
      /** Blocking call to the provider's pricing API. */
      Quote quote(String pair, BigDecimal amount);
      }
      record Quote(String provider, String pair, BigDecimal ask) {
      }

      src/main/java/fx/ BestRateWithFutures.java CompletableFuture: kendi executor'ı, sağlayıcı başına zaman aşımı ve hatada boş sonuç. Biri yavaşlarsa diğerleri beklemez.

      src/main/java/fx/BestRateWithFutures.java
      public class BestRateWithFutures {
      private final List<LiquidityProvider> providers;
      // Never the common ForkJoinPool for blocking I/O: its few threads would all wait.
      private final ExecutorService io = Executors.newVirtualThreadPerTaskExecutor();
      public BestRateWithFutures(List<LiquidityProvider> providers) {
      this.providers = providers;
      }
      public Optional<Quote> best(String pair, BigDecimal amount) {
      List<CompletableFuture<Optional<Quote>>> calls = providers.stream()
      .map(provider -> CompletableFuture
      .supplyAsync(() -> Optional.of(provider.quote(pair, amount)), io)
      .orTimeout(800, TimeUnit.MILLISECONDS) // a quote older than this is useless
      .exceptionally(ex -> Optional.empty())) // a failed provider is just "no quote"
      .toList();
      CompletableFuture.allOf(calls.toArray(CompletableFuture[]::new)).join();
      return calls.stream()
      .map(CompletableFuture::join) // already complete: no waiting
      .flatMap(Optional::stream)
      .min(Comparator.comparing(Quote::ask));
      }
      }

      src/main/java/fx/ BestRateStructured.java StructuredTaskScope: biri hata verirse diğerleri iptal edilir, blok bitince hiçbir alt görev arkada kalmaz.

      src/main/java/fx/BestRateStructured.java
      // Preview API in Java 21–24 (compile and run with --enable-preview).
      public class BestRateStructured {
      private final List<LiquidityProvider> providers;
      public BestRateStructured(List<LiquidityProvider> providers) {
      this.providers = providers;
      }
      // Here every provider is required (say, a large trade that must be split across
      // all of them): if one fails, the others are cancelled and the whole call fails.
      // The subtasks can never outlive this method.
      public Quote best(String pair, BigDecimal amount) throws InterruptedException, ExecutionException {
      try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
      List<StructuredTaskScope.Subtask<Quote>> tasks = providers.stream()
      .map(provider -> scope.fork(() -> provider.quote(pair, amount)))
      .toList();
      scope.joinUntil(Instant.now().plusMillis(800)); // one deadline for all of them
      scope.throwIfFailed();
      return tasks.stream()
      .map(StructuredTaskScope.Subtask::get)
      .min(Comparator.comparing(Quote::ask))
      .orElseThrow();
      } catch (TimeoutException ex) {
      throw new ExecutionException("providers too slow", ex);
      }
      }
      }

      src/main/java/fx/ Main.java İkisini çalıştıran program.

      src/main/java/fx/Main.java
      public class Main {
      public static void main(String[] args) throws Exception {
      List<LiquidityProvider> providers = List.of(
      new FakeProvider("banka-a", "38.4150", Duration.ofMillis(50)),
      new FakeProvider("banka-b", "38.3900", Duration.ofSeconds(3)), // misses the 800 ms budget
      new FakeProvider("banka-c", "38.4020", Duration.ofMillis(200)));
      BigDecimal amount = new BigDecimal("250000"); // EUR
      System.out.println(new BestRateWithFutures(providers).best("EURTRY", amount));
      // Optional[Quote[provider=banka-c, pair=EURTRY, ask=38.4020]] — the slow, cheaper one timed out
      try {
      new BestRateStructured(providers).best("EURTRY", amount);
      } catch (ExecutionException ex) {
      System.out.println("structured: " + ex.getMessage());
      // structured: providers too slow — and the slow call was cancelled, not left running
      }
      }
      }

      StructuredTaskScope.

      try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
      var user = scope.fork(() -> fetchUser(id));
      var orders = scope.fork(() -> fetchOrders(id));
      var recs = scope.fork(() -> fetchRecommendations(id));
      scope.join(); // hepsi bitene ya da biri patlayana kadar
      scope.throwIfFailed();
      return new Profile(user.get(), orders.get(), recs.get());
      } // try-with-resources: scope kapanmadan buradan çıkılamaz

      İki garanti veriyor:

      • Biri patlarsa diğerleri iptal edilir (ShutdownOnFailure).
      • Blok bitmeden hiçbir alt görev hayatta kalamaz. try-with-resources bunu dilin kendisiyle zorunlu kılar.

      ShutdownOnSuccess ise tersini yapar: ilk başarılı sonuç gelince kalanları iptal eder — birden çok sağlayıcıya aynı soruyu sorup en hızlısını almak için.

      Bu API Java 21’de önizleme olarak geldi ve sonraki sürümlerde şekil değiştirdi; kalıcı olan fikir, imza değil.

      Hızlı kontrolOrta

      `thenApply` ile `thenCompose` arasındaki fark nedir?

      Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

      Zincirleme ve hata

      Klasik soru, tek cümlelik cevap: thenApply map’tir, thenCompose flatMap.

      // fonksiyon düz bir değer döndürüyor → thenApply
      future.thenApply(user -> user.name()); // CompletableFuture<String>
      // fonksiyon zaten bir future döndürüyor → thenCompose
      future.thenCompose(user -> fetchOrders(user.id())); // CompletableFuture<List<Order>>
      future.thenApply(user -> fetchOrders(user.id())); // ❌ CompletableFuture<CompletableFuture<...>>

      thenCombine ise iki bağımsız future’ı birleştirir — zincirleme değil, paralel:

      user.thenCombine(orders, Profile::new); // ikisi de paralel çalışır

      Dikkat: thenCompose zincirlemek paralellik üretmez. Her adım bir öncekini bekler; simülatördeki “Sırayla” seçeneği tam olarak budur.

      Hata yönetimi. Bir aşama patlarsa, aşağıdaki tüm thenApply adımları atlanır ve istisna ilk exceptionally/handle’a kadar akar.

      MetotNe zaman çalışırNe döner
      exceptionallyYalnızca hata olduğundaYedek değer
      handleHer durumda(sonuç, hata) → yeni değer
      whenCompleteHer durumdaDeğeri değiştirmez, yalnızca gözlemler

      İki tuzak:

      • İstisna sarmalanır. join() CompletionException, get() ExecutionException atar. Asıl hata getCause() içindedir.
      • Kimse beklemezse istisna sessizce kaybolur. thenAccept ile biten ve hiç join edilmeyen bir zincirdeki hata hiçbir yere yazılmaz. Zinciri her zaman exceptionally veya whenComplete ile kapat.
      Kafam karıştı, daha basit anlat

      Sonraki adım düz bir değer döndürüyorsa thenApply yaz. Kendisi de bir future döndürüyorsa thenCompose yaz, yoksa kutunun içinde kutu olur.

      Hızlı kontrolOrta

      Bir `CompletableFuture` zincirinde atılan istisna ne olur?

      Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

      Tuzaklar

      Executor vermezsen supplyAsync ForkJoinPool.commonPool kullanır ve o havuzun boyutu çekirdek sayısı − 1’dir.

      Dört çekirdekli bir makinede üç bloklayan I/O çağrısı havuzu tamamen doldurur ve parallelStream dahil aynı havuzu kullanan her şey durur.

      supplyAsync(() -> blockingCall(), ioExecutor); // I/O için her zaman kendi executor'ın

      Java 21 ile daha iyi bir cevap var: Executors.newVirtualThreadPerTaskExecutor(). Virtual thread bloklandığında carrier threadBir virtual thread'i o an çalıştıran gerçek OS thread'i. Virtual thread blokladığında taşıyıcı serbest kalır ve başka bir virtual thread'e geçer.Sözlükte gör → thread’i serbest bırakır, dolayısıyla havuz doygunluğu sorunu ortadan kalkar.

      Kendini sına

      Önce hızlı bir ısınma: puan yok, kayıt yok.

      Şimşek turu1/5

      allOf, işlerden biri patlayınca ötekileri iptal eder.

      Soru 1/3Uzman

      Üç `CompletableFuture` `allOf` ile bekleniyor ve ikincisi 30 ms'de patlıyor. Diğer ikisine ne olur?

      Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

      Aklında kalacak üç şey

      1. 1 allOf sadece bekler: bir iş patlayınca hata verir ama diğer işleri iptal etmez, onlar çalışmaya devam eder.
      2. 2 Structured concurrency'nin getirdiği şey hız değil, düzen: hiçbir iş onu başlatan istekten uzun yaşamaz.
      3. 3 thenApply ile thenCompose farkı map ile flatMap farkıdır. Yanlışını seçersen elinde iç içe iki future kalır.
      Sonraki kapı Çöp toplayıcı çalışırken uygulaman neden bir anlığına donar? Garbage Collection — Üçlü Takas · 9 dk

      5 kart sonraki derste seni bekliyor

      0/5 kart bu dersten toplandı