En tancar la lliçó anterior, BiblioTech ja era correcte: el catàleg suporta dos empleats alhora, el registre de préstecs manté el seu invariant i el comptador no perd increments. Però el codi per arribar-hi era artesanal: panys que s'han de deixar anar en un finally, ordres d'adquisició que s'han de respectar per convenció, i un fil creat a mà per cada tasca. Els dos-cents avisos de venciment continuen necessitant dos-cents fils, dos-cents megabytes de piles i una espera de join() per cadascun.

Aquesta lliçó canvia el nivell d'abstracció. java.util.concurrent —dissenyat per Doug Lea i incorporat a Java 5— ofereix peces ja construïdes i provades per a tot el que a les lliçons anteriors feies a mà: pools de fils que es reutilitzen, tasques que retornen resultats i propaguen excepcions, programadors periòdics, i sincronitzadors que substitueixen wait/notify per alguna cosa que es pot raonar.

El canvi mental és aquest: deixes de pensar en fils i comences a pensar en tasques. Tu descrius què cal fer; l'executor decideix on i quan es fa. És la mateixa diferència que hi ha entre gestionar la memòria a mà i tenir un recol·lector de brossa.

En acabar, els 200 avisos de BiblioTech s'enviaran amb un pool de vuit fils en tres segons, amb barra de progrés, amb límit de concurrència sobre el recurs compartit i amb cancel·lació neta.

Contingut

  1. Per què java.util.concurrent substitueix la gestió manual
  2. Executor, ExecutorService i Executors
  3. Les fàbriques d'Executors i quan fer servir cadascuna
  4. L'advertiment sobre les cues il·limitades
  5. ThreadPoolExecutor a mida
  6. Triar la mida del pool
  7. El cicle de vida de l'executor
  8. Callable i Future
  9. Cancel·lació de tasques
  10. invokeAll i invokeAny
  11. ScheduledExecutorService
  12. CountDownLatch
  13. CyclicBarrier
  14. Semaphore
  15. Exchanger i taula comparativa de sincronitzadors
  16. ForkJoinPool i divideix i venceràs
  17. BiblioTech: 200 avisos amb pool, progrés i cancel·lació
  18. Errors Comuns i Consells
  19. Exercicis

  1. Per què java.util.concurrent substitueix la gestió manual

Recorda el codi de 08-02 per enviar avisos:

// EL QUE FEIES: un fil per tasca, gestio manual completa.
Thread[] fils = new Thread[200];
for (int i = 0; i < 200; i++) {
    fils[i] = new Thread(new TascaAvis(avisos.get(i)), "bibliotech-avisos-" + i);
    fils[i].start();
}
for (Thread f : fils) {
    f.join();
}
// I el resultat de cada tasca? Camps volatile i llegir-los a ma.
// I si una falla? Un gestor d'excepcions no capturades.
// I si vull cancellar-ho tot? interrupt() als 200, un per un.
// I si son 200.000 avisos? OutOfMemoryError.

Sis problemes, tots reals:

Problema Conseqüència
Un fil per tasca ~1 MB de pila i desenes de µs per fil; no escala
Sense límit de fils Amb moltes tasques, OutOfMemoryError: unable to create native thread
Els fils no es reutilitzen Es paga la creació una vegada i una altra
Runnable no retorna res Camps volatile i convencions implícites
Excepcions que no arriben a qui crida Fallades silencioses (08-02)
Cancel·lació manual Recórrer un array cridant interrupt()

java.util.concurrent resol els sis. El mateix codi passa a ser:

// EL QUE FARAS: descriure tasques, l'executor les colloca.
ExecutorService executor = Executors.newFixedThreadPool(8);
List<Future<ResultatAvis>> futurs = new ArrayList<>();

for (Avis a : avisos) {
    futurs.add(executor.submit(new TascaAvis(a)));   // Callable: retorna valor
}

for (Future<ResultatAvis> f : futurs) {
    ResultatAvis r = f.get();   // espera, i RELLANCA l'excepcio si n'hi va haver
}

executor.shutdown();

Vuit fils en lloc de dos-cents, resultats tipats, excepcions que arriben, i cancel·lació amb una crida.

  1. Executor, ExecutorService i Executors

Tres noms semblants que convé no confondre:

Nom Què és Paper
Executor Interfície amb un mètode: void execute(Runnable) L'abstracció mínima: "executa això"
ExecutorService Interfície que estén Executor Afegeix resultats (submit), apagada i espera
Executors Classe d'utilitat amb mètodes estàtics Fàbrica d'implementacions ja configurades

Executor és deliberadament minimalista, i aquesta és la seva virtut: separa què s'executa de com s'executa.

// La interficie completa. Un sol metode.
public interface Executor {
    void execute(Runnable ordre);
}

// Tres implementacions valides de la MATEIXA interficie:
Executor alMateixFil    = tasca -> tasca.run();
Executor unFilNou       = tasca -> new Thread(tasca).start();
Executor unPool         = Executors.newFixedThreadPool(8);

// El codi client no canvia:
alMateixFil.execute(() -> System.out.println("hola"));
unPool.execute(() -> System.out.println("hola"));

ExecutorService afegeix el que cal a la pràctica:

public interface ExecutorService extends Executor {

    // Enviar tasques i obtenir un Future
    Future<?> submit(Runnable tasca);
    <T> Future<T> submit(Callable<T> tasca);
    <T> Future<T> submit(Runnable tasca, T resultat);

    // Enviar-ne diverses
    <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasques);
    <T> T invokeAny(Collection<? extends Callable<T>> tasques);

    // Apagada
    void shutdown();
    List<Runnable> shutdownNow();
    boolean isShutdown();
    boolean isTerminated();
    boolean awaitTermination(long temps, TimeUnit unitat);
}

Diferència clau entre execute i submit, i és un dels paranys del tema:

// execute(Runnable): no retorna res. Si la tasca llanca una excepcio,
// va al gestor d'excepcions no capturades del fil (08-02) i
// normalment s'imprimeix a la consola.
executor.execute(() -> { throw new RuntimeException("fallada visible"); });

// submit(...): retorna un Future. Si la tasca llanca una excepcio,
// aquesta es DESA al Future i NO s'imprimeix enlloc.
// Si ningu no crida get(), la fallada desapareix SENSE DEIXAR RASTRE.
executor.submit(() -> { throw new RuntimeException("fallada INVISIBLE"); });

Aquest segon cas és una de les causes més freqüents de "la meva tasca no fa res i no surt cap error". La regla: si fas servir submit, crida get() o comprova el resultat d'alguna forma.

  1. Les fàbriques d'Executors i quan fer servir cadascuna

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;

// 1. Pool FIX: N fils, sempre els mateixos, cua illimitada.
ExecutorService fix = Executors.newFixedThreadPool(8);

// 2. Pool ELASTIC: crea fils segons necessitat, els recicla als 60 s.
ExecutorService elastic = Executors.newCachedThreadPool();

// 3. UN SOL FIL: garanteix ordre d'execucio.
ExecutorService unic = Executors.newSingleThreadExecutor();

// 4. PROGRAMAT: tasques retardades i periodiques.
ScheduledExecutorService programat = Executors.newScheduledThreadPool(2);

// 5. ROBATORI DE FEINA: ForkJoinPool amb parallelisme = nuclis.
ExecutorService robatori = Executors.newWorkStealingPool();
Fàbrica Fils Cua Quan fer-la servir Risc
newFixedThreadPool(n) Exactament n Il·limitada Càrrega estable, càlcul La cua creix sense límit
newCachedThreadPool() 0 … Integer.MAX_VALUE SynchronousQueue (sense capacitat) Moltes tasques curtes d'E/S Crea fils sense límit
newSingleThreadExecutor() 1 Il·limitada Tasques que han d'anar en ordre La cua creix sense límit
newScheduledThreadPool(n) n (creix si cal) Cua amb retard Feina periòdica Una tasca llarga retarda les altres
newWorkStealingPool() = nuclis Cues per fil Càlcul divisible No garanteix ordre
newVirtualThreadPerTaskExecutor() Un de virtual per tasca (Java 21) Milers de tasques d'E/S Veure 10-06

Casos d'ús a BiblioTech:

// Enviament d'avisos: E/S, carrega acotada i coneguda -> pool fix.
ExecutorService avisos = Executors.newFixedThreadPool(8);

// Escriptura del log d'auditoria: HA de conservar l'ordre -> un sol fil.
// A mes, en ser un unic fil, no cal sincronitzar el fitxer.
ExecutorService auditoria = Executors.newSingleThreadExecutor();

// Avis de venciments cada hora i copia de seguretat cada nit.
ScheduledExecutorService manteniment = Executors.newScheduledThreadPool(2);

El cas de newSingleThreadExecutor mereix un comentari: un executor d'un sol fil és confinament (08-04, apartat 17) implementat com a servei. Tot el que s'hi executa està serialitzat per construcció, així que el seu estat no necessita cap pany. És una tècnica molt potent i poc utilitzada.

  1. L'advertiment sobre les cues il·limitades

Les fàbriques d'Executors són còmodes i tenen un parany seriós que cal conèixer abans de fer-les servir en producció.

newFixedThreadPool i newSingleThreadExecutor fan servir una LinkedBlockingQueue sense límit de capacitat.

Si les tasques arriben més ràpid del que el pool les consumeix, la cua creix indefinidament. No hi ha cap mecanisme que ho freni. El resultat, després de minuts o hores de funcionament aparentment normal, és un OutOfMemoryError.

// BOMBA DE RELLOTGERIA
ExecutorService pool = Executors.newFixedThreadPool(4);

// Cada tasca triga 100 ms; n'arriben 1.000 per segon.
// El pool en processa 40/s. S'acumulen 960 tasques per segon en memoria.
// En 10 minuts: ~576.000 objectes en cua. OutOfMemoryError.
while (hiHaPeticions()) {
    pool.submit(new TascaLenta(seguentPeticio()));
}

newCachedThreadPool té el problema simètric i pitjor: la seva cua és una SynchronousQueue, que no emmagatzema res, així que cada tasca que arriba quan tots els fils estan ocupats provoca la creació d'un fil nou, sense límit (fins a Integer.MAX_VALUE). Amb una ràfega de deu mil tasques lentes, la JVM intenta crear deu mil fils i mor amb OutOfMemoryError: unable to create native thread.

La solució: una cua acotada i una política de rebuig explícita, que és el que es construeix a l'apartat següent.

