Cours 6 — Pipelines, WaitGroup, Mutex et Deadlock
Pipeline
Qu'est-ce qu'un pipeline? Ça sert à prendre la sortie d'une méthode et l'envoyer dans l'entrée d'une autre:
Ainsi on pourrait schématiser ceci:
Compteur ------> AuCarre -------> Afficheur
Comment réussir ça sous go? En utilisant deux channels:
package main
import (
"fmt"
)
func main() {
nombres := make(chan int)
carres := make(chan int)
// Compteur (déclarée anonyme et inline)
go func() {
for x := 0; x < 100; x++ {
nombres <- x // Envoi
}
}()
// Carré (déclarée anonyme et inline)
go func() {
for {
x := <- nombres // Réception
carres <- x * x // Envoi
}
}()
// Afficheur (goroutine main)
for {
fmt.Println(<-carres) // Réception
}
}
Il est recommandé de fermer les canaux lorsque nous savons que nous avons terminé de communiquer, n'oubliez pas que c'est synchrone, donc le close ne se produira pas tant que l'exécution ne sera pas terminée avec des boucles.
Il est possible de fermer les channels dès que l'émetteur n'en a plus besoin. Un simple appel à close(chan) permettra de le fermer. Ainsi, il est nécessaire de tester afin de s'assurer que le channel est toujours ouvert avant d'assigner la valeur:
x, ok := <- nombres
if !ok {
break // channel was closed and drained
}
carres <- x * x
Ceci est très laid et surtout horrible. Une autre fonctionnalité intéressante est l'utilisation du range plutôt que de tester à chaque fois:
for x := range nombres {
carres <- x * x
}
Piège fréquent en programmation réseau : channel jamais fermé dans un pipeline. Un range sur un channel ne se termine que lorsque ce channel est fermé — dans un pipeline de traitement de messages réseau (lecture → parsing → traitement), oublier de fermer le channel intermédiaire bloque silencieusement toute la chaîne en aval, indéfiniment :
func lireMessages(conn net.Conn, messages chan<- string) {
scanner := bufio.NewScanner(conn)
for scanner.Scan() {
messages <- scanner.Text()
}
// BUG : pas de close(messages) ici — quand la connexion se termine (scanner.Scan()
// retourne false), la goroutine ci-dessous reste bloquée pour toujours sur son "range"
}
func traiterMessages(messages <-chan string) {
for m := range messages { // ne sortira jamais de cette boucle si "messages" n'est jamais fermé
fmt.Println("reçu:", m)
}
fmt.Println("terminé") // n'est jamais atteint
}
La correction : fermer le channel de sortie dès que la source n'a plus rien à envoyer, typiquement avec defer close(messages) au début de lireMessages :
func lireMessages(conn net.Conn, messages chan<- string) {
defer close(messages) // garantit la fermeture même si la boucle se termine par une erreur
scanner := bufio.NewScanner(conn)
for scanner.Scan() {
messages <- scanner.Text()
}
}
Il est possible de passer les channels en paramètres afin d'obtenir un affichage plus consistant. On peut aussi identifier le sens du channel: in ou out:
package main
import (
"fmt"
)
// Compteur (déclarée anonyme et inline)
func Compteur(out chan<- int) {
for x := 0; x < 100; x++ {
out <- x
}
defer close(out)
}
// Carre (déclarée anonyme et inline)
func Carre(out chan<- int, in <-chan int) {
for x := range in {
out <- x * x
}
defer close(out)
}
// Afficheur (in main goroutine)
func Afficheur(in <-chan int) {
for x := range in {
fmt.Println(x)
}
}
func main() {
nombres := make(chan int)
carres := make(chan int)
go Compteur(nombres)
go Carre(carres, nombres)
Afficheur(carres)
}
Attendre la fin des goroutines (WaitGroup)
Un waitgroup est une primitive de synchronisation. Elle permet de forcer go à attendre un moment avant de continuer le code (donc synchronisation).
Au lieu de faire un fmt.Scan à la fin du main ou de faire un compteur par channel, on peut par programmation attendre que les goroutines terminent. Vous avez besoin du import "sync". Voici les lignes de codes importantes:
var wg sync.WaitGroup // Objet important, il faudra le passer en pointeur aux goroutines
wg.Add(1) // A appeler avant de partir la goroutine, ajoute '1' au compteur du WaitGroup
wg.Done() // A appeler à la fin de la goroutine, retire 1 au compteur du WaitGroup
wg.Wait() // Bloque jusqu'à ce que le compteur arrive à zéro
Exemple sans WaitGroup (waitGroup-1) — le main peut se terminer avant même que la goroutine ait eu la chance de s'exécuter.
Exemple avec WaitGroup (waitGroup-2) — corrige le problème ci-dessus en attendant explicitement la fin de la goroutine.
Quand l'utiliser
Pour un arrêt propre (graceful shutdown) d'un serveur — attendre que toutes les goroutines encore en train de traiter une connexion cliente terminent avant de couper le programme, plutôt que de les interrompre brutalement en pleine écriture :
var wg sync.WaitGroup
for {
conn, err := listener.Accept()
if err != nil {
break // ex.: listener.Close() appelé ailleurs pour déclencher l'arrêt
}
wg.Add(1)
go func(c net.Conn) {
defer wg.Done()
defer c.Close()
traiterClient(c)
}(conn)
}
wg.Wait() // attend que toutes les connexions en cours terminent proprement
fmt.Println("serveur arrêté")
Quand l'éviter
Si vous devez récupérer une valeur de retour ou une erreur individuelle de chaque goroutine — un WaitGroup ne fait qu'attendre, il ne transporte aucune donnée ; combinez-le avec un channel de résultats si vous devez agréger des réponses.
Piège fréquent en programmation réseau : envoi sur un channel déjà fermé. Contrairement à la lecture (qui retourne simplement la valeur zéro une fois le channel fermé), un envoi sur un channel fermé provoque un panic: send on closed channel. Le scénario classique : deux goroutines qui gèrent chacune une connexion cliente et qui pensent chacune être responsables de fermer un channel notifications partagé une fois « leur » client déconnecté :
notifications := make(chan string)
go func() { // goroutine A : gère client 1
defer close(notifications) // BUG : ferme un channel PARTAGÉ avec la goroutine B
// ... traiter client 1, envoyer des notifications ...
}()
go func() { // goroutine B : gère client 2
defer close(notifications) // BUG : close() d'un channel déjà fermé → panique
notifications <- "client 2 déconnecté" // et un envoi après l'autre close() panique aussi
}()
La règle à retenir : un seul propriétaire ferme un channel, jamais plusieurs émetteurs concurrents. Ici, la correction consiste à donner à chaque client son propre channel de notifications, ou à confier la fermeture à une troisième goroutine coordinatrice qui sait quand tous les clients sont partis — par exemple via le WaitGroup vu ci-dessus :
var wg sync.WaitGroup
notifications := make(chan string)
wg.Add(2)
go func() { defer wg.Done(); /* ... traiter client 1 ... */ }()
go func() { defer wg.Done(); /* ... traiter client 2 ... */ }()
go func() {
wg.Wait() // attend que les DEUX goroutines aient fini d'envoyer
close(notifications) // un seul endroit ferme le channel, une seule fois
}()
Interblocage : deadlock
Si vous oubliez de libérer un mutex ou si vos channels attendent des réponses sans rien recevoir, vous avez un interblocage. Un interblocage survient lorsque vos threads ou goroutines attendent tous après une ressource, mais que rien ne se passe.

