Al mòdul 5 vas construir el cor de BiblioTech amb ArrayList, HashMap, HashSet i ArrayDeque. Eren les eines correctes —per a un sol fil—. A 08-04 vas protegir el catàleg amb panys i va funcionar, però cada lectura paga un lock() encara que llegir no destorbi ningú, i a 08-05 ja vas fer servir AtomicInteger sense haver-lo explicat.
Aquesta lliçó tanca les dues escletxes. Comença on més fa mal: una demostració que un HashMap compartit entre dos fils no és que doni resultats estranys, és que es pot corrompre estructuralment i deixar un nucli al 100 % en un bucle infinit del qual no surt mai. A partir d'aquí recorre les tres generacions de solució —col·leccions sincronitzades, col·leccions concurrents, i variables atòmiques— amb la pregunta que les uneix: com s'aconsegueix que una operació composta sigui atòmica sense bloquejar tothom?
Pel camí se salda el deute del mòdul 5: la BlockingQueue que s'hi va mencionar i es va remetre aquí, i la implementació completa del patró productor-consumidor que es va descriure conceptualment i es va deixar sense escriure.
En acabar, el catàleg de BiblioTech farà servir ConcurrentHashMap, les seves estadístiques seran comptadors atòmics sense un sol pany, i la seva cua de reserves serà un productor-consumidor real amb apagada neta per píndola verinosa.
Contingut
- Un
HashMapcorromput per dos fils - La
ConcurrentModificationExceptionrevisitada - Primera generació: col·leccions sincronitzades
- El doble parany de les col·leccions sincronitzades
ConcurrentHashMap: com funciona per dins- Les operacions atòmiques compostes
- Iteradors dèbilment consistents i el
size()aproximat CopyOnWriteArrayListiCopyOnWriteArraySetConcurrentLinkedQueuei les cues concurrentsBlockingQueue: la família completa- El patró productor-consumidor
- La píndola verinosa
ConcurrentSkipListMapen una nota- Variables atòmiques i compare-and-swap
- Les operacions de les classes atòmiques
AtomicReferencei el problema ABALongAddersota contenció- BiblioTech: catàleg concurrent i estadístiques atòmiques
- Taula final de decisió
- Errors Comuns i Consells
- Exercicis
- Un
HashMap corromput per dos fils
HashMap corromput per dos filsHashMap no és segur per a diversos fils. La frase es repeteix a tot arreu; el que gairebé mai no s'explica és què significa exactament, i el significat és pitjor del que sona.
Recorda de 05-05 l'estructura interna: un array de cubetes, i a cada cubeta una llista enllaçada (o un arbre, si creix molt) amb les entrades que col·lisionen. Quan el mapa supera el factor de càrrega, es fa un rehash: es crea un array més gran i es redistribueixen totes les entrades.
Si dos fils fan put durant un rehash, les llistes enllaçades poden acabar formant un cicle. I una llista amb un cicle fa que un get() posterior no acabi mai.
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeUnit;
public class HashMapCorromput {
public static void main(String[] args) throws InterruptedException {
for (int intent = 1; intent <= 5; intent++) {
Map<Integer, String> mapa = new HashMap<>();
final int PER_FIL = 100_000;
Thread f1 = new Thread(() -> {
for (int i = 0; i < PER_FIL; i++) mapa.put(i, "A" + i);
}, "escriptor-1");
Thread f2 = new Thread(() -> {
for (int i = PER_FIL; i < PER_FIL * 2; i++) mapa.put(i, "B" + i);
}, "escriptor-2");
f1.start();
f2.start();
// Termini: si el HashMap es corromp, un fil es pot quedar
// en un bucle infinit dins de put() o de get().
f1.join(5000);
f2.join(5000);
int esperat = PER_FIL * 2;
if (f1.isAlive() || f2.isAlive()) {
System.out.printf("Intent %d: BUCLE INFINIT. "
+ "f1=%s f2=%s <-- estructura corrompuda%n",
intent, f1.getState(), f2.getState());
System.out.println(" (mira l'us de CPU: un nucli al 100%)");
System.exit(1);
}
System.out.printf("Intent %d: esperat %d, real %d, perdues %d%n",
intent, esperat, mapa.size(), esperat - mapa.size());
}
}
}Sortida típica:
Intent 1: esperat 200000, real 187341, perdues 12659
Intent 2: esperat 200000, real 193028, perdues 6972
Intent 3: BUCLE INFINIT. f1=RUNNABLE f2=RUNNABLE <-- estructura corrompuda
(mira l'us de CPU: un nucli al 100%)Els dos modes de fallada, de menys a més greu:
- Entrades perdudes. Dos
putsimultanis sobre la mateixa cubeta trepitgen l'un la feina de l'altre. El mapa acaba amb menys entrades de les inserides. És la condició de cursa de 08-04 aplicada a una estructura de dades. - Bucle infinit. Durant el rehash, dos fils poden deixar la llista enllaçada d'una cubeta apuntant-se a si mateixa. Un
get()que caigui en aquella cubeta recorre el cicle per sempre, consumint un nucli al 100 % i sense llançar cap excepció.
El segon cas és una anècdota famosa de la indústria: durant anys va ser una causa recurrent de servidors que es quedaven al 100 % de CPU sense motiu aparent, i el diagnòstic —bolcat de fils, apartat 12 de 08-03, amb diversos fils RUNNABLE dins de HashMap.get— és una història que explica qualsevol que porti temps en producció.
Nota tècnica. En Java 8+ la implementació del rehash va canviar i el cicle és molt més difícil de provocar que en Java 7, però la classe continua sense ser segura per a diversos fils i les pèrdues d'entrades es reprodueixen sense dificultat. No és un problema resolt: és un problema menys visible.
La conclusió: compartir un HashMap entre fils sense protecció no dona "resultats aproximats". Dona estructures trencades.
- La
ConcurrentModificationException revisitada
ConcurrentModificationException revisitadaA 05-02 vas veure el comportament fail-fast: els iteradors de les col·leccions clàssiques porten un comptador modCount i comproven a cada next() que ningú no hagi modificat la col·lecció per darrere.
List<String> llista = new ArrayList<>(List.of("a", "b", "c"));
for (String s : llista) {
llista.remove(s); // ConcurrentModificationException
}Allà era un sol fil. Ara apareix la versió concurrent, i és més traïdora:
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;
public class FailFastConcurrent {
public static void main(String[] args) throws InterruptedException {
List<String> cataleg = new ArrayList<>();
for (int i = 0; i < 10_000; i++) cataleg.add("978-" + i);
Thread lector = new Thread(() -> {
try {
while (!Thread.currentThread().isInterrupted()) {
int n = 0;
for (String isbn : cataleg) { // iteracio
n += isbn.length();
}
}
} catch (java.util.ConcurrentModificationException e) {
System.out.println("[lector] ConcurrentModificationException: "
+ "un altre fil ha modificat el cataleg mentre iterava");
}
}, "bibliotech-lector");
Thread escriptor = new Thread(() -> {
try {
while (!Thread.currentThread().isInterrupted()) {
cataleg.add("978-nou");
TimeUnit.MILLISECONDS.sleep(1);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "bibliotech-escriptor");
lector.start();
escriptor.start();
TimeUnit.SECONDS.sleep(2);
lector.interrupt();
escriptor.interrupt();
}
}Sortida (gairebé immediata):
L'important d'aquest exemple és com cal interpretar aquesta excepció. No és "un error de concurrència que cal capturar". És un detector d'errors que la col·lecció ofereix com a cortesia: t'avisa que estàs fent servir la col·lecció de forma insegura. Capturar l'excepció i reintentar és tapar el símptoma; la solució és canviar la col·lecció o protegir l'accés.
I hi ha un matís que s'oblida: el fail-fast no està garantit. La documentació és explícita: modCount no és volatile, així que un iterador pot no veure la modificació i retornar dades incoherents en lloc de llançar l'excepció. Pots tenir una cursa silenciosa.
- Primera generació: col·leccions sincronitzades
Java 1.2 va introduir els embolcalls sincronitzats de Collections:
import java.util.*;
Map<String, Material> mapa = Collections.synchronizedMap(new HashMap<>());
List<Material> llista = Collections.synchronizedList(new ArrayList<>());
Set<String> conjunt = Collections.synchronizedSet(new HashSet<>());Cada mètode de l'embolcall és un synchronized sobre el mateix embolcall:
// Aproximadament, el que fa Collections.synchronizedMap:
public V get(Object clau) {
synchronized (mutex) { return m.get(clau); }
}
public V put(K clau, V valor) {
synchronized (mutex) { return m.put(clau, valor); }
}Resol la corrupció de l'apartat 1 —els put ja no es trepitgen— però té un problema evident i dos paranys.
El problema evident: un únic pany global. Totes les operacions, incloses les lectures, es serialitzen. Amb setze fils llegint un mapa, quinze estan esperant. És un coll d'ampolla exacte.
HashtableiVectorsón la versió antiga del mateix, de Java 1.0, amb tots els seus mètodes sincronitzats. Tenen a més el defecte de 08-04: sincronitzen sobrethis, així que el seu pany està exposat. No els facis servir en codi nou.
- El doble parany de les col·leccions sincronitzades
Parany 1: iterar continua requerint sincronització manual.
Cada mètode individual està sincronitzat, però una iteració són moltes crides. Entre hasNext() i next() no hi ha cap pany, així que un altre fil pot modificar la col·lecció i provocar la ConcurrentModificationException de l'apartat 2.
List<Material> llista = Collections.synchronizedList(new ArrayList<>());
// MALAMENT: cada crida de l'iterador esta sincronitzada, pero la
// ITERACIO COMPLETA no. ConcurrentModificationException garantida.
for (Material m : llista) {
processar(m);
}
// BE: sincronitzar la iteracio sencera sobre el mateix embolcall.
// Es el que exigeix explicitament el javadoc de Collections.synchronizedList,
// i gairebe ningu no llegeix aquesta part.
synchronized (llista) {
for (Material m : llista) {
processar(m); // COMPTE: 'processar' s'executa AMB EL PANY PRES
} // (regla 3 de 08-04: res de codi alie aqui)
}Fixa't en el cost de la versió correcta: durant tota la iteració, ningú més no pot tocar la llista. Amb deu mil elements i un processament d'un mil·lisegon cadascun, són deu segons de bloqueig total.
Parany 2: les operacions compostes continuen sense ser atòmiques.
Aquest és el pitjor, perquè el codi sembla segur.
Map<String, Integer> prestecsPerIsbn = Collections.synchronizedMap(new HashMap<>());
// MALAMENT: comprovar-despres-actuar. Cada crida esta sincronitzada,
// pero el FORAT entre elles no ho esta.
if (!prestecsPerIsbn.containsKey(isbn)) { // (1) fil A: no existeix
prestecsPerIsbn.put(isbn, 1); // (3) fil A escriu 1
} // (2) fil B tambe va veure "no existeix"
// (4) fil B escriu 1 -> se'n perd un
// MALAMENT: llegir-modificar-escriure, el mateix problema que 'comptador++' (08-04)
Integer n = prestecsPerIsbn.get(isbn); // (1) llegeix 5
prestecsPerIsbn.put(isbn, n + 1); // (3) escriu 6
// un altre fil tambe va llegir 5 i
// tambe escriu 6: un prestec perdutSón les dues formes canòniques de condició de cursa de 08-04, i la col·lecció sincronitzada no fa res per evitar-les. L'única forma d'arreglar-ho amb aquesta generació és un pany extern:
// Correcte pero maldestre: pany extern sobre la colleccio.
synchronized (prestecsPerIsbn) {
Integer n = prestecsPerIsbn.get(isbn);
prestecsPerIsbn.put(isbn, n == null ? 1 : n + 1);
}Funciona, però has hagut de sortir de l'abstracció: la col·lecció "segura" no ho era per al que necessitaves, i ara la seguretat depèn que tot el codi de l'aplicació recordi prendre aquell pany. Amb això arribem a la segona generació.
ConcurrentHashMap: com funciona per dins
ConcurrentHashMap: com funciona per dinsConcurrentHashMap (Java 5, reescrit a Java 8) resol les tres coses alhora: no es corromp, no serialitza les lectures, i ofereix operacions compostes atòmiques.
Les seves dues idees centrals:
Idea 1: les lectures no bloquegen mai. Els nodes interns tenen els camps value i next declarats volatile. Gràcies a les garanties de visibilitat de 08-04, un get() pot llegir sense adquirir cap pany i tot i així veure un valor coherent i recent. Zero contenció entre lectors, i entre lectors i escriptors.
Idea 2: les escriptures bloquegen només la cubeta afectada. En lloc d'un pany global, es sincronitza sobre el primer node de la cubeta. Dues escriptures en cubetes diferents —que és el cas normal, perquè el hash les reparteix— no es destorben en absolut.
flowchart TB
subgraph SM["Collections.synchronizedMap"]
direction TB
C1["UN pany global"] --> T1["cubeta 0"]
C1 --> T2["cubeta 1"]
C1 --> T3["cubeta 2"]
C1 --> T4["cubeta 3"]
N1["Totes les operacions,<br/>lectures incloses,<br/>es serialitzen"]
end
subgraph CHM["ConcurrentHashMap"]
direction TB
L["Lectures: SENSE pany<br/>camps volatile"]
B0["cubeta 0<br/>pany propi"]
B1["cubeta 1<br/>pany propi"]
B2["cubeta 2<br/>pany propi"]
B3["cubeta 3<br/>pany propi"]
N2["Escriptures en cubetes<br/>diferents: en paral·lel"]
end
Demostració de la diferència de rendiment:
import java.util.*;
import java.util.concurrent.*;
public class ComparativaMapes {
static long mesurar(Map<Integer, String> mapa, int fils, int opsPerFil,
double proporcioLectura) throws InterruptedException {
// Precarregar perque les lectures encertin.
for (int i = 0; i < 10_000; i++) mapa.put(i, "valor-" + i);
CountDownLatch sortida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(fils);
for (int f = 0; f < fils; f++) {
new Thread(() -> {
try {
sortida.await();
ThreadLocalRandom atzar = ThreadLocalRandom.current();
for (int i = 0; i < opsPerFil; i++) {
int clau = atzar.nextInt(10_000);
if (atzar.nextDouble() < proporcioLectura) {
mapa.get(clau);
} else {
mapa.put(clau, "nou-" + i);
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
meta.countDown();
}
}, "acces-" + f).start();
}
long inici = System.nanoTime();
sortida.countDown(); // arrenquen tots alhora (08-05)
meta.await();
return (System.nanoTime() - inici) / 1_000_000;
}
public static void main(String[] args) throws InterruptedException {
final int FILS = 16;
final int OPS = 200_000;
System.out.printf("%d fils x %d operacions%n%n", FILS, OPS);
System.out.printf("%-28s | %10s | %10s%n", "Implementacio", "90% lectura", "50% lectura");
System.out.println("-----------------------------|------------|------------");
long s90 = mesurar(Collections.synchronizedMap(new HashMap<>()), FILS, OPS, 0.9);
long s50 = mesurar(Collections.synchronizedMap(new HashMap<>()), FILS, OPS, 0.5);
System.out.printf("%-28s | %8d ms | %8d ms%n", "synchronizedMap", s90, s50);
long c90 = mesurar(new ConcurrentHashMap<>(), FILS, OPS, 0.9);
long c50 = mesurar(new ConcurrentHashMap<>(), FILS, OPS, 0.5);
System.out.printf("%-28s | %8d ms | %8d ms%n", "ConcurrentHashMap", c90, c50);
System.out.printf("%nMillora: %.1fx (90%% lectura), %.1fx (50%% lectura)%n",
(double) s90 / c90, (double) s50 / c50);
}
}Sortida orientativa:
16 fils x 200000 operacions
Implementacio | 90% lectura | 50% lectura
-----------------------------|------------|------------
synchronizedMap | 3187 ms | 3402 ms
ConcurrentHashMap | 184 ms | 271 ms
Millora: 17.3x (90% lectura), 12.6x (50% lectura)Un ordre de magnitud de diferència, i creix amb el nombre de fils. Amb un sol fil, en canvi, els dos són gairebé iguals: el guany és d'escalabilitat, no de velocitat bruta.
- Les operacions atòmiques compostes
Aquesta és l'aportació més valuosa de ConcurrentHashMap, i la que resol el parany 2 de l'apartat 4: operacions compostes que són atòmiques de veritat.
| Mètode | Què fa atòmicament |
|---|---|
putIfAbsent(k, v) |
Insereix si la clau no hi és; retorna el valor previ o null |
computeIfAbsent(k, f) |
Si no hi és, calcula el valor amb f i insereix |
computeIfPresent(k, f) |
Si hi és, recalcula el valor amb f |
compute(k, f) |
Recalcula sempre; null com a resultat elimina l'entrada |
merge(k, v, f) |
Si no hi és posa v; si hi és, combina l'actual amb v amb f |
remove(k, v) |
Elimina només si el valor actual és v |
replace(k, vell, nou) |
Substitueix només si el valor actual és vell |
getOrDefault(k, d) |
Retorna el valor o d si no hi és (no modifica) |
Els mateixos casos de l'apartat 4, ara correctes:
import java.util.concurrent.ConcurrentHashMap;
import java.util.Map;
public class OperacionsAtomiquesMapa {
private final Map<String, Integer> prestecsPerIsbn = new ConcurrentHashMap<>();
private final Map<String, Fitxa> cacheFitxes = new ConcurrentHashMap<>();
/** Comprovar-despres-actuar, resolt: una sola operacio atomica. */
public boolean registrarPrimerPrestec(String isbn) {
// Retorna null si NO existia (i l'insereix), o el valor previ.
return prestecsPerIsbn.putIfAbsent(isbn, 1) == null;
}
/** Llegir-modificar-escriure, resolt amb merge. */
public void comptarPrestec(String isbn) {
// Si no existeix -> posa 1. Si existeix -> aplica Integer::sum
// entre el valor actual i l'1 que passem. Tot atomic.
prestecsPerIsbn.merge(isbn, 1, Integer::sum);
}
/** El mateix amb compute, mes explicit. */
public void comptarPrestecAlternatiu(String isbn) {
prestecsPerIsbn.compute(isbn, (clau, actual) ->
actual == null ? 1 : actual + 1);
}
/**
* Cache mandrosa, resolta amb computeIfAbsent.
* Substitueix les tres fases de l'exercici 3 de 08-04:
* la funcio s'executa COM A MOLT UNA VEGADA per clau, encara que
* deu fils la demanin alhora. Els altres nou esperen i
* reben el mateix objecte: no hi ha calcul duplicat.
*/
public Fitxa obtenirFitxa(String isbn) {
return cacheFitxes.computeIfAbsent(isbn, this::construirFitxa);
}
/** Eliminar nomes si el valor es l'esperat: comparar-i-eliminar. */
public boolean retornarSiEsUltim(String isbn) {
return prestecsPerIsbn.remove(isbn, 1);
}
/** Comptador de descarregues sense NullPointerException. */
public int descarreguesDe(String isbn) {
return prestecsPerIsbn.getOrDefault(isbn, 0);
}
private Fitxa construirFitxa(String isbn) { /* consulta cara */ return null; }
}L'advertiment crític sobre computeIfAbsent i compute: la funció que hi passes s'executa amb el pany de la cubeta pres. D'aquí tres prohibicions absolutes:
// PROHIBIT 1: modificar el MATEIX mapa dins de la funcio.
// Pot provocar un interbloqueig o corrompre l'estructura.
mapa.computeIfAbsent(k, clau -> {
mapa.put("altra", "cosa"); // MAI
return calcular(clau);
});
// PROHIBIT 2: operacions llargues o blocants.
// Bloqueja la cubeta sencera, i tota l'aplicacio se'n ressent
// (regla 2 de 08-04: res d'E/S dins del bloqueig).
mapa.computeIfAbsent(k, clau -> llegirDelDisc(clau)); // MALAMENT si triga
// PROHIBIT 3: cridar codi alie (oients, callbacks).
// Regla 3 de 08-04, aplicada aqui.
// CORRECTE quan el calcul es car: calcular fora i publicar amb putIfAbsent.
Fitxa calculada = llegirDelDisc(isbn); // sense cap pany
Fitxa establerta = mapa.putIfAbsent(isbn, calculada);
Fitxa resultat = (establerta != null) ? establerta : calculada;
// Hi pot haver calcul duplicat si dos fils coincideixen, pero
// el resultat es correcte i no es bloqueja la cubeta.
- Iteradors dèbilment consistents i el
size() aproximat
size() aproximatLes col·leccions concurrents canvien dos contractes respecte a les clàssiques, i cal conèixer-los.
Iteradors dèbilment consistents. Els iteradors de ConcurrentHashMap no llancen ConcurrentModificationException. Recorren l'estat del mapa en el moment en què es va crear l'iterador, i poden o no reflectir modificacions posteriors. No garanteixen veure els canvis, però garanteixen no trencar-se.
Fail-fast (HashMap) |
Dèbilment consistent (ConcurrentHashMap) |
|
|---|---|---|
| Modificació durant la iteració | ConcurrentModificationException |
Sense excepció |
| Veu els canvis posteriors | — | Potser sí, potser no |
| Recorre cada element | Sí, o falla | Sí, cada element present en començar |
| Bloqueja els escriptors | Només si sincronitzes a mà | Mai |
| És segur per a diversos fils | No | Sí |
Map<String, Material> cataleg = new ConcurrentHashMap<>();
// SEGUR: no llanca excepcio, no bloqueja ningu, i cap escriptor
// no es queda esperant que acabis de recorrer 100.000 entrades.
for (Map.Entry<String, Material> e : cataleg.entrySet()) {
processar(e.getValue());
}
// Si un altre fil afegeix una entrada mentre iteres, potser la veuras
// i potser no. El que NO passara es que el bucle falli.size() és aproximat. A ConcurrentHashMap, size(), isEmpty() i containsValue() retornen un valor que era correcte en algun instant recent, però que pot haver canviat abans que el facis servir.
// MALAMENT: comprovar-despres-actuar sobre una mida aproximada.
if (cataleg.size() < LIMIT) {
cataleg.put(isbn, material); // la mida ha pogut canviar entremig
}
// La mida d'una estructura concurrent es una dada ESTADISTICA,
// util per a logs i metriques, no per a decisions de control.
LOG.log(Level.INFO, "Cataleg amb ~{0} materials", cataleg.size());És una conseqüència inevitable del disseny: mantenir un comptador exacte exigiria un punt de sincronització global, i això és justament el que s'ha eliminat per guanyar escalabilitat.
CopyOnWriteArrayList i CopyOnWriteArraySet
CopyOnWriteArrayList i CopyOnWriteArraySetEstratègia radicalment diferent: cada modificació copia l'array sencer.
- Lectures: sense cap pany, sobre un array immutable. Rapidíssimes.
- Escriptures: sota pany, copien tot l'array. Cost O(n) per escriptura.
package com.nexussoftware.bibliotech.servei;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
/**
* Registre d'oients d'esdeveniments del cataleg.
*
* Cas d'us PERFECTE per a CopyOnWriteArrayList:
* - Es registren ~10 oients en arrencar i gairebe mai no canvien.
* - Es recorren a CADA operacio del cataleg: milers de vegades per minut.
* - Proporcio lectura/escriptura: 100.000 a 1.
*/
public class RegistreOients {
private final List<OientCataleg> oients = new CopyOnWriteArrayList<>();
public void registrar(OientCataleg o) { oients.add(o); }
public void desregistrar(OientCataleg o) { oients.remove(o); }
/**
* Notifica tots els oients.
*
* DOS avantatges decisius davant d'una llista sincronitzada:
* 1. NO hi ha pany durant el recorregut: els oients poden
* registrar o desregistrar oients sense interbloquejar-se.
* 2. Es la regla 3 de 08-04 —no cridar codi alie amb un
* pany— complerta automaticament per la colleccio.
*/
public void notificarAlta(Material m) {
for (OientCataleg o : oients) {
try {
o.materialAfegit(m);
} catch (RuntimeException e) {
// Un oient defectuos no ha d'impedir que els altres
// se n'assabentin (politica de degradacio de 06-07).
LOG.log(Level.WARNING, "Oient fallit: " + o, e);
}
}
}
}El detall elegant: l'iterador de CopyOnWriteArrayList treballa sobre una instantània de l'array en el moment de crear-se. És completament immune a les modificacions, no llança excepcions i no bloqueja ningú. A canvi, no admet remove() (llança UnsupportedOperationException) perquè modificar una instantània no tindria sentit.
Quan compensa i quan no:
| Situació | Fer servir còpia-en-escriure? |
|---|---|
| Oients d'esdeveniments | Sí, el cas canònic |
| Llista de configuració llegida constantment | Sí |
| Conjunt petit de regles de negoci | Sí |
| Llista amb milers d'elements i escriptures freqüents | No: cada escriptura copia milers de referències |
| Cua de treball | No: fes servir una BlockingQueue |
| Acumulador que creix en un bucle | No: O(n²) total |
// DESASTRE: 10.000 escriptures sobre una llista que creix.
// Cada add() copia tot l'array: 1 + 2 + ... + 10.000 ≈ 50 milions
// de copies de referencies. Segons on haurien de ser milisegons.
List<Material> llista = new CopyOnWriteArrayList<>();
for (int i = 0; i < 10_000; i++) {
llista.add(materials.get(i)); // O(n) cadascuna -> O(n²) total
}
ConcurrentLinkedQueue i les cues concurrents
ConcurrentLinkedQueue i les cues concurrentsConcurrentLinkedQueue és una cua FIFO no blocant i il·limitada, implementada amb l'algoritme de Michael-Scott basat en compare-and-swap (apartat 14): sense panys en absolut.
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
Queue<Reserva> pendents = new ConcurrentLinkedQueue<>();
pendents.offer(reserva); // afegir: no bloqueja mai, no falla mai
Reserva r = pendents.poll(); // treure: retorna NULL si esta buidaLa diferència clau amb una BlockingQueue: poll() retorna null si la cua està buida, en lloc d'esperar. Això obliga el consumidor a sondejar:
// ANTIPATRO amb ConcurrentLinkedQueue: sondeig (08-03, apartat 6).
while (!Thread.currentThread().isInterrupted()) {
Reserva r = pendents.poll();
if (r == null) {
Thread.sleep(100); // es desperta 10 vegades per segon per no res
continue;
}
atendre(r);
}Aquest bucle crema CPU quan no hi ha feina i afegeix fins a 100 ms de latència quan n'hi ha. Gairebé sempre el que vols és una BlockingQueue.
Fes servir ConcurrentLinkedQueue quan no necessitis esperar: acumular esdeveniments que un altre fil buidarà periòdicament, o una bossa de feina consultada de forma oportunista.
BlockingQueue: la família completa
BlockingQueue: la família completaAquí se salda el deute del mòdul 5. BlockingQueue és una cua les operacions de la qual esperen quan no es poden completar: take() espera si està buida, put() espera si està plena.
Els quatre grups d'operacions, que cal conèixer perquè l'elecció importa:
| Llança excepció | Retorna valor especial | Bloqueja | Amb termini | |
|---|---|---|---|---|
| Inserir | add(e) |
offer(e) → false |
put(e) |
offer(e, t, u) |
| Extreure | remove() |
poll() → null |
take() |
poll(t, u) |
| Examinar | element() |
peek() → null |
— | — |
Les implementacions:
| Implementació | Capacitat | Estructura | Quan fer-la servir |
|---|---|---|---|
ArrayBlockingQueue |
Acotada (fixa) | Array circular | L'opció per defecte: la cota dona contrapressió |
LinkedBlockingQueue |
Opcionalment acotada | Llista enllaçada | Més rendiment amb molts productors i consumidors |
SynchronousQueue |
0 | Sense emmagatzematge | Traspàs directe: cada put espera un take |
PriorityBlockingQueue |
Il·limitada | Monticle | Quan l'ordre el marca la prioritat, no l'arribada |
DelayQueue |
Il·limitada | Monticle per temps | Elements que no es poden treure fins a cert instant |
LinkedTransferQueue |
Il·limitada | Llista enllaçada | transfer(): esperar que el consumidor ho rebi |
Cadascuna en context:
import java.util.concurrent.*;
// 1. ARRAYBLOCKINGQUEUE: acotada. Si s'omple, el productor ESPERA.
// Aixo es CONTRAPRESSIO: el productor es frena al ritme del consumidor.
// Es el que evita l'OutOfMemoryError de les cues illimitades (08-05).
BlockingQueue<Reserva> reserves = new ArrayBlockingQueue<>(100);
// 2. LINKEDBLOCKINGQUEUE acotada: dos panys interns (un per al
// cap i un altre per a la cua), aixi que un productor i un consumidor
// poden treballar alhora. Millor amb molta concurrencia.
BlockingQueue<Avis> avisos = new LinkedBlockingQueue<>(500);
// 3. SYNCHRONOUSQUEUE: capacitat ZERO. Cada put() espera un take().
// Es el traspas ma a ma, sense magatzem. Es la cua que fa servir
// newCachedThreadPool (08-05) i per aixo crea fils sense limit.
BlockingQueue<Tasca> traspas = new SynchronousQueue<>();
// 4. PRIORITYBLOCKINGQUEUE: surt primer el mes "petit" segons el
// Comparator (05-09). Els avisos mes vencuts, primer.
BlockingQueue<Prestec> perUrgencia = new PriorityBlockingQueue<>(
100, Comparator.comparingLong(Prestec::diesDeRetard).reversed());
// 5. DELAYQUEUE: els elements implementen Delayed i no es poden
// treure fins que expiri el seu retard. Reintents amb espera.
DelayQueue<ReintentAvis> reintents = new DelayQueue<>();L'ArrayBlockingQueue acotada mereix un paràgraf, perquè resol un problema real de 08-05. Amb una cua il·limitada, un productor ràpid acumula tasques fins a exhaurir la memòria. Amb una d'acotada, quan la cua s'omple el productor es bloqueja a put() i deixa de produir fins que hi hagi forat. El sistema s'autoregula sense descartar res i sense créixer sense límit. És el mateix efecte que CallerRunsPolicy aconseguia en un pool.
- El patró productor-consumidor
Promès al mòdul 5, descrit conceptualment i remès aquí. Ara, complet.
La idea: uns fils (productors) generen feina i la posen en una cua; altres (consumidors) la treuen i la processen. La cua desacobla tots dos: no necessiten conèixer-se, ni anar al mateix ritme, ni coordinar-se.
sequenceDiagram
participant P1 as productor-1
participant P2 as productor-2
participant Q as BlockingQueue(10)
participant C1 as consumidor-1
participant C2 as consumidor-2
P1->>Q: put(reserva A)
P2->>Q: put(reserva B)
C1->>Q: take() -> A
Note over C1: processa A
C2->>Q: take() -> B
Note over C2: processa B
C1->>Q: take()
Note over C1,Q: cua buida: C1 BLOQUEJAT<br/>sense consumir CPU
P1->>Q: put(reserva C)
Q-->>C1: es desperta amb C
Note over P1,Q: si la cua s'omple,<br/>els productors es bloquegen<br/>a put(): CONTRAPRESSIÓ
Implementació completa per a BiblioTech:
package com.nexussoftware.bibliotech.servei;
import com.nexussoftware.bibliotech.domini.Reserva;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.logging.Level;
import java.util.logging.Logger;
/**
* Processador de reserves amb el patro productor-consumidor.
*
* Els empleats produeixen reserves des del menu (diversos fils);
* un grup de consumidors les aten (comprovar disponibilitat,
* notificar, registrar).
*
* La cua ACOTADA dona contrapressio: si les reserves arriben mes rapid
* del que s'atenen, el menu es frena en lloc d'acumular
* reserves en memoria fins a exhaurir-la.
*/
public class ProcessadorReserves implements AutoCloseable {
private static final Logger LOG = Logger.getLogger(ProcessadorReserves.class.getName());
private static final int CAPACITAT = 50;
private static final int CONSUMIDORS = 4;
private final BlockingQueue<Reserva> cua = new ArrayBlockingQueue<>(CAPACITAT);
private final ExecutorService consumidors;
private final AtomicInteger ateses = new AtomicInteger();
private final AtomicInteger rebutjades = new AtomicInteger();
private volatile boolean acceptantNoves = true;
public ProcessadorReserves() {
AtomicInteger n = new AtomicInteger(1);
this.consumidors = Executors.newFixedThreadPool(CONSUMIDORS,
r -> new Thread(r, "bibliotech-reserves-" + n.getAndIncrement()));
for (int i = 0; i < CONSUMIDORS; i++) {
consumidors.execute(this::bucleConsumidor);
}
}
// ---------- PRODUCTOR ----------
/**
* Encua una reserva. BLOQUEJA si la cua esta plena: es la
* contrapressio, i es una caracteristica, no un defecte.
*/
public void encuar(Reserva r) throws InterruptedException {
if (!acceptantNoves) {
throw new IllegalStateException("El processador s'esta apagant");
}
cua.put(r); // espera si esta plena
}
/**
* Variant que no espera indefinidament: si en 2 segons no hi ha
* forat, rebutja. Es l'apropiat per a una interficie d'usuari:
* millor dir "el sistema esta saturat" que deixar el menu penjat.
*/
public boolean encuarAmbTermini(Reserva r) throws InterruptedException {
boolean acceptada = cua.offer(r, 2, TimeUnit.SECONDS);
if (!acceptada) {
rebutjades.incrementAndGet();
LOG.log(Level.WARNING, "Reserva rebutjada per saturacio: {0}", r.id());
}
return acceptada;
}
// ---------- CONSUMIDOR ----------
private void bucleConsumidor() {
String jo = Thread.currentThread().getName();
LOG.log(Level.INFO, "[{0}] consumidor llest", jo);
try {
while (true) {
// take() BLOQUEJA sense consumir CPU si la cua esta buida.
// Es l'alternativa correcta al sondeig amb sleep (08-03).
Reserva r = cua.take();
// PINDOLA VERINOSA: senyal d'apagada (apartat 12).
if (r == Reserva.FI) {
LOG.log(Level.INFO, "[{0}] pindola rebuda, acabo", jo);
return;
}
try {
atendre(r);
ateses.incrementAndGet();
} catch (Exception e) {
// Una fallada individual NO ha de matar el consumidor:
// si mor, el pool perd capacitat silenciosament.
LOG.log(Level.WARNING, "[" + jo + "] reserva fallida: " + r.id(), e);
}
}
} catch (InterruptedException e) {
LOG.log(Level.INFO, "[{0}] consumidor interromput", jo);
Thread.currentThread().interrupt(); // restaurar (08-02)
}
}
private void atendre(Reserva r) throws InterruptedException {
TimeUnit.MILLISECONDS.sleep(120); // E/S simulada
if (r.id().hashCode() % 23 == 0) {
throw new IllegalStateException("material ja prestat");
}
}
// ---------- APAGADA ----------
/**
* Apagada ORDENADA: deixa d'acceptar, insereix una pindola per
* consumidor i espera. Tot el que ja era a la cua es processa.
*/
@Override
public void close() {
acceptantNoves = false;
try {
for (int i = 0; i < CONSUMIDORS; i++) {
cua.put(Reserva.FI); // una per consumidor
}
consumidors.shutdown();
if (!consumidors.awaitTermination(30, TimeUnit.SECONDS)) {
consumidors.shutdownNow();
}
} catch (InterruptedException e) {
consumidors.shutdownNow();
Thread.currentThread().interrupt();
}
LOG.log(Level.INFO, "Processador tancat: {0} ateses, {1} rebutjades",
new Object[] { ateses.get(), rebutjades.get() });
}
public int ateses() { return ateses.get(); }
public int rebutjades() { return rebutjades.get(); }
public int pendents() { return cua.size(); }
}Ús:
public class DemostracioReserves {
public static void main(String[] args) throws InterruptedException {
try (ProcessadorReserves processador = new ProcessadorReserves()) {
// Tres productors: tres empleats fent servir el menu alhora.
Thread[] productors = new Thread[3];
for (int p = 0; p < 3; p++) {
final String empleat = switch (p) {
case 0 -> "Marta Ruiz";
case 1 -> "Diego Alonso";
default -> "Nuria Vidal";
};
productors[p] = new Thread(() -> {
try {
for (int i = 1; i <= 40; i++) {
processador.encuar(new Reserva(
empleat + "-R" + i, "978-000000000" + (i % 3 + 1)));
TimeUnit.MILLISECONDS.sleep(20);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "menu-" + empleat.split(" ")[0]);
productors[p].start();
}
// Progres mentre treballen.
for (int t = 0; t < 10; t++) {
System.out.printf(" ateses=%d pendents=%d%n",
processador.ateses(), processador.pendents());
TimeUnit.MILLISECONDS.sleep(400);
}
for (Thread p : productors) p.join();
System.out.println("Tots els productors han acabat");
} // close(): pindoles verinoses i espera al buidatge
}
}Sortida:
ateses=0 pendents=3
ateses=12 pendents=14
ateses=25 pendents=26
ateses=38 pendents=39
ateses=51 pendents=50 <-- cua PLENA: els productors es frenen
ateses=64 pendents=50
ateses=77 pendents=43
ateses=90 pendents=30
ateses=103 pendents=17
ateses=113 pendents=7
Tots els productors han acabat
INFO: Processador tancat: 120 ateses, 0 rebutjadesLa línia clau és on pendents s'estanca en 50. La cua va arribar a la seva capacitat i els productors van començar a bloquejar-se a put(): van deixar de produir al ritme que volien i van passar a produir al ritme que el sistema podia absorbir. Això és contrapressió, i és la propietat que impedeix que un sistema s'enfonsi sota càrrega. Amb una cua il·limitada, pendents hauria continuat creixent fins a exhaurir la memòria.
- La píndola verinosa
Com se li diu a un consumidor bloquejat a take() que ja no hi haurà més feina? Hi ha dues formes, i una és millor.
Opció A: interrompre. Funciona —take() llança InterruptedException— però és brusca: si el consumidor era a mig processar un element, aquella feina es perd.
Opció B: la píndola verinosa (poison pill). S'insereix a la cua un element sentinella que significa "s'ha acabat". El consumidor el reconeix i acaba ordenadament, després d'haver processat tot el que hi havia al davant.
package com.nexussoftware.bibliotech.domini;
public record Reserva(String id, String isbn) {
/**
* PINDOLA VERINOSA: instancia sentinella que significa
* "no hi haura mes feina, acaba".
*
* Es compara amb == (identitat), no amb equals: es un objecte
* unic i irrepetible, i cap dada real no pot coincidir amb ell.
*/
public static final Reserva FI = new Reserva("__FI__", "__FI__");
}Les tres regles de la píndola verinosa:
1. Una píndola per consumidor. Cada consumidor en consumeix una i acaba; si n'insereixes una de sola amb quatre consumidors, tres es queden esperant per sempre.
2. La píndola va al final. Com que la cua és FIFO, tot l'inserit abans es processa abans. L'apagada és ordenada per construcció: no es perd ni un element.
3. Comparar amb ==, no amb equals. La píndola és una instància única; comparar per identitat és més ràpid i no es pot confondre amb una dada real que casualment sigui igual.
Compte amb la variant de "reinjectar la píndola", que es veu de vegades:
// Alternativa: una sola pindola que cada consumidor reinjecta.
if (r == Reserva.FI) {
cua.put(Reserva.FI); // passar-la al seguent
return;
}És enginyosa i perillosa amb una cua acotada: si la cua està plena, aquell put() es bloqueja i el consumidor no acaba mai. Amb offer() al seu lloc podries perdre la píndola. Una píndola per consumidor és més simple i sempre correcte.
Comparació:
| Píndola verinosa | Interrupció | |
|---|---|---|
| Feina pendent a la cua | Es processa | Es perd |
| Feina en curs | Acaba | S'interromp |
| Requereix un valor sentinella | Sí | No |
| Funciona si un consumidor està encallat | No | Sí |
| Quan fer-la servir | Apagada ordenada | Apagada urgent |
A la pràctica es fan servir totes dues: píndola primer, i shutdownNow() com a pla B si no acaben dins del termini. És el patró de dues fases de 08-05, aplicat aquí.
ConcurrentSkipListMap en una nota
ConcurrentSkipListMap en una notaConcurrentSkipListMap i ConcurrentSkipListSet són les versions concurrents i ordenades de TreeMap i TreeSet. Implementen NavigableMap, així que ofereixen firstKey, headMap, tailMap, ceilingKey i companyia, i són segures per a diversos fils sense panys.
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.NavigableMap;
// Prestecs ordenats per data de venciment (en milisegons).
NavigableMap<Long, Prestec> perVenciment = new ConcurrentSkipListMap<>();
perVenciment.put(vencimentMs, prestec);
// Tots els vencuts abans d'ara, en ordre, sense bloquejar ningu.
NavigableMap<Long, Prestec> vencuts =
perVenciment.headMap(System.currentTimeMillis(), true);Estan implementades amb llistes de salts (skip lists), una estructura probabilística que dona O(log n) sense necessitat de reequilibrar com un arbre —cosa que seria molt costosa de fer concurrentment—. Fes-les servir quan necessitis ordre i concurrència; si només necessites concurrència, ConcurrentHashMap és més ràpid.
- Variables atòmiques i compare-and-swap
Segona meitat de la lliçó. Les classes de java.util.concurrent.atomic resolen el problema del comptador de 08-01 sense panys.
import java.util.concurrent.atomic.AtomicInteger;
public class ComptadorAtomic {
private final AtomicInteger valor = new AtomicInteger(0);
public void incrementar() {
valor.incrementAndGet(); // ATOMIC, sense pany
}
public int valor() { return valor.get(); }
public static void main(String[] args) throws InterruptedException {
final int VOLTES = 1_000_000;
ComptadorAtomic c = new ComptadorAtomic();
Thread f1 = new Thread(() -> { for (int i = 0; i < VOLTES; i++) c.incrementar(); });
Thread f2 = new Thread(() -> { for (int i = 0; i < VOLTES; i++) c.incrementar(); });
f1.start(); f2.start(); f1.join(); f2.join();
System.out.println("Esperat: " + VOLTES * 2 + ", real: " + c.valor());
}
}Exacte, sempre. Com, sense pany?
La instrucció compare-and-swap
Els processadors moderns ofereixen una instrucció atòmica anomenada CAS (compare-and-swap, o CMPXCHG en x86). La seva semàntica, executada pel maquinari de forma indivisible:
«Mira aquesta posició de memòria. Si conté el valor
esperat, substitueix-lo pernoui digues-me que sí. Si no, no toquis res i digues-me que no.»
És l'operació que sosté tota la concurrència sense panys.
AtomicInteger a = new AtomicInteger(10);
// "Si val 10, posa'l a 11". Retorna true si ho ha fet.
boolean exit = a.compareAndSet(10, 11);El bucle de reintent que hi ha a sota
incrementAndGet() no és màgia: és un bucle CAS. La seva implementació, conceptualment:
/**
* El que fa incrementAndGet() per dins, simplificat.
* Es el patro BUCLE CAS, la base de tota la programacio sense panys.
*/
public int incrementAndGet() {
while (true) {
int actual = get(); // 1. llegir
int seguent = actual + 1; // 2. calcular
if (compareAndSet(actual, seguent)) { // 3. intentar escriure
return seguent; // exit: sortir
}
// Fallada: un altre fil ha canviat el valor entre 1 i 3.
// No s'ha perdut res: es torna a llegir i es reintenta.
}
}La diferència essencial amb un pany:
- Un pany diu: "que ningú més no toqui això mentre treballo". Els altres esperen bloquejats.
- El CAS diu: "ho intento; si algú se m'ha avançat, ho torno a intentar". Ningú no espera bloquejat; sempre hi ha algú progressant.
Aquesta última propietat s'anomena llibertat de bloqueig (lock-free): en qualsevol moment, almenys un fil avança. No hi pot haver interbloqueig, perquè no hi ha res a retenir.
sequenceDiagram
participant A as fil-A
participant M as AtomicInteger
participant B as fil-B
Note over M: valor = 10
A->>M: get() -> 10
B->>M: get() -> 10
A->>M: compareAndSet(10, 11)
Note over M: coincideix: valor = 11, retorna true
B->>M: compareAndSet(10, 11)
Note over M: NO coincideix (val 11): retorna false
Note over B: reintenta
B->>M: get() -> 11
B->>M: compareAndSet(11, 12)
Note over M: coincideix: valor = 12
Note over A,B: Dos increments, valor = 12.<br/>CAP de perdut.
Compara aquest diagrama amb el de 08-04: allà els dos fils escrivien 11 i es perdia un increment. Aquí el segon fil detecta que algú se li ha avançat i reintenta. Aquesta detecció és el que fa el CAS.
Cost: sota contenció extrema, un bucle CAS pot reintentar moltes vegades i malbaratar CPU. Amb contenció moderada, és més ràpid que un pany perquè no hi ha canvis de context ni suspensió de fils.
- Les operacions de les classes atòmiques
Les classes principals: AtomicInteger, AtomicLong, AtomicBoolean, AtomicReference<V>, més els arrays AtomicIntegerArray, AtomicLongArray i AtomicReferenceArray.
| Mètode | Què fa | Retorna |
|---|---|---|
get() / set(v) |
Llegir / escriure (com volatile) |
valor / void |
incrementAndGet() |
++v |
El valor nou |
getAndIncrement() |
v++ |
El valor anterior |
decrementAndGet() / getAndDecrement() |
--v / v-- |
nou / anterior |
addAndGet(d) / getAndAdd(d) |
Sumar d |
nou / anterior |
getAndSet(v) |
Escriure i retornar el que hi havia | anterior |
compareAndSet(esp, nou) |
CAS | boolean |
updateAndGet(f) |
Aplicar la funció f atòmicament |
nou |
getAndUpdate(f) |
Ídem | anterior |
accumulateAndGet(x, f) |
Combinar el valor actual amb x mitjançant f |
nou |
package com.nexussoftware.bibliotech.servei;
import java.util.concurrent.atomic.*;
/**
* Estadistiques de BiblioTech amb comptadors atomics.
* Sense un sol pany, i amb lectures que no bloquegen res.
*/
public class EstadistiquesBiblioTech {
private final AtomicLong prestecsTotals = new AtomicLong();
private final AtomicLong devolucionsTotals = new AtomicLong();
private final AtomicLong multesRecaptadesCentims = new AtomicLong();
private final AtomicInteger prestecsActius = new AtomicInteger();
private final AtomicBoolean modeManteniment = new AtomicBoolean(false);
private final AtomicLong maximSimultanis = new AtomicLong();
public void registrarPrestec() {
prestecsTotals.incrementAndGet();
int actius = prestecsActius.incrementAndGet();
// MAXIM HISTORIC amb accumulateAndGet: combina el valor
// actual amb 'actius' fent servir Math::max, atomicament.
// Escriure-ho amb get()+set() seria una cursa classica.
maximSimultanis.accumulateAndGet(actius, Math::max);
}
public void registrarDevolucio(long multaCentims) {
devolucionsTotals.incrementAndGet();
prestecsActius.decrementAndGet();
if (multaCentims > 0) {
multesRecaptadesCentims.addAndGet(multaCentims);
}
}
/**
* Aplica un descompte del 10% al total recaptat.
* updateAndGet aplica la funcio ATOMICAMENT, amb un bucle CAS
* per sota: si un altre fil modifica el valor entre la lectura i
* l'escriptura, la funcio es TORNA A EXECUTAR amb el valor nou.
*
* IMPORTANT: la funcio ha de ser PURA i rapida, perque es pot
* executar diverses vegades. Res d'efectes secundaris.
*/
public long aplicarDescompte() {
return multesRecaptadesCentims.updateAndGet(v -> (long) (v * 0.9));
}
/**
* Entrar en manteniment NOMES SI no hi erem ja.
* compareAndSet garanteix que, encara que deu fils ho demanin alhora,
* exactament UN rebi true i executi la preparacio.
* Es l'idioma "nomes una vegada" sense panys.
*/
public boolean entrarEnManteniment() {
if (modeManteniment.compareAndSet(false, true)) {
prepararManteniment(); // nomes un fil hi arriba
return true;
}
return false; // un altre se'ns ha avancat
}
public String resum() {
return String.format(
"prestecs=%d devolucions=%d actius=%d maxim=%d multes=%.2f EUR",
prestecsTotals.get(), devolucionsTotals.get(),
prestecsActius.get(), maximSimultanis.get(),
multesRecaptadesCentims.get() / 100.0);
}
private void prepararManteniment() { /* ... */ }
}Nota sobre
resum(): els sis valors es llegeixen en instants diferents, així que el conjunt no és una instantània coherent. Pot mostrarprestecs=100iactius=3de dos moments diferents. Per a mètriques és perfectament acceptable; si necessitessis coherència entre tots els comptadors, caldrien un pany o el patró d'estat immutable ambAtomicReferencede l'apartat següent.
AtomicReference i el problema ABA
AtomicReference i el problema ABAAtomicReference<V> aplica el CAS a una referència a objecte. Combinat amb la immutabilitat de 08-04, dona un patró molt potent: actualitzar un estat complet de forma atòmica i sense panys.
package com.nexussoftware.bibliotech.servei;
import java.util.concurrent.atomic.AtomicReference;
public class EstatBiblioTech {
/** Estat complet, immutable (record de 04-07). */
public record Estat(long prestecs, long devolucions,
long multesCentims, boolean manteniment) {
Estat ambPrestec() {
return new Estat(prestecs + 1, devolucions, multesCentims, manteniment);
}
Estat ambDevolucio(long multa) {
return new Estat(prestecs, devolucions + 1,
multesCentims + multa, manteniment);
}
}
private final AtomicReference<Estat> estat =
new AtomicReference<>(new Estat(0, 0, 0, false));
/**
* Lectura sense cap bloqueig i SEMPRE COHERENT: els quatre
* camps venen del mateix instant, perque son un sol objecte
* immutable. Es l'avantatge sobre quatre comptadors separats.
*/
public Estat instantania() {
return estat.get();
}
/** Actualitzacio atomica dels quatre camps alhora. */
public void registrarPrestec() {
estat.updateAndGet(Estat::ambPrestec);
}
public void registrarDevolucio(long multaCentims) {
estat.updateAndGet(e -> e.ambDevolucio(multaCentims));
}
}Aquesta combinació —estat immutable + AtomicReference + updateAndGet— és una de les tècniques més elegants de la concurrència en Java: lectures gratuïtes i sempre coherents, escriptures atòmiques sense panys, i cap possibilitat d'interbloqueig.
El problema ABA
Hi ha un cas patològic del CAS que convé conèixer.
El CAS comprova que el valor sigui l'esperat, no que no hagi canviat. Si un valor passa d'A a B i torna a A, un CAS que esperava A tindrà èxit, encara que entremig hagi passat alguna cosa important.
Fil 1: llegeix A ............................... CAS(A -> C): EXIT
Fil 2: llegeix A, CAS(A->B), CAS(B->A)
^
El fil 1 no s'assabenta que hi va haver dos canvis.
Si la seva decisio depenia que NO hagues passat res, es un bug.Amb comptadors d'enters això és innocu: si el comptador torna a valer 10, val 10 i punt. El problema apareix amb referències en estructures enllaçades: un node pot ser retirat, reciclat i reinserit, i un CAS sobre la seva referència tindria èxit sobre un node que ja no és el mateix lògicament.
La solució: afegir un segell o marca que canviï sempre.
import java.util.concurrent.atomic.AtomicStampedReference;
// Cada modificacio incrementa el SEGELL, encara que el valor torni a ser
// el mateix. El CAS comprova valor I segell, aixi que A-B-A es detecta.
AtomicStampedReference<Node> capcalera =
new AtomicStampedReference<>(nodeInicial, 0);
int[] segellActual = new int[1];
Node actual = capcalera.get(segellActual);
capcalera.compareAndSet(actual, nouNode,
segellActual[0], segellActual[0] + 1); // segell + 1
// Variant mes simple amb un sol bit boolea:
// AtomicMarkableReference<Node>A la pràctica, ABA només apareix si implementes estructures de dades sense panys a mà. Si fas servir ConcurrentHashMap i AtomicInteger, la biblioteca ja se n'ha ocupat. Coneix-lo per saber que existeix i per entendre per què AtomicStampedReference hi és.
LongAdder sota contenció
LongAdder sota contencióAtomicLong és excel·lent amb contenció moderada. Sota contenció extrema —setze fils incrementant el mateix comptador sense parar— el seu bucle CAS comença a fallar molt: els fils reintenten una vegada i una altra i, a més, tots escriuen sobre la mateixa línia de memòria cau, cosa que provoca invalidacions constants entre nuclis (cache line bouncing).
LongAdder (Java 8) resol això amb una idea senzilla: manté diverses cel·les internes i cada fil incrementa la seva. Només en cridar sum() se sumen totes.
import java.util.concurrent.atomic.LongAdder;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.*;
public class ComparativaComptadors {
static long mesurarAtomicLong(int fils, int ops) throws InterruptedException {
AtomicLong c = new AtomicLong();
return mesurar(fils, ops, c::incrementAndGet, c::get);
}
static long mesurarLongAdder(int fils, int ops) throws InterruptedException {
LongAdder c = new LongAdder();
return mesurar(fils, ops, c::increment, c::sum);
}
static long mesurar(int fils, int ops, Runnable increment,
java.util.function.LongSupplier lectura) throws InterruptedException {
CountDownLatch sortida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(fils);
for (int f = 0; f < fils; f++) {
new Thread(() -> {
try {
sortida.await();
for (int i = 0; i < ops; i++) increment.run();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally { meta.countDown(); }
}, "comptador-" + f).start();
}
long inici = System.nanoTime();
sortida.countDown();
meta.await();
long ns = System.nanoTime() - inici;
// Verificacio: tots dos han de donar el resultat exacte.
if (lectura.getAsLong() != (long) fils * ops) {
throw new AssertionError("resultat incorrecte");
}
return ns / 1_000_000;
}
public static void main(String[] args) throws InterruptedException {
final int OPS = 2_000_000;
System.out.printf("%-8s | %14s | %14s | %8s%n",
"Fils", "AtomicLong ms", "LongAdder ms", "Millora");
System.out.println("---------|----------------|----------------|---------");
for (int fils : new int[] { 1, 2, 4, 8, 16 }) {
long a = mesurarAtomicLong(fils, OPS / fils);
long l = mesurarLongAdder(fils, OPS / fils);
System.out.printf("%-8d | %14d | %14d | %7.1fx%n",
fils, a, l, (double) a / Math.max(l, 1));
}
}
}Sortida orientativa:
Fils | AtomicLong ms | LongAdder ms | Millora
---------|----------------|----------------|---------
1 | 12 | 18 | 0.7x
2 | 41 | 21 | 2.0x
4 | 96 | 19 | 5.1x
8 | 213 | 22 | 9.7x
16 | 487 | 26 | 18.7xEl que ensenya aquesta taula:
- Amb un fil,
AtomicLongguanya.LongAdderté més maquinària interna i no l'amortitza sense contenció. - L'avantatge creix amb els fils. A 16 fils, gairebé 19×.
AtomicLongescala malament: passar d'1 a 16 fils multiplica el temps per 40, encara que la feina total sigui la mateixa. És contenció pura.
AtomicLong |
LongAdder |
|
|---|---|---|
| Escriptura sota contenció | Es degrada | Excel·lent |
Lectura (get/sum) |
O(1) exacta | O(nre. de cel·les), aproximada si hi ha escriptures concurrents |
| Memòria | 1 valor | Diverses cel·les (creix amb la contenció) |
Admet compareAndSet |
Sí | No |
| Fer servir per a | Comptadors amb poca contenció; quan necessites CAS | Mètriques i estadístiques d'alta freqüència |
La regla: si només comptes i llegeixes el total de tant en tant —mètriques, comptadors de peticions, estadístiques—, LongAdder. Si necessites el valor exacte a cada operació o fer servir compareAndSet, AtomicLong. DoubleAdder, LongAccumulator i DoubleAccumulator completen la família; els Accumulator permeten una funció de combinació arbitrària.
- BiblioTech: catàleg concurrent i estadístiques atòmiques
La versió final del catàleg, sense un sol lock() explícit:
package com.nexussoftware.bibliotech.servei;
import com.nexussoftware.bibliotech.domini.Material;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.LongAdder;
/**
* Cataleg concurrent de BiblioTech.
*
* Davant de la versio amb ReadWriteLock de 08-04:
* - No hi ha cap pany explicit que deixar anar en un finally.
* - Les lectures no bloquegen RES, ni tan sols altres lectures.
* - Les escriptures en cubetes diferents ocorren en parallel.
* - Les operacions compostes es resolen amb els metodes atomics
* del mateix mapa, no encapsulant seccions critiques a ma.
*
* LIMIT HONEST: un invariant que abasti DIVERSES estructures
* (com el de RegistrePrestecs entre perId i perEmpleat) continua
* necessitant un pany. Les colleccions concurrents garanteixen
* l'atomicitat d'UNA operacio sobre UNA colleccio, no de dues.
*/
public class CatalegConcurrent {
/** Index principal per ISBN. Lectures sense bloqueig. */
private final ConcurrentMap<String, Material> perIsbn = new ConcurrentHashMap<>();
/** Materials agrupats per tipus. El valor es una llista concurrent. */
private final ConcurrentMap<TipusMaterial, List<Material>> perTipus =
new ConcurrentHashMap<>();
/** Oients d'esdeveniments: moltes lectures, gairebe cap escriptura. */
private final List<OientCataleg> oients = new CopyOnWriteArrayList<>();
// --- Estadistiques: LongAdder per la seva alta frequencia d'escriptura ---
private final LongAdder consultes = new LongAdder();
private final LongAdder encerts = new LongAdder();
private final LongAdder altes = new LongAdder();
private final AtomicLong ultimaModificacioMs = new AtomicLong();
// ---------- ESCRIPTURA ----------
/**
* Afegeix un material si el seu ISBN no hi era.
* putIfAbsent es ATOMIC: encara que deu fils afegeixin el mateix ISBN
* alhora, exactament un rep true. Es el comprovar-despres-actuar
* resolt sense panys.
*/
public boolean afegir(Material m) {
if (perIsbn.putIfAbsent(m.isbn(), m) != null) {
return false; // ja existia
}
// computeIfAbsent crea la llista nomes si no existeix, atomicament.
// CopyOnWriteArrayList perque aquestes llistes es recorren molt
// mes del que es modifiquen.
perTipus.computeIfAbsent(m.tipus(), t -> new CopyOnWriteArrayList<>()).add(m);
altes.increment();
ultimaModificacioMs.set(System.currentTimeMillis());
notificarAlta(m);
return true;
}
public boolean eliminar(String isbn) {
Material m = perIsbn.remove(isbn);
if (m == null) return false;
List<Material> llista = perTipus.get(m.tipus());
if (llista != null) llista.remove(m);
ultimaModificacioMs.set(System.currentTimeMillis());
return true;
}
// ---------- LECTURA ----------
/** Sense panys. Amb 100 fils consultant, cap no espera cap altre. */
public Material cercarPerIsbn(String isbn) {
consultes.increment();
Material m = perIsbn.get(isbn);
if (m != null) encerts.increment();
return m;
}
/**
* Retorna la llista d'un tipus. Com que es CopyOnWriteArrayList,
* qui crida la pot iterar amb total seguretat encara que un altre fil
* la modifiqui: el seu iterador treballa sobre una instantania immutable.
* No cal copiar defensivament, a diferencia de 08-04.
*/
public List<Material> perTipus(TipusMaterial tipus) {
return perTipus.getOrDefault(tipus, List.of());
}
/**
* Recorregut complet del cataleg.
* L'iterador es DEBILMENT CONSISTENT: no llanca
* ConcurrentModificationException i no bloqueja els escriptors.
*/
public void recorrer(java.util.function.Consumer<Material> accio) {
for (Material m : perIsbn.values()) {
accio.accept(m);
}
}
/** COMPTE: aproximat (apartat 7). Val per a metriques, no per a control. */
public int midaAproximada() {
return perIsbn.size();
}
// ---------- OIENTS ----------
public void registrarOient(OientCataleg o) { oients.add(o); }
private void notificarAlta(Material m) {
// Sense pany durant la notificacio: la regla 3 de 08-04
// ("no cridis mai codi alie amb un pany") es compleix
// automaticament gracies a CopyOnWriteArrayList.
for (OientCataleg o : oients) {
try {
o.materialAfegit(m);
} catch (RuntimeException e) {
LOG.log(Level.WARNING, "Oient fallit", e);
}
}
}
// ---------- METRIQUES ----------
public String metriques() {
long c = consultes.sum();
long a = encerts.sum();
return String.format("consultes=%d encerts=%d taxa=%.1f%% altes=%d materials~%d",
c, a, c == 0 ? 0.0 : 100.0 * a / c, altes.sum(), perIsbn.size());
}
}Prova d'esforç:
public class ProvaCatalegConcurrent {
public static void main(String[] args) throws InterruptedException {
final int FILS = 16;
final int OPS = 100_000;
CatalegConcurrent cataleg = new CatalegConcurrent();
CountDownLatch sortida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(FILS);
for (int f = 0; f < FILS; f++) {
final int id = f;
new Thread(() -> {
try {
sortida.await();
ThreadLocalRandom atzar = ThreadLocalRandom.current();
for (int i = 0; i < OPS; i++) {
if (atzar.nextInt(100) < 90) {
cataleg.cercarPerIsbn("978-" + atzar.nextInt(1000));
} else {
cataleg.afegir(new Llibre(
"978-" + atzar.nextInt(1000), "Titol", "Autor"));
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally { meta.countDown(); }
}, "cataleg-" + id).start();
}
long inici = System.nanoTime();
sortida.countDown();
meta.await();
long ms = (System.nanoTime() - inici) / 1_000_000;
System.out.println("Operacions : " + (FILS * OPS));
System.out.println("Temps : " + ms + " ms");
System.out.println("Rendiment : " + (FILS * OPS / Math.max(ms, 1)) + " ops/ms");
System.out.println(cataleg.metriques());
}
}Sortida orientativa:
Operacions : 1600000
Temps : 312 ms
Rendiment : 5128 ops/ms
Metriques : consultes=1439871 encerts=1438204 taxa=99.9% altes=1000 materials~1000Un milió sis-centes mil operacions amb setze fils en 312 ms, sense un sol lock(). I el detall que valida el disseny: altes=1000 amb 1000 ISBN possibles. Encara que setze fils van intentar afegir els mateixos ISBN repetidament, exactament mil altes van tenir èxit: el putIfAbsent va complir el seu contracte sota màxima contenció.
- Taula final de decisió
| Necessites… | Fes servir | Per què |
|---|---|---|
| Un comptador d'alta freqüència | LongAdder |
Escala sota contenció |
| Un comptador amb valor exacte a cada operació | AtomicInteger/AtomicLong |
CAS, exacte, sense panys |
| Un marcador amb "només una vegada" | AtomicBoolean.compareAndSet |
Exactament un guanyador |
| Un estat de diversos camps, coherent en llegir | AtomicReference + record |
Instantània immutable i atòmica |
| Un marcador simple d'aturada | volatile boolean |
El més barat que garanteix visibilitat |
| Un mapa compartit | ConcurrentHashMap |
Lectures sense bloqueig, compostes atòmiques |
| Un mapa compartit i ordenat | ConcurrentSkipListMap |
Navegable i concurrent |
| Una llista d'oients o de configuració | CopyOnWriteArrayList |
Iteració immune i sense bloqueig |
| Una llista gran amb escriptures freqüents | synchronizedList o un pany |
La còpia-en-escriure seria O(n²) |
| Traspassar feina entre fils | BlockingQueue |
take/put sense sondeig, amb contrapressió |
| Una cua sense espera | ConcurrentLinkedQueue |
No blocant, però obliga a sondejar |
| Un invariant entre diverses estructures | synchronized o Lock |
Les col·leccions concurrents no ho cobreixen |
| Moltes lectures i poques escriptures sobre estat propi | ReadWriteLock |
Lectors en paral·lel |
| Dades que no canvien | Objecte immutable (record) |
Zero sincronització, impossible corrompre |
| Dades utilitzades per un sol fil | Variable local / ThreadLocal |
La millor opció: no compartir |
L'ordre en què cal plantejar-se les opcions, de millor a pitjor:
- No compartir (local, confinat).
- Compartir immutable (
record). - Atòmic (
AtomicX,LongAdder). - Col·lecció concurrent (
ConcurrentHashMap,BlockingQueue). - Pany (
synchronized,Lock,ReadWriteLock).
Errors Comuns i Consells
Error 1: compartir un HashMap entre fils. No dona resultats "aproximats": perd entrades i es pot corrompre fins a provocar un bucle infinit amb un nucli al 100 %.
Error 2: creure que Collections.synchronizedMap fa segur el teu codi. Cada mètode ho és; iterar i les operacions compostes, no. És l'error més freqüent de la primera generació.
Error 3: iterar una col·lecció sincronitzada sense sincronitzar la iteració. ConcurrentModificationException. I si sincronitzes, bloqueges tothom durant tot el recorregut.
Error 4: fer get + put sobre un ConcurrentHashMap. El mapa és segur; la teva seqüència de dues crides no. Fes servir merge, compute, computeIfAbsent o putIfAbsent.
Error 5: fer servir size() d'una col·lecció concurrent per prendre decisions. És aproximat per disseny. Val per a mètriques, no per a control.
Error 6: fer feina pesada o blocant dins de computeIfAbsent. S'executa amb la cubeta bloquejada. Calcula fora i publica amb putIfAbsent.
Error 7: modificar el mateix mapa dins de la funció de compute. Pot interbloquejar o corrompre l'estructura. Prohibit.
Error 8: fer servir CopyOnWriteArrayList per acumular en un bucle. Cada add copia l'array: O(n²) total. És per a moltes lectures i gairebé cap escriptura.
Error 9: sondejar una ConcurrentLinkedQueue amb poll() + sleep(). Crema CPU i afegeix latència. Fes servir una BlockingQueue i take().
Error 10: fer servir una BlockingQueue il·limitada com a cua de treball. Sense cota no hi ha contrapressió, i el productor ràpid acaba exhaurint la memòria.
Error 11: inserir una sola píndola verinosa amb diversos consumidors. Els altres esperen per sempre. Una per consumidor.
Error 12: passar una funció amb efectes secundaris a updateAndGet. Es pot executar diverses vegades pel bucle CAS. Ha de ser pura i ràpida.
Error 13: creure que les col·leccions concurrents eliminen la necessitat de panys. Garanteixen l'atomicitat d'una operació sobre una col·lecció. Un invariant entre dues estructures continua necessitant un pany.
Consell 1: ConcurrentHashMap per defecte per a qualsevol mapa compartit. És més ràpid, més segur i amb millor API que les alternatives. No hi ha motiu per no fer-lo servir.
Consell 2: aprèn merge i computeIfAbsent de memòria. Resolen el 90 % dels casos de comprovar-després-actuar en una línia llegible i atòmica.
Consell 3: acota les teves cues. La cota és la diferència entre un sistema que es degrada amb elegància i un que cau.
Consell 4: LongAdder per a mètriques, AtomicLong per a lògica. Si només comptes, LongAdder. Si necessites el valor exacte a cada pas o compareAndSet, AtomicLong.
Consell 5: AtomicReference + record és la teva millor eina per a estat compartit de diversos camps. Lectures gratuïtes i coherents, escriptures atòmiques, impossible interbloquejar.
Consell 6: la píndola verinosa dona una apagada ordenada; la interrupció, una d'urgent. Fes servir la primera i guarda la segona com a pla B amb termini.
Exercicis
Exercici 1: Les tres generacions, mesurades
Escriu TresGeneracions que compari HashMap sense protecció, Collections.synchronizedMap i ConcurrentHashMap sota la mateixa càrrega: 12 fils, 100.000 operacions cadascun, 80 % lectures i 20 % escriptures sobre un espai de 5.000 claus. Per a cada implementació mesura el temps, comprova si el nombre final d'entrades és l'esperat, i captura qualsevol excepció. Fes servir un CountDownLatch de porta de sortida i un termini per si el HashMap entra en bucle infinit. Explica els tres resultats.
Exercici 2: Productor-consumidor complet amb contrapressió
Implementa PipelineImportacio, un processament en dues etapes per a BiblioTech:
- Etapa 1 (2 fils productors): llegeixen "línies" d'un catàleg simulat i les posen en una
ArrayBlockingQueue<String>de capacitat 20. - Etapa 2 (4 fils consumidors): treuen línies, les converteixen en
Material(amb unsleepde 30 ms) i les afegeixen a unConcurrentHashMap.
Requisits: comptadors amb LongAdder per a línies llegides, materials creats i línies descartades; un fil monitor que imprimeixi cada 300 ms la mida de la cua i els comptadors, demostrant que la cua s'omple i frena els productors; apagada amb píndola verinosa (una per consumidor); i verificació final que no s'ha perdut ni una línia.
Exercici 3: Comptador d'estadístiques, quatre implementacions
Escriu ComparativaEstadistiques que implementi el mateix comptador de préstecs de quatre formes: (a) long amb synchronized, (b) AtomicLong, (c) LongAdder, i (d) AtomicReference<Estat> amb un record immutable de tres camps actualitzat amb updateAndGet. Sotmet cadascuna a 16 fils × 500.000 increments, verifica que totes donen el resultat exacte, i mesura el temps. Afegeix una segona mesura amb un sol fil per mostrar la inversió de resultats. Comenta quina implementació triaries per a mètriques d'alta freqüència i quina per a un estat de negoci que s'ha de llegir de forma coherent.
Solucions
Solució a l'Exercici 1
import java.util.*;
import java.util.concurrent.*;
public class TresGeneracions {
static final int FILS = 12;
static final int OPS = 100_000;
static final int CLAUS = 5_000;
record Resultat(String nom, long ms, int entrades, String incidencia) { }
static Resultat mesurar(String nom, Map<Integer, String> mapa)
throws InterruptedException {
CountDownLatch sortida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(FILS);
// volatile no n'hi ha prou per acumular text des de diversos fils:
// fem servir una cua concurrent per recollir incidencies.
Queue<String> incidencies = new ConcurrentLinkedQueue<>();
for (int f = 0; f < FILS; f++) {
Thread t = new Thread(() -> {
try {
sortida.await();
ThreadLocalRandom atzar = ThreadLocalRandom.current();
for (int i = 0; i < OPS; i++) {
int clau = atzar.nextInt(CLAUS);
if (atzar.nextInt(100) < 80) {
mapa.get(clau);
} else {
mapa.put(clau, "v" + i);
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (RuntimeException e) {
incidencies.add(e.getClass().getSimpleName());
} finally {
meta.countDown();
}
}, nom + "-" + f);
t.setDaemon(true); // dimoni: si entra en bucle infinit,
t.start(); // no impedira que la JVM acabi
}
long inici = System.nanoTime();
sortida.countDown();
// Termini: el HashMap sense proteccio pot no acabar MAI.
boolean completat = meta.await(20, TimeUnit.SECONDS);
long ms = (System.nanoTime() - inici) / 1_000_000;
String inc = completat
? (incidencies.isEmpty() ? "-" : incidencies.peek())
: "NO HA ACABAT (bucle infinit o corrupcio)";
int entrades;
try {
entrades = mapa.size();
} catch (RuntimeException e) {
entrades = -1;
}
return new Resultat(nom, ms, entrades, inc);
}
public static void main(String[] args) throws InterruptedException {
System.out.printf("%d fils x %d ops (80%% lectura) sobre %d claus%n%n",
FILS, OPS, CLAUS);
List<Resultat> resultats = new ArrayList<>();
resultats.add(mesurar("HashMap", new HashMap<>()));
resultats.add(mesurar("synchronizedMap", Collections.synchronizedMap(new HashMap<>())));
resultats.add(mesurar("ConcurrentHashMap", new ConcurrentHashMap<>()));
System.out.printf("%-20s | %8s | %10s | %-40s%n",
"Implementacio", "ms", "entrades", "incidencia");
System.out.println("---------------------|----------|------------|"
+ "------------------------------------------");
for (Resultat r : resultats) {
System.out.printf("%-20s | %8d | %10d | %-40s%n",
r.nom(), r.ms(), r.entrades(), r.incidencia());
}
System.out.println();
System.out.println("Esperat: " + CLAUS + " entrades (totes les claus tocades)");
}
}Sortida orientativa:
12 fils x 100000 ops (80% lectura) sobre 5000 claus
Implementacio | ms | entrades | incidencia
---------------------|----------|------------|------------------------------------------
HashMap | 20003 | 4211 | NO HA ACABAT (bucle infinit o corrupcio)
synchronizedMap | 1876 | 5000 | -
ConcurrentHashMap | 147 | 5000 | -
Esperat: 5000 entrades (totes les claus tocades)Els tres resultats:
HashMapno va acabar en 20 segons i el seusize()diu 4211 en lloc de 5000: s'han perdut entrades i almenys un fil es va quedar en un bucle. Marcar-los com a dimoni va ser essencial perquè el programa pogués acabar.synchronizedMapés correcte però lent: 1876 ms, amb dotze fils serialitzats per un únic pany, incloses les 80 % de lectures que no es destorbarien entre si.ConcurrentHashMapés correcte i 12× més ràpid que el sincronitzat: les lectures no bloquegen res i les escriptures només competeixen quan cauen a la mateixa cubeta.
Solució a l'Exercici 2
package com.nexussoftware.bibliotech.persistencia;
import java.util.concurrent.*;
import java.util.concurrent.atomic.LongAdder;
public class PipelineImportacio {
private static final String PINDOLA = "__FI__";
private static final int CAPACITAT = 20;
private static final int PRODUCTORS = 2;
private static final int CONSUMIDORS = 4;
private static final int LINIES_PER_PRODUCTOR = 150;
// Cua ACOTADA: si els consumidors no donen l'abast, els productors
// es bloquegen a put(). Contrapressio.
private final BlockingQueue<String> cua = new ArrayBlockingQueue<>(CAPACITAT);
private final ConcurrentMap<String, String> materials = new ConcurrentHashMap<>();
// LongAdder: escriptures molt frequents, lectures ocasionals.
private final LongAdder llegides = new LongAdder();
private final LongAdder creats = new LongAdder();
private final LongAdder descartades = new LongAdder();
private volatile boolean enMarxa = true;
// ---------- ETAPA 1: PRODUCTORS ----------
private void produir(int idProductor) {
try {
for (int i = 0; i < LINIES_PER_PRODUCTOR; i++) {
String linia = "978-" + idProductor + String.format("%05d", i)
+ ";Titol " + i + ";Autor";
cua.put(linia); // BLOQUEJA si la cua esta plena
llegides.increment();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// ---------- ETAPA 2: CONSUMIDORS ----------
private void consumir() {
String jo = Thread.currentThread().getName();
try {
while (true) {
String linia = cua.take(); // BLOQUEJA si esta buida
// Pindola verinosa: comparacio per IDENTITAT.
if (linia == PINDOLA) {
System.out.println(" [" + jo + "] pindola rebuda, acabo");
return;
}
try {
TimeUnit.MILLISECONDS.sleep(30); // conversio "cara"
String[] camps = linia.split(";");
if (camps.length < 3) {
throw new IllegalArgumentException("linia incompleta");
}
materials.put(camps[0], camps[1]);
creats.increment();
} catch (IllegalArgumentException e) {
descartades.increment(); // degradar (06-07)
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// ---------- ORQUESTRACIO ----------
public void executar() throws InterruptedException {
ExecutorService productors = Executors.newFixedThreadPool(PRODUCTORS,
new AnomenadorFils("pipeline-productor"));
ExecutorService consumidors = Executors.newFixedThreadPool(CONSUMIDORS,
new AnomenadorFils("pipeline-consumidor"));
for (int c = 0; c < CONSUMIDORS; c++) consumidors.execute(this::consumir);
CountDownLatch produccioAcabada = new CountDownLatch(PRODUCTORS);
for (int p = 0; p < PRODUCTORS; p++) {
final int id = p;
productors.execute(() -> {
try { produir(id); }
finally { produccioAcabada.countDown(); } // SEMPRE
});
}
// Monitor: demostra que la cua s'omple i frena els productors.
Thread monitor = new Thread(() -> {
try {
while (enMarxa) {
System.out.printf(" [monitor] cua=%2d/%d llegides=%d creats=%d%n",
cua.size(), CAPACITAT, llegides.sum(), creats.sum());
TimeUnit.MILLISECONDS.sleep(300);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "pipeline-monitor");
monitor.setDaemon(true);
monitor.start();
// 1. Esperar que els productors acabin.
produccioAcabada.await();
System.out.println(" Produccio acabada; injectant pindoles");
// 2. UNA pindola PER CONSUMIDOR, al final de la cua:
// tot l'anterior es processa abans.
for (int c = 0; c < CONSUMIDORS; c++) cua.put(PINDOLA);
// 3. Apagada ordenada en dues fases.
productors.shutdown();
consumidors.shutdown();
boolean ok = consumidors.awaitTermination(30, TimeUnit.SECONDS);
if (!ok) consumidors.shutdownNow();
enMarxa = false;
monitor.interrupt();
// 4. Verificacio.
long total = LINIES_PER_PRODUCTOR * PRODUCTORS;
System.out.println();
System.out.println("=== RESULTAT ===");
System.out.println("Linies produides : " + llegides.sum() + " (esperat " + total + ")");
System.out.println("Materials creats : " + creats.sum());
System.out.println("Descartades : " + descartades.sum());
System.out.println("Al mapa : " + materials.size());
System.out.println("Cua al final : " + cua.size() + " (ha de ser 0)");
System.out.println("Sense perdues : "
+ (creats.sum() + descartades.sum() == total));
}
static class AnomenadorFils implements ThreadFactory {
private final String prefix;
private int n = 1;
AnomenadorFils(String prefix) { this.prefix = prefix; }
@Override public synchronized Thread newThread(Runnable r) {
return new Thread(r, prefix + "-" + n++);
}
}
public static void main(String[] args) throws InterruptedException {
new PipelineImportacio().executar();
}
}Sortida (fragment):
[monitor] cua=20/20 llegides= 42 creats= 22
[monitor] cua=20/20 llegides= 82 creats= 62
[monitor] cua=20/20 llegides=122 creats=102
[monitor] cua=20/20 llegides=162 creats=142
...
Produccio acabada; injectant pindoles
[pipeline-consumidor-2] pindola rebuda, acabo
[pipeline-consumidor-1] pindola rebuda, acabo
[pipeline-consumidor-4] pindola rebuda, acabo
[pipeline-consumidor-3] pindola rebuda, acabo
=== RESULTAT ===
Linies produides : 300 (esperat 300)
Materials creats : 300
Descartades : 0
Al mapa : 300
Cua al final : 0 (ha de ser 0)
Sense perdues : trueTres coses que demostra la sortida:
cua=20/20de forma sostinguda: la cua està permanentment plena, així que els productors passen la major part del temps bloquejats aput(). Produeixen al ritme dels consumidors, no al seu. Amb una cua il·limitada,cuahauria crescut fins a 300 i tota la memòria intermèdia s'hauria reservat de cop.Cua al final: 0iSense perdues: true: la píndola verinosa, en anar al final d'una cua FIFO, garanteix que tota la feina anterior es processa abans de l'apagada. Ni una línia perduda.- Els quatre consumidors acaben, un per píndola. Amb una sola píndola, tres s'haurien quedat bloquejats a
take()per sempre iawaitTerminationhauria expirat.
Solució a l'Exercici 3
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
public class ComparativaEstadistiques {
interface Comptador {
void registrarPrestec();
long total();
String nom();
}
/** (a) long protegit amb synchronized. */
static class AmbSynchronized implements Comptador {
private long prestecs = 0;
public synchronized void registrarPrestec() { prestecs++; }
public synchronized long total() { return prestecs; }
public String nom() { return "synchronized"; }
}
/** (b) AtomicLong: bucle CAS per sota. */
static class AmbAtomicLong implements Comptador {
private final AtomicLong prestecs = new AtomicLong();
public void registrarPrestec() { prestecs.incrementAndGet(); }
public long total() { return prestecs.get(); }
public String nom() { return "AtomicLong"; }
}
/** (c) LongAdder: celles separades, suma en llegir. */
static class AmbLongAdder implements Comptador {
private final LongAdder prestecs = new LongAdder();
public void registrarPrestec() { prestecs.increment(); }
public long total() { return prestecs.sum(); }
public String nom() { return "LongAdder"; }
}
/** (d) AtomicReference sobre un record immutable de TRES camps. */
static class AmbAtomicReference implements Comptador {
record Estat(long prestecs, long devolucions, long multes) {
Estat ambPrestec() {
return new Estat(prestecs + 1, devolucions, multes);
}
}
private final AtomicReference<Estat> estat =
new AtomicReference<>(new Estat(0, 0, 0));
public void registrarPrestec() {
// updateAndGet reintenta si un altre fil s'ha avancat.
// La funcio ha de ser PURA: es pot executar diverses vegades.
estat.updateAndGet(Estat::ambPrestec);
}
public long total() { return estat.get().prestecs(); }
/** AVANTATGE UNIC: instantania COHERENT dels tres camps. */
public Estat instantania() { return estat.get(); }
public String nom() { return "AtomicReference+record"; }
}
static long mesurar(Comptador c, int fils, int perFil) throws InterruptedException {
CountDownLatch sortida = new CountDownLatch(1);
CountDownLatch meta = new CountDownLatch(fils);
for (int f = 0; f < fils; f++) {
new Thread(() -> {
try {
sortida.await();
for (int i = 0; i < perFil; i++) c.registrarPrestec();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally { meta.countDown(); }
}, "est-" + f).start();
}
long inici = System.nanoTime();
sortida.countDown();
meta.await();
long ms = (System.nanoTime() - inici) / 1_000_000;
long esperat = (long) fils * perFil;
if (c.total() != esperat) {
throw new AssertionError(c.nom() + " INCORRECTE: "
+ c.total() + " != " + esperat);
}
return ms;
}
static Comptador nova(int tipus) {
return switch (tipus) {
case 0 -> new AmbSynchronized();
case 1 -> new AmbAtomicLong();
case 2 -> new AmbLongAdder();
default -> new AmbAtomicReference();
};
}
public static void main(String[] args) throws InterruptedException {
final int TOTAL = 8_000_000;
// Escalfament (JIT).
for (int t = 0; t < 4; t++) mesurar(nova(t), 4, 50_000);
for (int fils : new int[] { 1, 16 }) {
System.out.printf("%n=== %d fil(s), %d increments en total ===%n",
fils, TOTAL);
System.out.printf("%-26s | %8s | %14s%n", "Implementacio", "ms", "inc/ms");
System.out.println("---------------------------|----------|---------------");
for (int t = 0; t < 4; t++) {
Comptador c = nova(t);
long ms = mesurar(c, fils, TOTAL / fils);
System.out.printf("%-26s | %8d | %14d%n",
c.nom(), ms, TOTAL / Math.max(ms, 1));
}
}
// Demostracio de l'avantatge unic d'AtomicReference.
AmbAtomicReference ar = new AmbAtomicReference();
mesurar(ar, 8, 100_000);
AmbAtomicReference.Estat foto = ar.instantania();
System.out.printf("%nInstantania COHERENT: prestecs=%d devolucions=%d multes=%d%n",
foto.prestecs(), foto.devolucions(), foto.multes());
System.out.println("Els tres camps venen del MATEIX instant, cosa que");
System.out.println("tres comptadors independents no poden garantir.");
}
}Sortida orientativa:
=== 1 fil(s), 8000000 increments en total ===
Implementacio | ms | inc/ms
---------------------------|----------|---------------
synchronized | 61 | 131147
AtomicLong | 48 | 166666
LongAdder | 72 | 111111
AtomicReference+record | 284 | 28169
=== 16 fil(s), 8000000 increments en total ===
Implementacio | ms | inc/ms
---------------------------|----------|---------------
synchronized | 1842 | 4343
AtomicLong | 918 | 8714
LongAdder | 94 | 85106
AtomicReference+record | 1531 | 5225
Instantania COHERENT: prestecs=800000 devolucions=0 multes=0
Els tres camps venen del MATEIX instant, cosa que
tres comptadors independents no poden garantir.Anàlisi completa:
- Amb un fil,
AtomicLongguanya iLongAdderperd. Sense contenció, les cel·les múltiples deLongAddersón maquinària que no s'amortitza. Confirma que el millor comptador depèn de la contenció, no n'hi ha cap d'absolut. - Amb 16 fils,
LongAdderés 20× més ràpid queAtomicLongi gairebé 20× més quesynchronized. És exactament el cas per al qual va ser dissenyat. AtomicReference+recordés el més lent en tots dos casos, i té sentit: cada increment crea un objecte nou i el seu bucle CAS reintenta sota contenció. Però mira l'última línia de la sortida: és l'únic que pot donar una instantània coherent dels tres camps. Els altres tres, amb tres comptadors separats, donarien valors d'instants diferents.synchronizedescala pitjor que tots: passa de 61 ms a 1842 ms, un factor de 30, amb la mateixa feina total. És contenció de pany pura, amb canvis de context i suspensió de fils.
Què triar: per a mètriques d'alta freqüència —consultes al catàleg, préstecs per segon—, LongAdder sense dubtar-ho. Per a un estat de negoci que s'ha de llegir de forma coherent —el resum que es mostra a l'usuari o es persisteix—, AtomicReference amb un record immutable, acceptant el seu cost més gran d'escriptura a canvi de lectures gratuïtes i sempre consistents.
Conclusió
Has tancat les dues escletxes que quedaven obertes i has saldat el deute del mòdul 5.
Saps exactament què significa que HashMap no sigui segur per a diversos fils. No és que doni resultats aproximats: perd entrades i pot corrompre la seva estructura fins al punt que un get() entri en un bucle infinit i deixi un nucli al 100 % sense llançar cap excepció. I coneixes la versió concurrent del fail-fast de 05-02, amb el matís que gairebé mai no es diu: la ConcurrentModificationException no està garantida, perquè modCount no és volatile, així que l'alternativa a l'excepció no és que tot vagi bé, sinó que la cursa passi desapercebuda.
Coneixes les tres generacions de solució. Les col·leccions sincronitzades de Collections, que arreglen la corrupció amb un pany global —serialitzant fins i tot les lectures— i porten amb elles el doble parany: iterar continua exigint sincronització manual del recorregut complet, i les operacions compostes —comprovar-després-actuar i llegir-modificar-escriure— continuen sense ser atòmiques, exactament igual que a 08-04. I les col·leccions concurrents, amb ConcurrentHashMap al capdavant: lectures sense cap pany gràcies als camps volatile dels seus nodes, escriptures bloquejades només a la cubeta afectada, i un ordre de magnitud de diferència en un escenari de setze fils.
Tens l'aportació més valuosa de ConcurrentHashMap: les operacions atòmiques compostes —putIfAbsent, computeIfAbsent, compute, merge, getOrDefault, remove(k,v), replace(k,v1,v2)— que són la forma correcta de resoldre comprovar-després-actuar, amb l'advertiment crític que la funció s'executa amb la cubeta bloquejada: res de feina llarga, res de tocar el mateix mapa, res de codi aliè. I coneixes els dos contractes que canvien: els iteradors dèbilment consistents, que no llancen excepció i no bloquegen ningú a canvi de no garantir que vegin els canvis posteriors, i el size() aproximat, que és una dada estadística i mai una base per decidir.
Saps quan fer servir CopyOnWriteArrayList —oients d'esdeveniments, configuració, regles: moltíssimes lectures i gairebé cap escriptura— i per què fer-la servir per acumular en un bucle és O(n²). Saps que ConcurrentLinkedQueue no bloqueja però obliga a sondejar, i que això gairebé sempre és la resposta equivocada.
I tens les BlockingQueue completes, promeses des del mòdul 5: els seus quatre grups d'operacions —excepció, valor especial, blocant, amb termini— i les seves sis implementacions, amb ArrayBlockingQueue acotada com a opció per defecte perquè la cota és el que dona contrapressió: quan la cua s'omple, el productor es bloqueja a put() i passa a produir al ritme que el sistema pot absorbir, en lloc d'acumular fins a exhaurir la memòria. I amb elles, el patró productor-consumidor complet que el mòdul 5 va descriure i va deixar sense escriure, amb take() que espera sense cremar CPU, amb les fallades individuals que no maten el consumidor, i amb la píndola verinosa i les seves tres regles: una per consumidor, al final de la cua —de manera que l'apagada és ordenada per construcció i no es perd ni un element—, i comparada per identitat.
I coneixes la programació sense panys. La instrucció compare-and-swap del processador —«si conté l'esperat, substitueix-lo i digues-me que sí; si no, no toquis res»— i el bucle de reintent que hi ha sota d'incrementAndGet. Amb la diferència essencial que la defineix: un pany diu "que ningú més no toqui això" i els altres esperen bloquejats; un CAS diu "ho intento i, si algú se m'ha avançat, ho repeteixo", i sempre hi ha algú progressant, cosa que fa impossible l'interbloqueig. Domines les operacions d'AtomicInteger, AtomicLong, AtomicBoolean i AtomicReference —inclosos updateAndGet i accumulateAndGet, la funció dels quals ha de ser pura perquè es pot executar diverses vegades—, l'idioma compareAndSet(false, true) que garanteix exactament un guanyador, el problema ABA i la seva solució amb AtomicStampedReference, i LongAdder, que a setze fils supera AtomicLong en un factor de vint i a un sol fil perd: la prova que el millor comptador depèn de la contenció.
Sobretot, tens el patró que combina el millor de dues lliçons: estat immutable amb record + AtomicReference + updateAndGet, que dona lectures gratuïtes, sempre coherents entre tots els camps, i escriptures atòmiques sense cap pany. És la tècnica més elegant de tota la concurrència en Java.
BiblioTech ja no té ni un sol lock() al seu catàleg. CatalegConcurrent fa servir ConcurrentHashMap per als seus índexs, CopyOnWriteArrayList per als seus oients —cosa que compleix automàticament la regla de 08-04 de no cridar codi aliè amb un pany a la mà—, i LongAdder per a les seves mètriques. Un milió sis-centes mil operacions amb setze fils en tres-cents mil·lisegons, amb altes=1000 exactes sobre mil ISBN possibles: el putIfAbsent complint el seu contracte sota màxima contenció. I ProcessadorReserves implementa el productor-consumidor real, amb contrapressió visible a la sortida —la cua estancada a la seva capacitat, frenant els productors— i apagada ordenada sense perdre ni una reserva.
Amb el límit declarat amb honestedat: les col·leccions concurrents garanteixen l'atomicitat d'una operació sobre una col·lecció, no de dues. L'invariant de RegistrePrestecs entre perId i perEmpleat continua necessitant el pany de 08-04, i això no és un defecte de la biblioteca: és la frontera real del que es pot resoldre sense exclusió mútua.
I tens la taula de decisió que ordena tot el mòdul, amb la seva jerarquia: no compartir, compartir immutable, atòmic, col·lecció concurrent, pany — en aquest ordre, baixant només quan el nivell anterior no hi arriba.
Però fixa't en el que continua faltant. BiblioTech ja fa coses en paral·lel, però cada vegada que necessita el resultat d'alguna cosa, algú es queda esperant: future.get() bloqueja. Si volguessis encadenar tres passos —consultar el catàleg, calcular les multes i exportar l'informe—, hauries de fer get() entre cadascun, i el fil que orquestra passaria gairebé tot el temps aturat. Future et diu "aquí arribarà un resultat", però l'única forma de fer-lo servir és preguntar i esperar. No hi ha manera de dir "quan estigui llest, fes això altre", ni de combinar dos resultats independents, ni de definir què fer si alguna cosa falla enmig de la cadena.
A la propera lliçó, Tasques Asíncrones amb CompletableFuture, es tanca el mòdul amb la resposta a això. Veuràs les quatre limitacions de Future que van motivar el seu successor; el model de composició asíncrona, en què declares la cadena sencera per endavant i cap fil no espera ningú; la creació amb supplyAsync i runAsync i per què convé passar el teu propi Executor en lloc de fer servir el pool comú; la diferència entre thenApply i thenCompose —amb el CompletableFuture<CompletableFuture<T>> que apareix quan t'equivoques—; la combinació amb thenCombine, allOf i anyOf; la gestió d'errors amb exceptionally, handle i whenComplete, amb la CompletionException que embolcalla la causa; els terminis d'orTimeout i completeOnTimeout; i els paranys que fan que una cadena asíncrona mal escrita sigui pitjor que el codi blocant que substitueix. En acabar, BiblioTech tindrà una cadena que consulta, calcula i exporta sense bloquejar el menú en cap moment, i el mòdul quedarà tancat.
Curs de Programació en Java
Mòdul 1: Introducció a Java
- Introducció a Java
- Configuració de l'entorn de desenvolupament
- Sintaxi i estructura bàsica
- Variables i tipus de dades
- Operadors
- Entrada i sortida per consola
- El teu primer programa complet: BiblioTech
Mòdul 2: Flux de control
- Sentències condicionals
- Bucles
- Sentències switch
- Break i continue
- Depuració i traces d'execució
- Projecte: menú interactiu de BiblioTech
Mòdul 3: Programació orientada a objectes
- Introducció a la POO
- Classes i objectes
- Mètodes
- Constructors
- Herència
- Polimorfisme
- Encapsulament
- Abstracció
- La classe Object: equals, hashCode i toString
Mòdul 4: Programació orientada a objectes avançada
- Interfícies
- Classes abstractes
- Classes internes
- Classes anònimes
- Expressions lambda
- Interfícies funcionals i referències a mètodes
- Enumeracions i registres
Mòdul 5: Estructures de dades i col·leccions
- Arrays
- El framework de col·leccions
- ArrayList
- LinkedList
- HashMap
- HashSet
- Cua i Deque
- Pila
- Ordenació i cerca en col·leccions
Mòdul 6: Gestió d'excepcions
- Introducció a les excepcions
- Bloc try-catch
- Throw i throws
- Excepcions personalitzades
- Bloc finally
- Try-with-resources i AutoCloseable
- Estratègies de gestió d'errors i logging
Mòdul 7: Entrada/sortida de fitxers
- Lectura de fitxers
- Escriptura de fitxers
- Fluxos de fitxers
- BufferedReader i BufferedWriter
- Serialització
- L'API NIO.2: Path i Files
- Formats d'intercanvi: CSV i Properties
Mòdul 8: Multifil i concurrència
- Introducció al multifil
- Creació de fils
- Cicle de vida d'un fil
- Sincronització
- Utilitats de concurrència
- Col·leccions concurrents i variables atòmiques
- Tasques asíncrones amb CompletableFuture
Mòdul 9: Xarxes
- Introducció a les xarxes
- Sockets
- ServerSocket
- DatagramSocket i DatagramPacket
- URL i HttpURLConnection
- El client HTTP modern
Mòdul 10: Temes avançats
- Genèrics
- Anotacions
- Reflexió
- Característiques de Java 8: Streams i Optional
- Dates i hores amb java.time
- Java 9 i més enllà
- Memòria, recol·lecció de brossa i rendiment
Mòdul 11: Frameworks i llibreries de Java
- Introducció als frameworks de Java
- Spring Framework
- Hibernate
- JUnit
- Maven
- Proves avançades amb Mockito
- Llibreries essencials de l'ecosistema