Regla professional. Les fàbriques d'Executors estan bé per a codi d'exemple, eines internes i tasques d'arrencada. Per a un servei que rep càrrega externa, construeix el teu ThreadPoolExecutor amb cua acotada. És la recomanació explícita de les guies d'estil de Google i de la majoria de manuals de la indústria.

  1. ThreadPoolExecutor a mida

Totes les fàbriques anteriors retornen, per sota, un ThreadPoolExecutor. Construir-lo directament dona control sobre els sis paràmetres que importen.

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

public class ExecutorBiblioTech {

    /**
     * ThreadFactory que posa nom als fils. Sense ella, els fils del pool
     * s'anomenen "pool-1-thread-1", que en un bolcat (08-03) no diu res.
     */
    static class FabricaFilsAnomenats implements ThreadFactory {
        private final String prefix;
        private final boolean dimoni;
        private final AtomicInteger comptador = new AtomicInteger(1);

        FabricaFilsAnomenats(String prefix, boolean dimoni) {
            this.prefix = prefix;
            this.dimoni = dimoni;
        }

        @Override
        public Thread newThread(Runnable r) {
            Thread t = new Thread(r, prefix + "-" + comptador.getAndIncrement());
            t.setDaemon(dimoni);
            // Xarxa de seguretat: qualsevol excepcio que escapi d'una tasca
            // enviada amb execute() queda registrada (06-07, 08-02).
            t.setUncaughtExceptionHandler((fil, error) ->
                    System.err.println("[" + fil.getName() + "] fallada no capturada: " + error));
            return t;
        }
    }

    public static ThreadPoolExecutor crear() {
        return new ThreadPoolExecutor(
                4,                                    // 1. mida del nucli
                8,                                    // 2. mida maxima
                60L, TimeUnit.SECONDS,                // 3. vida dels fils extra
                new ArrayBlockingQueue<>(100),        // 4. cua ACOTADA
                new FabricaFilsAnomenats("bibliotech-avisos", false),     // 5.
                new ThreadPoolExecutor.CallerRunsPolicy()                 // 6.
        );
    }
}

Els sis paràmetres, un a un:

1. corePoolSize — fils que es mantenen vius encara que estiguin ociosos. El pool els crea a mesura que arriben tasques i no els destrueix (llevat que facis allowCoreThreadTimeOut(true)).

2. maximumPoolSize — sostre absolut de fils.

3. keepAliveTime — quant sobreviu un fil per damunt del nucli estant ociós.

4. workQueue — on esperen les tasques. La decisió més important.

5. threadFactory — com es creen els fils: nom, dimoni, prioritat, gestor d'excepcions.

6. RejectedExecutionHandler — què fer quan la cua està plena i no es poden crear més fils.

L'algoritme que segueix el pool davant d'una tasca nova —i que és on gairebé tothom s'equivoca en raonar:

flowchart TD
    A["Arriba una tasca"] --> B{"fils actius<br/>menor que corePoolSize?"}
    B -->|Sí| C["Crear un fil nou<br/>i executar"]
    B -->|No| D{"hi cap a la cua?"}
    D -->|Sí| E["Encuar<br/>ESPERA el seu torn"]
    D -->|No| F{"fils actius<br/>menor que maximumPoolSize?"}
    F -->|Sí| G["Crear un fil extra<br/>i executar"]
    F -->|No| H["REBUTJAR<br/>RejectedExecutionHandler"]

La conseqüència contraintuïtiva: el pool prefereix encuar abans que crear fils extra. Amb core=4, max=100 i una cua il·limitada, mai no es crearan més de 4 fils: la cua no s'omple mai, així que el pas que crea fils extra no s'assoleix mai. És exactament el que passa amb newFixedThreadPool, i explica per què maximumPoolSize només serveix d'alguna cosa si la cua està acotada.

Les quatre polítiques de rebuig:

Política Comportament Quan fer-la servir
AbortPolicy (per defecte) Llança RejectedExecutionException Quan perdre una tasca és inacceptable i vols assabentar-te'n
CallerRunsPolicy El fil que va enviar executa la tasca Excel·lent: crea contrapressió natural
DiscardPolicy Descarta en silenci Gairebé mai: les fallades silencioses són verí
DiscardOldestPolicy Descarta la més antiga en cua Dades on el més recent val més (mètriques)

CallerRunsPolicy mereix una explicació perquè és la més útil i la menys evident: quan el pool està saturat, la tasca l'executa el mateix fil que va cridar submit. Aquest fil es queda ocupat i, per tant, deixa de produir tasques noves durant aquell temps. El resultat és un mecanisme de contrapressió automàtic: el productor s'alenteix al ritme del consumidor, sense cues infinites i sense descartar res.

  1. Triar la mida del pool

Reprenent la distinció de 08-01 entre tasques limitades per CPU i per E/S:

Tipus de tasca Fórmula Exemple a BiblioTech
Limitada per CPU N o N + 1 Recalcular totes les multes
Limitada per E/S N × (1 + espera/càlcul) Enviar 200 avisos
Mixta Mesurar, o separar en dos pools Importar (llegir + validar)
public class DimensionamentPool {

    private static final int NUCLIS = Runtime.getRuntime().availableProcessors();

    /** Calcul pur: mes fils que nuclis nomes afegeix canvis de context. */
    public static int perCalcul() {
        return NUCLIS;
    }

    /**
     * E/S: la formula de Brian Goetz.
     *   fils = nuclis * utilitzacioObjectiu * (1 + espera / calcul)
     * Amb utilitzacio 1.0, 300 ms d'espera i 20 ms de calcul en 8 nuclis:
     *   8 * 1.0 * (1 + 15) = 128 fils.
     * A la practica s'acota a un maxim raonable: 128 fils son 128 MB
     * de piles, i el recurs remot probablement no aguanti 128 peticions.
     */
    public static int perEs(double msEspera, double msCalcul) {
        double ratio = msEspera / msCalcul;
        int suggerit = (int) Math.ceil(NUCLIS * (1 + ratio));
        return Math.min(suggerit, 64);       // sostre de seny
    }

    public static void main(String[] args) {
        System.out.println("Nuclis              : " + NUCLIS);
        System.out.println("Pool de calcul      : " + perCalcul());
        System.out.println("Pool d'avisos       : " + perEs(300, 20));
        System.out.println("Pool d'importacio   : " + perEs(50, 30));
    }
}

El consell més important sobre dimensionament: separa els pools per tipus de feina. Un únic pool compartit entre tasques de càlcul i d'E/S és el pitjor dels dos mons: les tasques d'E/S ocupen fils sense fer servir CPU, i les de càlcul bloquegen les d'E/S.

// BE: tres pools amb propositos i mides diferents, i aillats
// entre si: una saturacio del pool d'avisos no afecta el de calcul.
private final ExecutorService poolCalcul = Executors.newFixedThreadPool(NUCLIS);
private final ExecutorService poolEs     = Executors.newFixedThreadPool(32);
private final ScheduledExecutorService poolProgramat =
        Executors.newScheduledThreadPool(2);

Aquesta idea —aïllament per mampares (bulkhead)— ve dels compartiments estancs d'un vaixell: si un s'inunda, la resta continua a flot. Es reprèn a 12-07.

  1. El cicle de vida de l'executor

Un ExecutorService té tres estats: en marxa, apagant-se i acabat.

Mètode Què fa
shutdown() Deixa d'acceptar tasques noves; acaba les pendents. No bloqueja
shutdownNow() Deixa d'acceptar; interromp les que corren; retorna les no començades
isShutdown() S'ha cridat algun dels dos?
isTerminated() Han acabat ja totes les tasques?
awaitTermination(t, u) Bloqueja fins que acabi tot o s'exhaureixi el termini

L'error crític: si no apagues l'executor, la JVM no acaba. Els fils d'un pool són no dimoni per defecte, així que continuen vius esperant feina indefinidament (08-01, apartat 8). El teu main acaba i el procés es queda allà.

El patró d'apagada correcte, tal com el recomana la documentació d'ExecutorService:

import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;

public final class AturadaOrdenada {

    /**
     * Aturada en dues fases:
     *   1. shutdown() + espera cortes: deixar acabar el que esta en curs.
     *   2. Si no acaben, shutdownNow() (interrupcio) + espera curta.
     *   3. Si tot i aixi no acaben, registrar-ho: hi ha tasques que no
     *      responen a la interrupcio, i aixo es un BUG d'aquestes tasques
     *      (el catch buit d'InterruptedException de 08-02).
     */
    public static void aturar(ExecutorService executor, long segons) {
        executor.shutdown();                     // fase 1: no accepta mes
        try {
            if (!executor.awaitTermination(segons, TimeUnit.SECONDS)) {
                executor.shutdownNow();          // fase 2: interrompre
                if (!executor.awaitTermination(segons, TimeUnit.SECONDS)) {
                    System.err.println("L'executor no ha acabat: "
                            + "hi ha tasques que ignoren la interrupcio");
                }
            }
        } catch (InterruptedException e) {
            // Ens han interromput a NOSALTRES mentre esperavem.
            executor.shutdownNow();
            Thread.currentThread().interrupt();  // restaurar el marcador (08-02)
        }
    }
}

Detall important de shutdownNow(): interromp, no mata. Una tasca que s'empassa la InterruptedException continua corrent tan tranquil·la. Tot el protocol de 08-02 es cobra aquí: un catch buit fa que la teva aplicació no es pugui apagar.

Java 19+: ExecutorService és AutoCloseable. Des de Java 19 es pot fer servir amb try-with-resources (06-06), i close() fa un shutdown() i espera que acabin les tasques:

try (ExecutorService executor = Executors.newFixedThreadPool(8)) {
    for (Avis a : avisos) {
        executor.submit(new TascaAvis(a));
    }
}   // close(): shutdown() + espera indefinida

És molt més net, però espera indefinidament: si una tasca no acaba, el teu try no en surt mai. Per a codi amb terminis, el patró de dues fases continua sent el correcte.

  1. Callable i Future

A 08-02 va quedar pendent el problema: Runnable.run() retorna void i no pot llançar excepcions comprovades. Callable ho resol:

@FunctionalInterface
public interface Callable<V> {
    V call() throws Exception;      // retorna V i POT llancar
}
Runnable Callable<V>
Mètode void run() V call() throws Exception
Retorna valor No
Excepcions comprovades No
S'envia amb execute o submit submit

Future<V> és l'objecte que representa "un resultat que arribarà". El <V> és simplement el tipus del resultat: un Future<Informe> promet un Informe (els genèrics, a fons, a 10-01).