package main
func main() {
println("Channels are in Deadlock")
Channel := make(chan int)
Channel2 := make(chan int)
<- Channel
<- Channel2
}
Concurrence critique (Race condition)
La concurrence critique est un problème qui arrive lorsque plusieurs goroutines accèdent et modifient une même variable, par exemple:
func main() {
var nb int = 0
for i := 0; i < 1000; i++ {
go func() {
nb++
}(nb)
}
println(nb)
}
Ce magnifique code affiche un nombre entre 900 et 950. Pourquoi est-ce que je n'obtiens pas 1000? L'instruction '++' se fait en 3 étapes: lire la valeur, l'incrémenter et l'affecter. Si 2 goroutines s'exécutent en même temps, nous aurons:
- go1: Lis la valeur de nb : 20
- go1: Incrémente la valeur en mémoire : 21
- go2: Lis la valeur de nb : 20
- go1: Affecte la nouvelle valeur de n : 21
- go2: Incrémente la valeur en mémoire : 21
- go2: Affecte la nouvelle valeur de n : 21
La solution que nous avons utilisée à présent est de limiter la manipulation de la variable à une seule goroutine, en utilisant un channel. Ça fonctionne mais il existe des outils plus adaptés.
Solutions
Opération atomique (Atomic Operations)
Dans le paquet "sync/atomic", l'opération atomique permet de modifier une variable (nombre entier) sans problème de concurrence critique. Par exemple, dans l'exemple précédent au lieu de "n++" je vais mettre:
func main() {
var nb int32 = 0
wg := sync.WaitGroup{}
for i := 0; i < 1000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
atomic.AddInt32(&nb, 1)
}()
}
wg.Wait()
println(nb)
}
Avec cette modification ma boucle va me donner le bon chiffre.
Les opérations atomiques ne fonctionnent que sur les types: int32, int64, uint32, uint64. L'opération atomique utilise une fonction spéciale du processeur: "test-and-set". Le processeur va donc faire les 3 étapes (lecture, modification et sauvegarde) en 1 opération. Cette opération est plus lente qu'une simple affectation.
J'ai roulé un benchmark sur une simple boucle qui utilise "n++" ou "atomic.AddInt32(&n, 1)". L'affectation simple a pris 8 μs (microseconde, 1000 nanosecondes) tandis que l'opération atomique a pris 266 μs, c'est 32 fois plus lent.
Il y a également des méthodes de Load / Store, CompareAndSwap (CAS) et Swap.
Mutex (Mutual Exclusion)
Dans le paquet "sync", un mutex est un verrou. Une seule goroutine à la fois peut accéder au verrou. C'est comme une toilette, quand la porte est verrouillée tu attends ton tour...
var n int = 0
var m sync.Mutex
for i := 0; i < 1000; i++ {
go func(nb *int, m *sync.Mutex) {
m.Lock()
nb++
m.Unlock()
}(&n, &m)
}
Le mutex utilise l'opération atomique de processeur "test-and-set" (comme les opérations atomiques). Pour faire simple, le Mutex crée un booléen (barré ou non) et utilise une opération atomique pour mettre ou enlever le verrou.
L'opération atomique n'effectue qu'une seule opération tandis qu'avec un Mutex nous pouvons faire aussi compliqué que nous le désirons. Le Mutex est bloquant, utilisez-le seulement pour le code qui en a besoin.
Côté performance, le Mutex est encore plus lent que l'opération atomique pour une addition simple. L'opération atomique prenait environ 266 μs, tandis que le mutex prend 354 μs.
Le Mutex est disponible dans tous les langages, c'est quelque chose de standard. Dans la plupart des langages, le mutex peut seulement être libéré par le Thread qui a obtenu le verrou.
Lien avec le deadlock : un Lock() sans Unlock() correspondant (oubli, return prématuré, panique) est l'une des causes les plus fréquentes d'interblocage — d'où l'habitude de systématiquement écrire defer m.Unlock() juste après m.Lock().
Quand l'utiliser
Dès qu'une donnée est lue et modifiée par plusieurs goroutines en même temps — typiquement la table des sessions clientes d'un serveur (voir le piège concurrent map read and map write du Cours 2) :
type Serveur struct {
mu sync.Mutex
clients map[string]net.Conn
}
func (s *Serveur) Ajouter(id string, conn net.Conn) {
s.mu.Lock()
defer s.mu.Unlock()
s.clients[id] = conn
}
Sans ce Mutex, deux goroutines qui traitent deux connexions en parallèle et modifient s.clients en même temps provoquent le crash fatal error: concurrent map read and map write.
Piège fréquent en programmation réseau : oublier de déverrouiller. Le danger n'est pas d'oublier le Unlock() dans le cas simple ci-dessus, mais dès qu'un chemin de sortie anticipé apparaît dans la méthode — une erreur réseau, par exemple :
func (s *Serveur) Deconnecter(id string) error {
s.mu.Lock()
conn, ok := s.clients[id]
if !ok {
return errors.New("client inconnu") // BUG : s.mu reste verrouillé pour toujours
}
conn.Close()
delete(s.clients, id)
s.mu.Unlock()
return nil
}
Le prochain appel à Deconnecter (ou à toute autre méthode qui prend s.mu) bloque indéfiniment : deadlock. La correction consiste à defer le déverrouillage immédiatement après le verrouillage, avant même d'écrire la suite de la fonction :
func (s *Serveur) Deconnecter(id string) error {
s.mu.Lock()
defer s.mu.Unlock() // s'exécute peu importe le chemin de sortie
conn, ok := s.clients[id]
if !ok {
return errors.New("client inconnu")
}
conn.Close()
delete(s.clients, id)
return nil
}
Quand l'éviter
Si une seule goroutine à la fois accède à la donnée (ex.: elle vit uniquement à l'intérieur de la goroutine qui gère une connexion et n'est jamais partagée) — ajouter un Mutex qui ne protège rien alourdit le code sans bénéfice.
RWMutex
Le RWMutex est un type particulier de Mutex en Go. Il contient 2 types de verrous, un verrou "booléen" en écriture et un verrou "numérique" en lecture. Plusieurs goroutines peuvent accéder au verrou en lecture en même temps, lorsqu'on obtient un verrou en lecture Go incrémente le nombre de lecteurs de manière atomique. Pour obtenir le verrou en écriture, le nombre de lecteurs doit être à zéro. Lorsqu'il y a un verrou en écriture, aucune autre goroutine ne peut accéder au mutex en lecture ou écriture jusqu'à ce que l'écrivain le libère.
package main
import (
"fmt"
"sync"
"time"
)
func main() {
var m sync.RWMutex
var data int
var wg sync.WaitGroup
// 1 écriture
wg.Add(1)
go func() {
defer wg.Done()
m.Lock()
data = 42
time.Sleep(50 * time.Millisecond) // simulons un calcul long
m.Unlock()
}()
// 100 lectures
for i := 0; i < 100; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
m.RLock()
fmt.Printf("Goroutine lecture %d: %d\n", id, data)
time.Sleep(10 * time.Millisecond) // lecture "lente"
m.RUnlock()
}(i)
}
wg.Wait()
}
Les avantages du RWMutex :
- Lecture simultanée : toutes les 100 goroutines peuvent lire en même temps (tant qu'aucune écriture n'est en cours).
- Écriture rare : seule la goroutine d'écriture bloque tout le monde pendant son calcul.
- Performance : si on remplace RWMutex par un simple Mutex, alors chaque lecture bloquerait les autres, ce qui ralentirait énormément le programme.
Quand l'utiliser
Quand la donnée protégée est lue beaucoup plus souvent qu'elle n'est modifiée — un cache de configuration rechargé occasionnellement (ex.: toutes les 30 secondes depuis un fichier) mais consulté à chaque requête entrante :
type Config struct {
mu sync.RWMutex
valeurs map[string]string
}
func (c *Config) Lire(cle string) string {
c.mu.RLock()
defer c.mu.RUnlock()
return c.valeurs[cle]
}
func (c *Config) Recharger(nouvelles map[string]string) {
c.mu.Lock()
defer c.mu.Unlock()
c.valeurs = nouvelles
}
Quand l'éviter
Si les lectures et les écritures sont à peu près aussi fréquentes, ou si la section critique est très courte — la gestion plus complexe du RWMutex (deux compteurs internes au lieu d'un simple booléen) le rend alors plus lent qu'un Mutex ordinaire, dans le même esprit que la comparaison Mutex/atomique plus haut.
Piège fréquent en programmation réseau : copier une struct qui contient un sync.Mutex. Passer une struct de session par valeur plutôt que par pointeur copie aussi son mutex — chaque copie protège alors sa propre copie du verrou, pas la donnée partagée :
type Session struct {
mu sync.Mutex
compte int
}
// BUG : Session est reçue par valeur, donc "s" à l'intérieur de Incrementer
// est une copie complète, avec son propre mutex indépendant de l'original
func Incrementer(s Session) {
s.mu.Lock()
s.compte++
s.mu.Unlock()
}
Le verrou pris dans Incrementer ne protège jamais la vraie Session partagée entre goroutines : la race condition persiste silencieusement, sans même que go run -race s'en étonne côté logique métier. La correction consiste à toujours passer (et stocker) ce genre de struct par pointeur :
func Incrementer(s *Session) {
s.mu.Lock()
s.compte++
s.mu.Unlock()
}
go vet détecte ce cas précis (un sync.Mutex copié par valeur) et le signale automatiquement — un bon réflexe est de lancer go vet ./... régulièrement.
Détecter les races avec l'outil Go
Plutôt que de repérer les races à l'œil, Go fournit un détecteur intégré qui instrumente les accès mémoire concurrents à l'exécution :
go run -race main.go
Piège fréquent : partager la variable de boucle directement avec une goroutine, plutôt que de la lui passer en paramètre :
// ❌ i est partagée : sa valeur au moment où la goroutine s'exécute n'est pas garantie
for i := 0; i < n; i++ {
go func() { fmt.Println(i) }()
}
// ✅ i est copiée dans j au moment du lancement
for i := 0; i < n; i++ {
go func(j int) { fmt.Println(j) }(i)
}
À retenir
- Toujours protéger les écritures concurrentes (
sync.Mutex,sync.RWMutexousync/atomic). - Ne jamais partager une variable de boucle directement dans une goroutine — la passer en paramètre.
- Utiliser
go run -racependant le développement pour détecter automatiquement les conditions de course. - Toujours
defer m.Unlock()juste aprèsm.Lock()pour éviter les deadlocks par oubli.