public interface Future<V> {
    V get() throws InterruptedException, ExecutionException;
    V get(long temps, TimeUnit unitat) throws ..., TimeoutException;
    boolean cancel(boolean interrompreSiEstaCorrent);
    boolean isCancelled();
    boolean isDone();
}

Exemple complet amb BiblioTech:

package com.nexussoftware.bibliotech.persistencia;

import java.nio.file.Path;
import java.util.concurrent.*;

public class ImportacioAmbFuture {

    public static void main(String[] args) {

        ExecutorService executor = Executors.newFixedThreadPool(2);
        try {
            // Callable<Informe>: retorna un Informe i pot llancar
            // excepcions comprovades. Totes dues coses impossibles amb Runnable.
            Callable<Informe> tasca = () -> {
                ImportadorCataleg imp = new ImportadorCataleg();
                return imp.importar(Path.of("dades/inventari.csv"));   // throws IOException
            };

            Future<Informe> futur = executor.submit(tasca);

            // El fil principal NO esta bloquejat: pot continuar treballant.
            System.out.println("[main] importacio llancada, atenc el menu");
            mostrarMenu();

            // isDone() no bloqueja: permet sondejar sense esperar.
            while (!futur.isDone()) {
                System.out.println("[main] important...");
                TimeUnit.MILLISECONDS.sleep(500);
            }

            // get() amb TERMINI: no facis servir mai get() sense termini en produccio.
            Informe informe = futur.get(30, TimeUnit.SECONDS);
            System.out.println("[main] importats: " + informe.importats());

        } catch (TimeoutException e) {
            System.err.println("[main] la importacio ha excedit el termini");

        } catch (ExecutionException e) {
            // La tasca ha llancat una excepcio. get() l'EMBOLCALLA en
            // ExecutionException; l'original es a getCause().
            // Es exactament l'encadenament de causes de 06-03.
            Throwable causa = e.getCause();
            System.err.println("[main] la importacio ha fallat: " + causa.getMessage());
            if (causa instanceof java.io.IOException) {
                System.err.println("[main] problema de fitxer; es conserva el cataleg anterior");
            }

        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();

        } finally {
            AturadaOrdenada.aturar(executor, 10);
        }
    }

    static void mostrarMenu() { /* ... */ }
}

Els tres punts que cal retenir:

1. get() sense termini bloqueja indefinidament. Si la tasca es penja, el teu fil es penja amb ella. En producció, sempre get(temps, unitat).

2. ExecutionException embolcalla la causa. L'excepció real és a getCause(). És l'encadenament de 06-03, i cal desembolcallar-la per prendre decisions:

try {
    Informe i = futur.get(30, TimeUnit.SECONDS);
} catch (ExecutionException e) {
    Throwable causa = e.getCause();
    // Recuperar el tipus original per aplicar la politica de 06-07
    if (causa instanceof FormatInvalidException fie) {
        registrarIDegradar(fie);
    } else if (causa instanceof java.io.IOException ioe) {
        avortarAmbCodi(ioe);
    } else {
        throw new BiblioTechException("Fallada inesperada a la importacio", causa);
    }
}

3. Les tres excepcions de get() signifiquen coses diferents:

Excepció Significat Reacció
ExecutionException La tasca ha fallat Desembolcallar amb getCause() i aplicar la política
TimeoutException La tasca continua corrent Decidir: esperar més o cancel(true)
InterruptedException Tu has estat interromput Restaurar el marcador; la tasca continua viva

Compte amb l'última fila: una TimeoutException no cancel·la res. La tasca continua executant-se. Si vols que s'aturi, cal cridar cancel(true) explícitament.

  1. Cancel·lació de tasques

Future.cancel(boolean interrompreSiEstaCorrent) és la cancel·lació de 08-02 exposada com a mètode.

Future<Informe> futur = executor.submit(tasca);

// cancel(false): nomes evita que COMENCI si encara es a la cua.
// Si ja esta corrent, la deixa acabar.
futur.cancel(false);

// cancel(true): a mes, INTERROMP el fil que l'executa.
// Es literalment un fil.interrupt() sobre el fil del pool.
futur.cancel(true);

cancel(true) funciona només si la tasca coopera. Tot el protocol de 08-02 s'aplica igual: comprovar isInterrupted(), capturar InterruptedException, no empassar-se-la. Una tasca que ignora la interrupció és incancel·lable, estigui en un pool o no.

/**
 * Tasca cancellable dins d'un pool. Identica al patro de 08-02:
 * cancel(true) del Future es tradueix en interrupt() sobre aquest fil.
 */
public class TascaAvisCancellable implements Callable<ResultatAvis> {

    private final Avis avis;

    public TascaAvisCancellable(Avis avis) { this.avis = avis; }

    @Override
    public ResultatAvis call() throws Exception {
        for (int intent = 1; intent <= 3; intent++) {

            // Punt de cancellacio a cada volta.
            if (Thread.currentThread().isInterrupted()) {
                throw new InterruptedException("avis cancellat abans de l'intent " + intent);
            }
            try {
                enviar(avis);                        // pot llancar InterruptedException
                return ResultatAvis.exit(avis);
            } catch (EnviamentFallitException e) {
                if (intent == 3) {
                    return ResultatAvis.fallada(avis, e.getMessage());
                }
                TimeUnit.MILLISECONDS.sleep(100L * intent);   // retroces
            }
        }
        return ResultatAvis.fallada(avis, "reintents exhaurits");
    }

    private void enviar(Avis a) throws Exception { /* ... */ }
}

Comprovació de l'estat després de cancel·lar:

Future<Informe> f = executor.submit(tasca);
TimeUnit.SECONDS.sleep(2);

boolean cancellada = f.cancel(true);
System.out.println("cancel() retorna  : " + cancellada);   // false si ja havia acabat
System.out.println("isCancelled()     : " + f.isCancelled());
System.out.println("isDone()          : " + f.isDone());   // true: cancellada tambe es "feta"

try {
    f.get();
} catch (CancellationException e) {
    // get() sobre un Future CANCELLAT llanca CancellationException,
    // que es NO comprovada. Es facil oblidar-se de capturar-la.
    System.out.println("confirmat: la tasca ha estat cancellada");
}

  1. invokeAll i invokeAny

Dos mètodes per enviar una col·lecció de tasques de cop.

invokeAll: executa totes i bloqueja fins que totes acabin. Retorna la llista de Future, tots ja completats.

List<Callable<ResultatAvis>> tasques = new ArrayList<>();
for (Avis a : avisos) {
    tasques.add(new TascaAvisCancellable(a));
}

// Bloqueja fins que TOTES acabin. Els Future retornats
// ja estan complets: get() no bloqueja.
List<Future<ResultatAvis>> resultats = executor.invokeAll(tasques);

int exits = 0, fallades = 0;
for (Future<ResultatAvis> f : resultats) {
    try {
        if (f.get().correcte()) exits++; else fallades++;
    } catch (ExecutionException e) {
        fallades++;
        LOG.log(Level.WARNING, "avis fallit", e.getCause());
    }
}

Amb termini global, que és el recomanable:

// Termini GLOBAL per al conjunt. Les que no acabin a temps
// es CANCELLEN automaticament, i el seu get() llancara CancellationException.
List<Future<ResultatAvis>> resultats =
        executor.invokeAll(tasques, 30, TimeUnit.SECONDS);

invokeAny: retorna el resultat de la primera que acabi amb èxit i cancel·la les altres. Útil quan hi ha diverses formes d'obtenir el mateix i vols la més ràpida.

// Tres fonts per al mateix cataleg; ens val la primera que respongui.
List<Callable<Cataleg>> fonts = List.of(
        () -> carregarDesDeCacheLocal(),
        () -> carregarDesDelFitxerPrincipal(),
        () -> carregarDesDeCopiaSeguretat()
);

// Retorna el primer resultat correcte; cancella els altres dos.
// Si TOTES fallen, llanca ExecutionException amb l'ultima causa.
Cataleg c = executor.invokeAny(fonts, 10, TimeUnit.SECONDS);
invokeAll invokeAny
Espera Totes La primera amb èxit
Retorna List<Future<T>> T (el valor, no un Future)
Amb les altres Res Les cancel·la
Si alguna falla El seu Future llança a get() S'ignora, llevat que fallin totes
Cas d'ús Feina en lot Fonts redundants, la més ràpida

  1. ScheduledExecutorService

Per a feina retardada o periòdica. Substitueix Timer/TimerTask, que són de Java 1.3 i tenen dos defectes greus: un sol fil per a totes les tasques, i una excepció no capturada mata el temporitzador sencer i les altres tasques no es tornen a executar mai.

import java.util.concurrent.*;

public class MantenimentBiblioTech {

    private final ScheduledExecutorService programador =
            Executors.newScheduledThreadPool(2, r -> {
                Thread t = new Thread(r, "bibliotech-programador");
                t.setDaemon(true);      // manteniment: no ha d'impedir el tancament
                return t;
            });

    public void arrencar() {

        // 1. UNA vegada, d'aqui a 5 segons.
        programador.schedule(
                () -> System.out.println("[manteniment] arrencada completada"),
                5, TimeUnit.SECONDS);

        // 2. PERIODICA a ritme fix: cada hora des de l'instant d'inici.
        programador.scheduleAtFixedRate(
                this::avisarVenciments,
                0, 1, TimeUnit.HOURS);

        // 3. PERIODICA amb retard fix: 24 h des que ACABA l'anterior.
        programador.scheduleWithFixedDelay(
                this::copiaSeguretat,
                1, 24, TimeUnit.HOURS);
    }

    /**
     * REGLA D'OR: una tasca periodica ha de capturar TOTA excepcio.
     * Si en deixa escapar una, la JVM CANCELLA la planificacio i la tasca
     * NO ES TORNA A EXECUTAR MAI, en silenci, sense cap avis.
     * Es la fallada mes traidora d'aquesta API.
     */
    private void avisarVenciments() {
        try {
            int enviats = serveiAvisos.enviarVenciments();
            LOG.log(Level.INFO, "Avisos de venciment enviats: {0}", enviats);
        } catch (Exception e) {
            LOG.log(Level.SEVERE, "Fallada en avisar dels venciments", e);
            // NO rellancar: si surt, s'acaba la planificacio per sempre.
        }
    }

    private void copiaSeguretat() {
        try {
            gestorCopies.copiar();
        } catch (Exception e) {
            LOG.log(Level.SEVERE, "Fallada a la copia de seguretat", e);
        }
    }
}

scheduleAtFixedRate davant de scheduleWithFixedDelay — la diferència importa:

scheduleAtFixedRate(t, i, p, u) scheduleWithFixedDelay(t, i, p, u)
El període es mesura des de L'inici de l'execució anterior El final de l'execució anterior
Si la tasca dura més que el període Les següents es retarden i encadenen sense solapar-se Sempre hi ha p de separació real
Ritme Constant (intenta complir l'horari) Variable (depèn del que duri)
Fer servir per a Feina amb horari: informes horaris, mètriques Feina la durada de la qual varia: còpies, neteja
Tasca que dura 3 s, periode 5 s:

scheduleAtFixedRate:   [---3s---]__2s__[---3s---]__2s__[---3s---]
                       0         3     5         8     10
                       inici cada 5 s exactes

scheduleWithFixedDelay:[---3s---]____5s____[---3s---]____5s____[--
                       0         3         8         11        16
                       5 s DESPRES d'acabar

El detall que arruïna sistemes de producció: si una tasca de scheduleAtFixedRate llança una excepció no capturada, la planificació es cancel·la silenciosament. La tasca no es torna a executar mai, no hi ha error al log llevat del que hi posis tu, i ningú no se n'assabenta fins que algú pregunta per què fa tres setmanes que no arriben els avisos. Embolcalla sempre el cos en un try/catch (Exception).

  1. CountDownLatch

Un tancament de compte enrere: un comptador que només baixa. Els fils que esperen a await() es desbloquegen quan arriba a zero. És d'un sol ús: no es pot reiniciar.

CountDownLatch tancament = new CountDownLatch(5);   // 5 esdeveniments pendents

tancament.countDown();          // resta 1 (mai per sota de 0)
tancament.await();              // bloqueja fins a arribar a 0
tancament.await(10, TimeUnit.SECONDS);   // amb termini; retorna false si expira
tancament.getCount();           // compte actual (per mostrar el progres)

Ús 1: esperar que acabin N tasques.

import java.util.concurrent.*;

public class EsperarNTasques {

    public static void main(String[] args) throws InterruptedException {

        final int TASQUES = 200;
        ExecutorService pool = Executors.newFixedThreadPool(8);
        CountDownLatch acabades = new CountDownLatch(TASQUES);

        for (int i = 0; i < TASQUES; i++) {
            final int n = i;
            pool.execute(() -> {
                try {
                    enviarAvis(n);
                } catch (Exception e) {
                    LOG.log(Level.WARNING, "avis " + n + " fallit", e);
                } finally {
                    // OBLIGATORI al finally: si una tasca falla i no
                    // fa countDown, l'await() de main no torna MAI.
                    acabades.countDown();
                }
            });
        }

        // Progres mentre esperem, sense sondejar amb sleep a cegues.
        while (!acabades.await(500, TimeUnit.MILLISECONDS)) {
            long pendents = acabades.getCount();
            System.out.printf("\rProgres: %d/%d (%.0f%%)",
                    TASQUES - pendents, TASQUES,
                    100.0 * (TASQUES - pendents) / TASQUES);
        }
        System.out.println("\nTots els avisos processats");
        AturadaOrdenada.aturar(pool, 10);
    }
}

El countDown() va SEMPRE en un finally. Si una tasca falla abans de cridar-lo, el comptador no arriba mai a zero i qui espera es queda bloquejat per sempre. És l'error número u de CountDownLatch.

Ús 2: porta de sortida — arrencar N fils exactament alhora. Molt útil per a proves de concurrència, perquè maximitza el solapament:

/**
 * Dos tancaments: un per donar la sortida a tots alhora, l'altre per
 * esperar que tots acabin. Es l'esquelet d'una prova d'estres
 * de concurrencia ben feta: sense la porta, els primers fils
 * acabarien abans que arrenquessin els ultims, i no hi hauria
 * solapament real (l'efecte que vas veure a l'exercici 2 de 08-01).
 */
public static long provaEstres(int fils, Runnable tasca) throws InterruptedException {

    CountDownLatch sortida = new CountDownLatch(1);       // la pistola
    CountDownLatch acabats = new CountDownLatch(fils);    // la meta

    for (int i = 0; i < fils; i++) {
        new Thread(() -> {
            try {
                sortida.await();     // tots esperen aqui
                tasca.run();
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            } finally {
                acabats.countDown();
            }
        }, "estres-" + i).start();
    }

    long inici = System.nanoTime();
    sortida.countDown();             // ja! els N arrenquen alhora
    acabats.await();
    return System.nanoTime() - inici;
}

  1. CyclicBarrier

Una barrera cíclica: N fils s'esperen mútuament en un punt; quan arriba l'últim, tots continuen i la barrera es reinicia. Aquesta última paraula és la diferència amb CountDownLatch.

import java.util.concurrent.*;

public class ProcessamentPerFases {

    public static void main(String[] args) {

        final int TREBALLADORS = 4;

        // L'accio de barrera l'executa l'ULTIM fil a arribar,
        // abans d'alliberar els altres. Ideal per consolidar la fase.
        CyclicBarrier barrera = new CyclicBarrier(TREBALLADORS, () ->
                System.out.println("--- fase completada pels 4, consolidant ---"));

        for (int i = 0; i < TREBALLADORS; i++) {
            final int id = i;
            new Thread(() -> {
                try {
                    for (int fase = 1; fase <= 3; fase++) {
                        System.out.printf("  treballador %d processa la fase %d%n", id, fase);
                        TimeUnit.MILLISECONDS.sleep(100 + id * 50);
                        barrera.await();      // espera els altres tres
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } catch (BrokenBarrierException e) {
                    // Si un fil s'interromp o expira, la barrera es TRENCA
                    // i TOTS els altres reben aquesta excepcio. Es correcte:
                    // sense tots els participants, la fase no te sentit.
                    System.err.println("barrera trencada: " + e);
                }
            }, "treballador-" + i).start();
        }
    }
}
CountDownLatch CyclicBarrier
Reutilitzable No (un sol ús) (es reinicia)
Qui espera Uns esperen, altres compten Tots els participants esperen
El comptador Baixa amb countDown() Puja en arribar a await()
Acció en completar-se No en té , Runnable opcional
Si un participant falla Els altres continuen esperant Es trenca per a tothom
Cas típic Esperar que acabi un lot Simulacions i càlculs per fases

A BiblioTech, CountDownLatch és el que necessites gairebé sempre: els 200 avisos són independents i ningú no ha d'esperar ningú a mig camí. CyclicBarrier encaixa en càlculs iteratius per fases, on cada iteració necessita que totes les anteriors hagin acabat.

  1. Semaphore

Un semàfor manté un nombre de permisos. acquire() en pren un (i bloqueja si no n'hi ha); release() en torna un. És l'eina per limitar la concurrència sobre un recurs.

package com.nexussoftware.bibliotech.servei;

import java.util.concurrent.*;

/**
 * Limita a 3 les exportacions simultanies.
 *
 * Motiu: cada exportacio obre el fitxer complet del cataleg, el
 * formata en memoria i l'escriu. Tres alhora saturen el disc i
 * multipliquen l'us de memoria; amb vint, el servidor s'arrossega.
 * El semafor posa un sostre dur sense limitar la resta del pool.
 */
public class ServeiExportacio {

    private static final int MAX_SIMULTANIES = 3;

    // 3 permisos. El 'true' activa l'equitat: els que porten mes temps
    // esperant passen abans, cosa que evita la inanicio (08-04).
    private final Semaphore permisos = new Semaphore(MAX_SIMULTANIES, true);

    private final ExecutorService pool = Executors.newFixedThreadPool(16);

    public Future<Path> exportar(String format) {
        return pool.submit(() -> {

            // Espera un permis. Si hi ha 3 exportacions en curs, aquest
            // fil es bloqueja aqui fins que una acabi.
            permisos.acquire();
            try {
                System.out.printf("[%s] exportant (%d permisos lliures)%n",
                        Thread.currentThread().getName(),
                        permisos.availablePermits());
                return generarFitxer(format);        // feina pesada
            } finally {
                // OBLIGATORI al finally: un permis no retornat es
                // perd PER SEMPRE. Despres de tres fallades, el semafor
                // queda a 0 i ningu no torna a exportar mai.
                permisos.release();
            }
        });
    }

    /** Variant que no espera: rebutja en lloc d'encuar. */
    public Future<Path> exportarSenseEsperar(String format) {
        return pool.submit(() -> {
            if (!permisos.tryAcquire(5, TimeUnit.SECONDS)) {
                throw new ServeiSaturatException(
                        "Hi ha " + MAX_SIMULTANIES + " exportacions en curs; reintenteu-ho");
            }
            try {
                return generarFitxer(format);
            } finally {
                permisos.release();
            }
        });
    }

    private Path generarFitxer(String format) throws Exception {
        TimeUnit.SECONDS.sleep(2);
        return Path.of("informes/cataleg." + format);
    }
}

Usos habituals del semàfor:

Ús Com
Limitar concurrència sobre un recurs new Semaphore(n), acquire/release
Convertir una col·lecció en acotada Un permís per forat disponible
Exclusió mútua simple new Semaphore(1) — però prefereix un pany
Senyalitzar entre fils Un fil fa release, un altre acquire

Diferència amb un pany, que es pregunta sempre: un pany és propietat del fil que el va prendre i només ell el pot deixar anar; un semàfor no té amo, així que un fil pot adquirir un permís i un altre alliberar-lo. Això el fa més flexible i també més fàcil de fer servir malament —d'aquí el finally obligatori—.

  1. Exchanger i taula comparativa de sincronitzadors

Exchanger<V> permet que dos fils intercanviïn objectes en un punt de trobada. Tots dos criden exchange(objecte) i cadascun rep el de l'altre.

Exchanger<List<Material>> intercanvi = new Exchanger<>();

// Fil lector: omple una memoria intermedia i la intercanvia per una buida.
List<Material> memoriaPropia = new ArrayList<>();
// ... omplir ...
memoriaPropia = intercanvi.exchange(memoriaPropia);   // rep la buida de l'altre

// Fil processador: processa una plena i en retorna una de buida.
List<Material> aProcessar = intercanvi.exchange(new ArrayList<>());

És un nínxol estret —el patró de doble memòria intermèdia entre exactament dos fils— i a la pràctica gairebé sempre es resol millor amb una BlockingQueue (08-06). Convé conèixer-lo per si apareix.

Taula comparativa completa:

Sincronitzador Participants Reutilitzable Per a què
CountDownLatch Uns compten, altres esperen No Esperar que passin N esdeveniments
CyclicBarrier N, tots iguals Sincronitzar fases entre N fils
Semaphore Qualsevol Limitar accessos concurrents
Exchanger Exactament 2 Intercanvi d'objectes per parelles
Phaser (Java 7) Variable, dinàmic Barrera amb participants que entren i surten

Phaser és una CyclicBarrier més flexible que permet registrar i desregistrar participants sobre la marxa. És potent i poc freqüent; menciona'l i no el facis servir llevat que la necessitat sigui evident.

  1. ForkJoinPool i divideix i venceràs

ForkJoinPool (Java 7) està pensat per a tasques que es poden partir recursivament en subtasques independents: el model divideix i venceràs.

La seva característica distintiva és el robatori de feina (work stealing): cada fil té la seva pròpia cua de subtasques, i quan es queda sense feina, roba de la cua d'un altre fil per l'extrem oposat. Això reparteix la càrrega automàticament sense un coordinador central i amb molt poca contenció.

flowchart TD
    A["Calcular multes de 100.000 prestecs"] --> B["Tros 1: 0-50.000"]
    A --> C["Tros 2: 50.000-100.000"]
    B --> D["0-25.000"]
    B --> E["25.000-50.000"]
    C --> F["50.000-75.000"]
    C --> G["75.000-100.000"]
    D --> H["... fins al llindar<br/>càlcul directe"]
    E --> H
    F --> H
    G --> H
    H --> I["Combinar resultats<br/>cap amunt"]

RecursiveTask<V> per a tasques que retornen valor; RecursiveAction per a les que no.

package com.nexussoftware.bibliotech.servei;

import java.util.List;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.RecursiveTask;

/**
 * Calcula la suma de multes d'una llista de prestecs dividint
 * la feina recursivament.
 *
 * RecursiveTask<Double>: el <Double> es el tipus del resultat
 * que produeix aquesta tasca (10-01 per als generics).
 */
public class CalculMultesParallel extends RecursiveTask<Double> {

    /**
     * LLINDAR: per sota d'aquesta mida es calcula directament.
     * Triar-lo be es el mes important de tot el patro:
     *  - massa baix: el cost de crear i coordinar subtasques
     *    supera el del calcul, i va MES LENT que en sequencial;
     *  - massa alt: no hi ha prou trossos per a tots els nuclis.
     * Regla practica: entre 100 i 10.000 elements, i MESURAR.
     */
    private static final int LLINDAR = 1_000;

    private final List<Prestec> prestecs;
    private final int des;
    private final int fins;
    private final CalculadoraMultes calculadora;

    public CalculMultesParallel(List<Prestec> prestecs, int des, int fins,
                                CalculadoraMultes calculadora) {
        this.prestecs = prestecs;
        this.des = des;
        this.fins = fins;
        this.calculadora = calculadora;
    }

    @Override
    protected Double compute() {
        int mida = fins - des;

        // CAS BASE: tros petit, calcul sequencial directe.
        if (mida <= LLINDAR) {
            double suma = 0;
            for (int i = des; i < fins; i++) {
                suma += calculadora.calcular(prestecs.get(i));
            }
            return suma;
        }

        // CAS RECURSIU: partir per la meitat.
        int mig = des + mida / 2;
        CalculMultesParallel esq =
                new CalculMultesParallel(prestecs, des, mig, calculadora);
        CalculMultesParallel dre =
                new CalculMultesParallel(prestecs, mig, fins, calculadora);

        // PATRO CORRECTE: bifurcar UNA i calcular l'altra en AQUEST fil.
        // Fer esq.fork() + dre.fork() + esq.join() + dre.join()
        // desaprofita el fil actual, que es quedaria esperant.
        esq.fork();                       // a la cua: un altre fil la pot robar
        double resultatDre = dre.compute();    // aquest fil treballa
        double resultatEsq = esq.join();       // recollir la bifurcada

        return resultatEsq + resultatDre;
    }

    // --- Us ---
    public static double calcularTotes(List<Prestec> prestecs,
                                       CalculadoraMultes calculadora) {
        // El pool COMU: compartit per tota la JVM, amb
        // (nuclis - 1) fils. El fan servir tambe els streams parallels (10-04).
        ForkJoinPool comu = ForkJoinPool.commonPool();
        return comu.invoke(new CalculMultesParallel(
                prestecs, 0, prestecs.size(), calculadora));
    }
}

El pool comú (ForkJoinPool.commonPool()) és una instància compartida per tota la JVM, amb nuclis - 1 fils, que es fa servir per defecte per a les tasques fork/join i per a CompletableFuture (08-07) i els streams paral·lels (10-04).

El seu gran perill: com que és compartit per tota l'aplicació, una tasca que hi bloquegi —una lectura de fitxer lenta, un get() que espera— roba un fil a tots els altres usuaris del pool. Amb nuclis - 1 fils, n'hi ha prou amb unes poques tasques blocants per deixar-lo inservible.

// PROHIBIT al pool comu: operacions blocants.
ForkJoinPool.commonPool().submit(() -> {
    return Files.readAllLines(camiEnorme);   // bloqueja un fil compartit
});

// CORRECTE: pool propi per a la feina blocant.
ForkJoinPool poolPropi = new ForkJoinPool(4);
try {
    poolPropi.invoke(new CalculMultesParallel(...));
} finally {
    poolPropi.shutdown();
}

Nota sobre parallelStream(). Els streams paral·lels de l'API de col·leccions fan servir exactament aquest ForkJoinPool.commonPool() per sota, i converteixen el patró anterior en una sola línia. No són matèria d'aquest mòdul: s'expliquen a 10-04, juntament amb tota l'API de Streams, on també es discuteix quan compensen i quan són contraproduents.

Nota sobre fils virtuals. Tota l'aritmètica de dimensionament de pools d'aquesta lliçó —N fils per a CPU, N × (1 + espera/càlcul) per a E/S, cues acotades, polítiques de rebuig— existeix perquè un fil de plataforma és car. Els fils virtuals de Java 21 costen nanosegons i uns centenars de bytes, i amb ells el patró recomanat per a tasques d'E/S passa a ser Executors.newVirtualThreadPerTaskExecutor(): un fil per tasca, sense pool i sense dimensionar res. Per a tasques de CPU, en canvi, els pools acotats continuen sent el correcte. S'expliquen a 10-06.

  1. BiblioTech: 200 avisos amb pool, progrés i cancel·lació

Tot junt. És el cas B de 08-01, resolt de veritat.

package com.nexussoftware.bibliotech.servei;

import com.nexussoftware.bibliotech.domini.Prestec;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;
import java.util.logging.Logger;

/**
 * Enviament massiu d'avisos de venciment.
 *
 * Propietats:
 *  - Pool ACOTAT amb fils anomenats i cua acotada.
 *  - SEMAFOR que limita a 5 els accessos simultanis al fitxer d'avisos.
 *  - COUNTDOWNLATCH per mostrar el progres i esperar el final.
 *  - Cancellacio neta mitjancant shutdownNow() -> interrupcio cooperativa.
 *  - Resultat agregat amb comptadors atomics (detall a 08-06).
 */
public class ServeiAvisos implements AutoCloseable {

    private static final Logger LOG = Logger.getLogger(ServeiAvisos.class.getName());

    private static final int FILS = 8;
    private static final int MAX_ESCRIPTURES_SIMULTANIES = 5;

    private final ThreadPoolExecutor pool;
    private final Semaphore accesFitxer =
            new Semaphore(MAX_ESCRIPTURES_SIMULTANIES, true);

    public ServeiAvisos() {
        AtomicInteger n = new AtomicInteger(1);
        this.pool = new ThreadPoolExecutor(
                FILS, FILS,
                0L, TimeUnit.MILLISECONDS,
                new ArrayBlockingQueue<>(500),          // cua ACOTADA
                r -> {
                    Thread t = new Thread(r, "bibliotech-avisos-" + n.getAndIncrement());
                    t.setUncaughtExceptionHandler((f, e) ->
                            LOG.log(Level.SEVERE, "Fallada no capturada a " + f.getName(), e));
                    return t;
                },
                new ThreadPoolExecutor.CallerRunsPolicy());  // contrapressio
    }

    /** Resultat agregat de l'enviament. Immutable (08-04). */
    public record ResumEnviament(int total, int enviats, int fallits,
                                 int cancellats, long milisegons) {

        public double percentatgeExit() {
            return total == 0 ? 0 : 100.0 * enviats / total;
        }
    }

    /**
     * Envia tots els avisos en parallel.
     *
     * @param terminiSegons termini global; en exhaurir-se es cancella el pendent
     */
    public ResumEnviament enviarTots(List<Prestec> vencuts, long terminiSegons)
            throws InterruptedException {

        final int total = vencuts.size();
        final CountDownLatch acabats = new CountDownLatch(total);
        final AtomicInteger enviats = new AtomicInteger();
        final AtomicInteger fallits = new AtomicInteger();

        long inici = System.nanoTime();
        List<Future<?>> futurs = new ArrayList<>(total);

        // 1. Enviar totes les tasques al pool.
        for (Prestec p : vencuts) {
            futurs.add(pool.submit(() -> {
                try {
                    // Punt de cancellacio abans de comencar.
                    if (Thread.currentThread().isInterrupted()) {
                        return;
                    }
                    // El semafor limita l'acces concurrent al fitxer.
                    accesFitxer.acquire();
                    try {
                        escriureAvis(p);
                        enviats.incrementAndGet();
                    } finally {
                        accesFitxer.release();     // SEMPRE
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();     // restaurar (08-02)
                } catch (Exception e) {
                    fallits.incrementAndGet();
                    LOG.log(Level.WARNING, "Avis fallit per a " + p.id(), e);
                } finally {
                    // SEMPRE, o l'await() no torna mai.
                    acabats.countDown();
                }
            }));
        }

        // 2. Progres mentre s'espera, amb termini global.
        long limitNs = TimeUnit.SECONDS.toNanos(terminiSegons);
        boolean completat = false;

        while (!(completat = acabats.await(300, TimeUnit.MILLISECONDS))) {
            long fets = total - acabats.getCount();
            System.out.printf("\r  [%s] %d/%d avisos (%.0f%%)",
                    barra(fets, total), fets, total, 100.0 * fets / total);

            if (System.nanoTime() - inici > limitNs) {
                System.out.println("\n  Termini exhaurit: cancellant el pendent");
                for (Future<?> f : futurs) {
                    f.cancel(true);        // interromp o descarta de la cua
                }
                break;
            }
        }

        long ms = (System.nanoTime() - inici) / 1_000_000;
        int cancellats = total - enviats.get() - fallits.get();

        System.out.printf("\r  [%s] %d/%d avisos (100%%)%n",
                barra(total, total), total, total);

        return new ResumEnviament(total, enviats.get(), fallits.get(), cancellats, ms);
    }

    private static String barra(long fets, long total) {
        int amplada = 30;
        int plens = total == 0 ? 0 : (int) (amplada * fets / total);
        return "#".repeat(plens) + "-".repeat(amplada - plens);
    }

    /** Simula l'escriptura de l'avis: 250 ms d'E/S. */
    private void escriureAvis(Prestec p) throws InterruptedException {
        TimeUnit.MILLISECONDS.sleep(250);
        if (p.id().hashCode() % 37 == 0) {
            throw new IllegalStateException("empleat sense adreca de contacte");
        }
    }

    /** AutoCloseable: aturada en dues fases (06-06 + apartat 7). */
    @Override
    public void close() {
        pool.shutdown();
        try {
            if (!pool.awaitTermination(30, TimeUnit.SECONDS)) {
                pool.shutdownNow();
                if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
                    LOG.severe("Hi ha tasques d'avis que ignoren la interrupcio");
                }
            }
        } catch (InterruptedException e) {
            pool.shutdownNow();
            Thread.currentThread().interrupt();
        }
    }
}

Ús:

package com.nexussoftware.bibliotech.presentacio;

import java.util.List;

public class DemostracioAvisos {

    public static void main(String[] args) throws InterruptedException {

        List<Prestec> vencuts = generarVencuts(200);

        System.out.println("=== ENVIAMENT D'AVISOS DE VENCIMENT ===");
        System.out.println("Avisos pendents: " + vencuts.size());
        System.out.println("Temps sequencial estimat: "
                + (vencuts.size() * 250 / 1000) + " s");
        System.out.println();

        // try-with-resources: close() atura el pool ordenadament.
        try (ServeiAvisos servei = new ServeiAvisos()) {

            ServeiAvisos.ResumEnviament r = servei.enviarTots(vencuts, 60);

            System.out.println();
            System.out.println("=== RESUM ===");
            System.out.println("Total      : " + r.total());
            System.out.println("Enviats    : " + r.enviats());
            System.out.println("Fallits    : " + r.fallits());
            System.out.println("Cancellats : " + r.cancellats());
            System.out.println("Temps      : " + r.milisegons() + " ms");
            System.out.printf ("Exit       : %.1f%%%n", r.percentatgeExit());
            System.out.printf ("Acceleracio: %.1fx sobre l'enviament sequencial%n",
                    (vencuts.size() * 250.0) / r.milisegons());
        }
    }
}

Sortida:

=== ENVIAMENT D'AVISOS DE VENCIMENT ===
Avisos pendents: 200
Temps sequencial estimat: 50 s

  [##############################] 200/200 avisos (100%)

=== RESUM ===
Total      : 200
Enviats    : 194
Fallits    : 6
Cancellats : 0
Temps      : 10247 ms
Exit       : 97.0%
Acceleracio: 4.9x sobre l'enviament sequencial

Anàlisi honesta del resultat, que és on hi ha l'aprenentatge:

  1. L'acceleració és 4,9× i no 8×, encara que el pool tingui 8 fils. La causa és el semàfor de 5 permisos: només 5 avisos poden escriure alhora, així que el paral·lelisme real està limitat per aquest 5, no pels 8 fils. 200 × 250 ms / 5 ≈ 10 s, exactament el que s'ha mesurat. El coll d'ampolla no és el pool: és el recurs protegit, i això és la norma, no l'excepció, en sistemes reals.
  2. Si puges el semàfor a 8, l'acceleració puja a ~8× i el temps baixa a 6,3 s. Si el puges a 20 sense pujar els fils, no canvia res: el límit passa a ser el pool. Dimensionar és trobar quin recurs és l'escàs.
  3. Les 6 fallades no van aturar l'enviament. Cada tasca captura les seves pròpies excepcions i fa countDown() al finally: la política de degradació de 06-07, ara en concurrència.
  4. Vuit fils, no dos-cents. 8 MB de piles en lloc de 200 MB, i vuit creacions de fil en lloc de dues-centes.

Errors Comuns i Consells

Error 1: no apagar l'executor. Els fils del pool són no dimoni; la JVM no acaba. És la causa número u de "el meu programa no acaba".

Error 2: fer servir newFixedThreadPool amb càrrega externa. Cua il·limitada → OutOfMemoryError després d'hores d'aparent normalitat. Cua acotada i política de rebuig.

Error 3: fer servir newCachedThreadPool amb moltes tasques lentes. Crea fils sense límit fins a tombar la JVM.

Error 4: enviar amb submit i no cridar get(). L'excepció de la tasca queda desada al Future i no apareix enlloc. Fallada silenciosa total. Fes servir execute si no t'interessa el resultat, o comprova el Future.

Error 5: get() sense termini en producció. Si la tasca es penja, qui crida es penja. Sempre get(temps, unitat).

Error 6: no desembolcallar ExecutionException. L'excepció real és a getCause(). Registrar l'ExecutionException tal qual amaga la informació útil.

Error 7: deixar que una tasca de scheduleAtFixedRate llanci una excepció. La planificació es cancel·la per sempre, en silenci. Embolcalla el cos en try/catch (Exception).

Error 8: countDown() fora del finally. Si la tasca falla abans, qui espera es bloqueja indefinidament.

Error 9: release() d'un semàfor fora del finally. Els permisos perduts no es recuperen; després d'unes fallades, el semàfor queda a zero i ningú no hi passa mai més.

Error 10: creure que maximumPoolSize s'assoleix amb cua il·limitada. El pool encua abans de crear fils extra: amb cua il·limitada mai no passa de corePoolSize.

Error 11: bloquejar al ForkJoinPool.commonPool(). És compartit per tota la JVM i té nuclis - 1 fils. Unes poques tasques blocants l'inutilitzen per a tothom.

Error 12: fer servir un únic pool per a càlcul i per a E/S. Es destorben mútuament. Pools separats per tipus de feina.

Consell 1: posa nom als fils amb una ThreadFactory. pool-1-thread-3 no diu res en un bolcat; bibliotech-avisos-3 sí. És el consell de 08-02 aplicat als pools.

Consell 2: fes servir CallerRunsPolicy quan no puguis perdre tasques. Crea contrapressió automàtica: el productor es frena sol.

Consell 3: mesura abans de dimensionar. Les fórmules són punts de partida. La mida real depèn de la teva càrrega, el teu maquinari i les teves dependències externes.

Consell 4: separa pools per tipus de feina (mampares). Aïlla les fallades: que la saturació dels avisos no impedeixi generar un informe.

Consell 5: un executor d'un sol fil és confinament gratis. Tot el que hi corre està serialitzat, així que el seu estat no necessita panys. Molt útil per a escriptures en fitxer i per a logs ordenats.

Consell 6: recorda que cancel(true) és interrupt(). Sense tasques cooperatives (08-02), la cancel·lació no funciona, estiguis en un pool o no.

Exercicis

Exercici 1: Comparativa de pools

Escriu ComparativaPools que executi la mateixa càrrega —300 tasques que dormen 100 ms cadascuna, simulant E/S— amb cinc estratègies: (a) un fil per tasca a mà, (b) newFixedThreadPool(4), (c) newFixedThreadPool(32), (d) newCachedThreadPool(), i (e) un ThreadPoolExecutor a mida amb 16 fils, cua de 50 i CallerRunsPolicy. Per a cadascuna mesura el temps total amb System.nanoTime() i compta quants fils diferents van executar tasques (fent servir un Set sincronitzat de noms de fil). Imprimeix una taula i explica els resultats.

Exercici 2: Importació amb Callable, termini i cancel·lació

Reescriu la importació de catàleg de BiblioTech com a Callable<Informe> executat en un ExecutorService. El programa ha de:

  1. Llançar la importació i continuar mostrant un menú simulat.
  2. Mostrar el progrés mentre !futur.isDone().
  3. Aplicar un termini de 5 segons amb get(5, TimeUnit.SECONDS).
  4. En exhaurir-se el termini, cridar cancel(true) i verificar-ho amb isCancelled().
  5. Distingir al catch els tres finals: ExecutionException (desembolcallant la causa), TimeoutException i InterruptedException.
  6. Apagar l'executor amb el patró de dues fases.

Exercici 3: Sala de lectura amb semàfor i tancament

Modela la sala de lectura de BiblioTech: hi caben 4 empleats alhora i hi ha 20 empleats que hi volen entrar. Escriu SalaDeLectura amb:

  1. Un Semaphore(4, true) que controla l'aforament, amb entrar()/sortir() i el release al finally.
  2. Un CountDownLatch de "porta de sortida" que faci que els 20 fils intentin entrar exactament alhora, i un altre de "meta" per esperar que tots hagin acabat.
  3. Un comptador atòmic de l'aforament instantani, i una comprovació que falli sorollosament si mai supera 4.
  4. Un ScheduledExecutorService que imprimeixi l'aforament actual cada 200 ms mentre duri la simulació, i que s'apagui al final.

Cada empleat roman a la sala un temps aleatori entre 200 i 600 ms.

Solucions

Solució a l'Exercici 1

import java.util.Collections;
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.*;

public class ComparativaPools {

    static final int TASQUES = 300;
    static final long DURACIO_MS = 100;

    /** Tasca d'E/S simulada que registra quin fil l'ha executada. */
    static Runnable tasca(Set<String> filsUsats) {
        return () -> {
            filsUsats.add(Thread.currentThread().getName());
            try {
                TimeUnit.MILLISECONDS.sleep(DURACIO_MS);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        };
    }

    record Resultat(String estrategia, long ms, int fils) { }

    /** (a) Un fil per tasca, gestionat a ma. */
    static Resultat unFilPerTasca() throws InterruptedException {
        Set<String> fils = Collections.synchronizedSet(new HashSet<>());
        Thread[] ts = new Thread[TASQUES];
        long inici = System.nanoTime();
        for (int i = 0; i < TASQUES; i++) {
            ts[i] = new Thread(tasca(fils), "manual-" + i);
            ts[i].start();
        }
        for (Thread t : ts) t.join();
        return new Resultat("un fil per tasca",
                (System.nanoTime() - inici) / 1_000_000, fils.size());
    }

    /** (b)-(e) Qualsevol ExecutorService. */
    static Resultat ambExecutor(String nom, ExecutorService ex)
            throws InterruptedException {
        Set<String> fils = Collections.synchronizedSet(new HashSet<>());
        long inici = System.nanoTime();
        for (int i = 0; i < TASQUES; i++) {
            ex.execute(tasca(fils));
        }
        ex.shutdown();
        ex.awaitTermination(2, TimeUnit.MINUTES);
        return new Resultat(nom,
                (System.nanoTime() - inici) / 1_000_000, fils.size());
    }

    public static void main(String[] args) throws InterruptedException {

        System.out.printf("%d tasques de %d ms d'E/S simulada%n", TASQUES, DURACIO_MS);
        System.out.printf("Sequencial seria: %d ms%n%n", TASQUES * DURACIO_MS);

        Resultat[] resultats = {
            unFilPerTasca(),
            ambExecutor("fixed(4)",  Executors.newFixedThreadPool(4)),
            ambExecutor("fixed(32)", Executors.newFixedThreadPool(32)),
            ambExecutor("cached",    Executors.newCachedThreadPool()),
            ambExecutor("a mida (16, cua 50, CallerRuns)",
                new ThreadPoolExecutor(16, 16, 0L, TimeUnit.MILLISECONDS,
                        new ArrayBlockingQueue<>(50),
                        new ThreadPoolExecutor.CallerRunsPolicy()))
        };

        System.out.printf("%-38s | %8s | %8s%n", "Estrategia", "ms", "fils");
        System.out.println("---------------------------------------|----------|---------");
        for (Resultat r : resultats) {
            System.out.printf("%-38s | %8d | %8d%n", r.estrategia(), r.ms(), r.fils());
        }
    }
}

Sortida orientativa (màquina de 8 nuclis):

300 tasques de 100 ms d'E/S simulada
Sequencial seria: 30000 ms

Estrategia                             |       ms |     fils
---------------------------------------|----------|---------
un fil per tasca                       |      147 |      300
fixed(4)                               |     7621 |        4
fixed(32)                              |     1024 |       32
cached                                 |      163 |      300
a mida (16, cua 50, CallerRuns)        |     1938 |       17

Anàlisi, fila a fila:

  • Un fil per tasca (147 ms) és el més ràpid. Amb 300 tasques purament d'espera, tenir 300 fils esperant alhora és òptim en temps. Però són 300 MB de piles reservades; amb 30.000 tasques, OutOfMemoryError. Ràpid i no escalable.
  • fixed(4) (7,6 s) és el més lent: només 4 tasques alhora, 300/4 = 75 tandes × 100 ms. Dimensionar un pool d'E/S com si fos de CPU és l'error clàssic.
  • fixed(32) (1,0 s) és un bon equilibri: 32 fils, 32 MB, i gairebé 30× d'acceleració.
  • cached (163 ms) va crear 300 fils, igual que la versió manual: la seva SynchronousQueue no emmagatzema, així que cada tasca que arriba amb tots ocupats provoca un fil nou. Ràpid aquí, i una bomba amb càrrega més gran.
  • El pool a mida (1,9 s) va fer servir 17 fils: els 16 del pool més el fil de main. Aquell fil extra és CallerRunsPolicy en acció: quan la cua de 50 es va omplir, main va executar tasques ell mateix. Això és contrapressió funcionant —el productor es va frenar sol— i és per això que la memòria no va créixer mai.

Solució a l'Exercici 2

package com.nexussoftware.bibliotech.persistencia;

import java.nio.file.Path;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;
import java.util.logging.Logger;

public class ImportacioAmbTermini {

    private static final Logger LOG = Logger.getLogger(ImportacioAmbTermini.class.getName());

    /** Informe immutable del resultat. */
    public record Informe(int importats, int descartats, long ms) { }

    /**
     * Importacio com a Callable<Informe>: retorna valor i pot
     * llancar excepcions comprovades. Es cancellable: comprova la
     * interrupcio a cada linia (08-02).
     */
    static class TascaImportacio implements Callable<Informe> {

        private final Path fitxer;
        private final AtomicInteger progres = new AtomicInteger();

        TascaImportacio(Path fitxer) { this.fitxer = fitxer; }

        int progres() { return progres.get(); }

        @Override
        public Informe call() throws Exception {
            long inici = System.nanoTime();
            int importats = 0, descartats = 0;

            for (int linia = 0; linia < 50_000; linia++) {

                // PUNT DE CANCELLACIO: cancel(true) del Future es
                // tradueix en interrupt() sobre aquest fil.
                if (Thread.currentThread().isInterrupted()) {
                    throw new InterruptedException(
                            "importacio cancellada a la linia " + linia);
                }

                if (linia % 97 == 0) descartats++; else importats++;
                progres.incrementAndGet();

                if (linia % 1000 == 0) {
                    TimeUnit.MILLISECONDS.sleep(30);   // simula validacio pesada
                }
            }
            return new Informe(importats, descartats,
                    (System.nanoTime() - inici) / 1_000_000);
        }
    }

    public static void main(String[] args) {

        ExecutorService executor = Executors.newSingleThreadExecutor(
                r -> new Thread(r, "bibliotech-importador"));

        TascaImportacio tasca = new TascaImportacio(Path.of("dades/inventari.csv"));
        Future<Informe> futur = executor.submit(tasca);

        try {
            // 1-2. El fil principal continua viu: aten el menu i pinta progres.
            System.out.println("[main] importacio llancada; el menu continua actiu");
            while (!futur.isDone()) {
                System.out.printf("\r[main] progres: %d linies", tasca.progres());
                TimeUnit.MILLISECONDS.sleep(300);
            }

            // 3. Termini de 5 segons.
            Informe informe = futur.get(5, TimeUnit.SECONDS);
            System.out.printf("%n[main] COMPLETADA: %d importats, %d descartats, %d ms%n",
                    informe.importats(), informe.descartats(), informe.ms());

        } catch (TimeoutException e) {
            // 4. El termini NO cancella per si sol: cal demanar-ho.
            System.out.printf("%n[main] termini exhaurit despres de %d linies; cancellant%n",
                    tasca.progres());
            boolean demanada = futur.cancel(true);
            System.out.println("[main] cancel() retorna   : " + demanada);
            System.out.println("[main] isCancelled()      : " + futur.isCancelled());
            System.out.println("[main] el cataleg conserva el seu contingut anterior");

        } catch (ExecutionException e) {
            // 5. La tasca ha fallat: la causa real es a getCause() (06-03).
            Throwable causa = e.getCause();
            if (causa instanceof InterruptedException) {
                System.out.println("\n[main] la importacio s'ha aturat per cancellacio");
            } else {
                LOG.log(Level.SEVERE, "La importacio ha fallat", causa);
                System.out.println("\n[main] fallada: " + causa.getMessage());
            }

        } catch (InterruptedException e) {
            // Ens han interromput a NOSALTRES; la tasca continua viva.
            System.out.println("\n[main] espera interrompuda");
            futur.cancel(true);
            Thread.currentThread().interrupt();

        } finally {
            // 6. Aturada en dues fases.
            executor.shutdown();
            try {
                if (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
                    executor.shutdownNow();
                    if (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
                        LOG.severe("L'importador ignora la interrupcio");
                    }
                }
            } catch (InterruptedException e) {
                executor.shutdownNow();
                Thread.currentThread().interrupt();
            }
            System.out.println("[main] executor aturat; la JVM pot acabar");
        }
    }
}

Sortida (quan s'exhaureix el termini):

[main] importacio llancada; el menu continua actiu
[main] progres: 41893 linies
[main] termini exhaurit despres de 42017 linies; cancellant
[main] cancel() retorna   : true
[main] isCancelled()      : true
[main] el cataleg conserva el seu contingut anterior
[main] executor aturat; la JVM pot acabar

Tres punts:

  • TimeoutException no cancel·la res per si sola. La tasca continuava corrent tan feliç; va caldre cridar cancel(true). Aquest és l'error d'expectativa més freqüent amb Future.
  • cancel(true) va funcionar perquè la tasca coopera. Sense la comprovació d'isInterrupted() al bucle, la importació hauria continuat fins al final i l'executor no s'hauria pogut apagar.
  • L'aturada en dues fases al final garanteix que la JVM pugui acabar. Sense ella, el fil bibliotech-importador continuaria viu.

Solució a l'Exercici 3

package com.nexussoftware.bibliotech.servei;

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

public class SalaDeLectura {

    private static final int AFORAMENT = 4;
    private static final int EMPLEATS = 20;

    // Semafor EQUITATIU: els que porten mes estona esperant entren abans,
    // cosa que evita la inanicio (08-04, apartat 12).
    private final Semaphore aforament = new Semaphore(AFORAMENT, true);

    // Aforament instantani, per comprovar l'invariant.
    private final AtomicInteger dins = new AtomicInteger();
    private final AtomicInteger maximObservat = new AtomicInteger();
    private volatile boolean invariantTrencat = false;

    /** Un empleat entra, hi roman una estona i surt. */
    public void visitar(String empleat) throws InterruptedException {

        aforament.acquire();      // bloqueja si ja n'hi ha 4 a dins
        try {
            int actual = dins.incrementAndGet();

            // Registre del maxim observat, amb bucle CAS (08-06).
            maximObservat.updateAndGet(m -> Math.max(m, actual));

            if (actual > AFORAMENT) {
                invariantTrencat = true;
                System.err.println("!!! INVARIANT TRENCAT: " + actual + " a dins !!!");
            }

            System.out.printf("  [+] %-12s entra  (dins=%d, esperant=%d)%n",
                    empleat, actual, aforament.getQueueLength());

            TimeUnit.MILLISECONDS.sleep(
                    ThreadLocalRandom.current().nextInt(200, 601));

            System.out.printf("  [-] %-12s surt   (dins=%d)%n",
                    empleat, dins.get() - 1);

        } finally {
            dins.decrementAndGet();
            // OBLIGATORI al finally: un permis perdut no torna mai,
            // i despres de quatre perdues la sala quedaria tancada per sempre.
            aforament.release();
        }
    }

    public static void main(String[] args) throws InterruptedException {

        SalaDeLectura sala = new SalaDeLectura();

        // Dos tancaments: la "pistola de sortida" i la "meta".
        CountDownLatch sortida = new CountDownLatch(1);
        CountDownLatch meta    = new CountDownLatch(EMPLEATS);

        // Monitor periodic de l'aforament.
        ScheduledExecutorService monitor = Executors.newSingleThreadScheduledExecutor(
                r -> {
                    Thread t = new Thread(r, "bibliotech-monitor-aforament");
                    t.setDaemon(true);
                    return t;
                });

        monitor.scheduleAtFixedRate(() -> {
            try {
                System.out.printf("      [monitor] aforament=%d/%d  cua=%d%n",
                        sala.dins.get(), AFORAMENT, sala.aforament.getQueueLength());
            } catch (Exception e) {
                // Una excepcio no capturada CANCELLARIA la planificacio
                // per sempre i en silenci.
                e.printStackTrace();
            }
        }, 200, 200, TimeUnit.MILLISECONDS);

        // 20 fils que esperen la sortida i arrenquen alhora.
        for (int i = 1; i <= EMPLEATS; i++) {
            final String nom = "empleat-" + String.format("%02d", i);
            new Thread(() -> {
                try {
                    sortida.await();           // tots esperen aqui
                    sala.visitar(nom);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } finally {
                    meta.countDown();          // SEMPRE
                }
            }, nom).start();
        }

        System.out.println("=== SALA DE LECTURA (aforament " + AFORAMENT + ") ===");
        TimeUnit.MILLISECONDS.sleep(200);      // que tots arribin a await()

        long inici = System.nanoTime();
        sortida.countDown();                   // alhora!
        meta.await();
        long ms = (System.nanoTime() - inici) / 1_000_000;

        monitor.shutdown();
        monitor.awaitTermination(2, TimeUnit.SECONDS);

        System.out.println();
        System.out.println("=== RESUM ===");
        System.out.println("Empleats atesos     : " + EMPLEATS);
        System.out.println("Aforament maxim      : " + sala.maximObservat.get()
                + " (limit " + AFORAMENT + ")");
        System.out.println("Invariant respectat  : " + !sala.invariantTrencat);
        System.out.println("Temps total          : " + ms + " ms");
        System.out.println("Dins al final        : " + sala.dins.get());
        System.out.println("Permisos disponibles : " + sala.aforament.availablePermits()
                + " (ha de ser " + AFORAMENT + ")");
    }
}

Sortida (fragment):

=== SALA DE LECTURA (aforament 4) ===
  [+] empleat-03   entra  (dins=1, esperant=16)
  [+] empleat-07   entra  (dins=2, esperant=15)
  [+] empleat-01   entra  (dins=3, esperant=14)
  [+] empleat-12   entra  (dins=4, esperant=13)
      [monitor] aforament=4/4  cua=16
  [-] empleat-03   surt   (dins=3)
  [+] empleat-05   entra  (dins=4, esperant=15)
  ...

=== RESUM ===
Empleats atesos     : 20
Aforament maxim      : 4 (limit 4)
Invariant respectat  : true
Temps total          : 2143 ms
Dins al final        : 0
Permisos disponibles : 4 (ha de ser 4)

Quatre observacions:

  1. L'aforament màxim observat és exactament 4, mai 5, en qualsevol execució. El semàfor compleix el seu contracte.
  2. Els permisos disponibles al final tornen a ser 4. És la comprovació que cap release() no s'ha perdut. Si treus el finally i provoques una excepció a dins, veuràs aquell número baixar i la sala quedar-se tancada.
  3. La porta de sortida amb CountDownLatch maximitza la contenció. Sense ella, els primers empleats haurien entrat i sortit abans que arrenquessin els últims, i no hi hauria hagut cua. És la tècnica correcta per a una prova de concurrència.
  4. Els 2,1 segons totals quadren: 20 empleats × ~400 ms de mitjana / 4 simultanis ≈ 2 s. L'aforament, no el nombre de fils, determina el temps — la mateixa lliçó del semàfor d'avisos de l'apartat 17.

Conclusió

Has canviat de nivell d'abstracció: ja no gestiones fils, descrius tasques.

Saps per què java.util.concurrent substitueix la gestió manual: un fil per tasca no escala, no hi ha límit que impedeixi crear-ne deu mil, els fils no es reutilitzen, Runnable no retorna res, les excepcions no arriben a qui crida i cancel·lar és recórrer un array. Els sis problemes desapareixen amb un executor.

Coneixes la trinitat Executor / ExecutorService / Executors i la diferència entre execute —l'excepció es veu— i submit —l'excepció es desa al Future i desapareix si ningú no crida get()—, que és un dels paranys més cars del tema. Saps què fa cada fàbrica d'Executors, i sobretot coneixes l'advertiment de les cues il·limitades: newFixedThreadPool acumula tasques fins a l'OutOfMemoryError, i newCachedThreadPool crea fils sense límit fins al mateix final per un altre camí. Per això saps construir un ThreadPoolExecutor a mida, amb els seus sis paràmetres, la seva ThreadFactory que posa nom als fils, la seva cua acotada i la seva política de rebuig —amb CallerRunsPolicy com la més útil, perquè converteix la saturació en contrapressió automàtica—. I entens l'algoritme contraintuïtiu del pool: encua abans de crear fils extra, així que maximumPoolSize no serveix de res si la cua és il·limitada.

Saps dimensionar: N per a càlcul, N × (1 + espera/càlcul) per a E/S, i —el consell que més rendeix— pools separats per tipus de feina, l'aïllament per mampares que impedeix que la saturació d'una part enfonsi les altres. I saps apagar un executor amb el patró de dues fases —shutdown, esperar, shutdownNow, esperar, registrar—, recordant que si no l'apagues la JVM no acaba, i que shutdownNow() interromp però no mata: una tasca que s'empassa la InterruptedException continua viva, i tot el protocol de 08-02 es cobra aquí.

Callable i Future resolen el que va quedar obert a 08-02: una tasca que retorna un valor i propaga excepcions comprovades. Amb les tres excepcions de get() ben distingides —ExecutionException que embolcalla la causa real a getCause(), TimeoutException que no cancel·la res i deixa la tasca corrent, i InterruptedException que t'afecta a tu i no a la tasca—, la regla de no fer servir mai get() sense termini en producció, i cancel(true) que no és màgia sinó un interrupt() sobre el fil del pool. Més invokeAll per a lots amb termini global i invokeAny per quedar-te amb la font més ràpida.

Coneixes ScheduledExecutorService i la diferència entre scheduleAtFixedRate —període des de l'inici, ritme constant— i scheduleWithFixedDelay —període des del final, separació real garantida—, amb l'advertiment que arruïna sistemes: una excepció no capturada en una tasca periòdica cancel·la la planificació per sempre, en silenci.

I tens els sincronitzadors: CountDownLatch d'un sol ús, per esperar N esdeveniments i per donar la sortida a N fils alhora, amb el countDown() sempre en un finally; CyclicBarrier reutilitzable, amb la seva acció de barrera i la seva ruptura per a tothom si un participant falla; Semaphore per posar un sostre dur a la concurrència sobre un recurs, amb el release() sempre en un finally perquè un permís perdut no torna mai; i Exchanger i Phaser reconeguts encara que rarament necessaris. Més ForkJoinPool amb RecursiveTask, el seu robatori de feina, el patró correcte fork() + compute() + join(), la importància crítica del llindar, i l'avís sobre el pool comú: és compartit per tota la JVM i bloquejar-lo perjudica tothom.

BiblioTech envia els seus 200 avisos en deu segons amb vuit fils, amb barra de progrés mitjançant CountDownLatch, amb un semàfor que limita a cinc els accessos simultanis al fitxer, amb termini global i cancel·lació neta, amb cada fallada individual registrada sense aturar el lot, i amb AutoCloseable perquè el try-with-resources apagui el pool. I l'anàlisi del resultat ensenya la lliçó més valuosa de l'apartat: l'acceleració va ser 4,9× i no 8× perquè el coll d'ampolla no era el pool, sinó el semàfor de cinc permisos. Dimensionar bé és trobar quin és el recurs escàs — i gairebé mai no és el que et penses.

Però queda una escletxa. ServeiAvisos fa servir AtomicInteger per als comptadors sense haver-los explicat, i el Cataleg que vas fer segur a 08-04 continua basant-se en un HashMap protegit amb panys: cada lectura paga un lock(), encara que llegir no destorbi ningú. I la CuaReserves del mòdul 5 continua sent un ArrayDeque que dos fils corromprien, amb el patró productor-consumidor promès des de llavors i encara sense implementar.

A la propera lliçó, Col·leccions Concurrents i Variables Atòmiques, es tanca aquesta escletxa. Veuràs una demostració real d'un HashMap corromput per dos fils —amb el bucle infinit que pot deixar un nucli al 100 % per sempre— i la ConcurrentModificationException del fail-fast de 05-02; les tres generacions de solució, des de Collections.synchronizedMap amb el seu parany —iterar continua requerint sincronització manual i les operacions compostes continuen sense ser atòmiques— fins a ConcurrentHashMap amb les seves operacions atòmiques compostes putIfAbsent, computeIfAbsent, compute i merge, que són la forma correcta de resoldre comprovar-després-actuar; CopyOnWriteArrayList per als oients d'esdeveniments; les BlockingQueue en totes les seves variants i, per fi, la implementació completa del patró productor-consumidor promès al mòdul 5, aplicat a la cua de reserves i amb píndola verinosa per a l'apagada; i les variables atòmiques amb l'operació compare-and-swap explicada per dins, el seu bucle de reintent, el problema ABA i LongAdder per a comptadors sota molta contenció. En acabar tindràs la taula de decisió definitiva: quan synchronized, quan un Lock, quan un atòmic, quan una col·lecció concurrent i quan simplement un objecte immutable.

Curs de Programació en Java

Mòdul 1: Introducció a Java

Mòdul 2: Flux de control

Mòdul 3: Programació orientada a objectes

Mòdul 4: Programació orientada a objectes avançada

Mòdul 5: Estructures de dades i col·leccions

Mòdul 6: Gestió d'excepcions

Mòdul 7: Entrada/sortida de fitxers

Mòdul 8: Multifil i concurrència

Mòdul 9: Xarxes

Mòdul 10: Temes avançats

Mòdul 11: Frameworks i llibreries de Java

Mòdul 12: Construcció d'aplicacions del món real

© Copyright 2026. Tots els drets reservats