Cours 8 — Channel avancé et condition de course
Channel de channel et pool de workers
En Go, les channels permettent la communication entre goroutines — ils servent aussi de brique de base pour organiser un pool de workers, une façon de traiter un grand nombre de tâches avec un nombre limité de goroutines actives. Ce chapitre part de l'implémentation la plus simple, puis introduit une variante plus avancée (le channel de channel) une fois que sa nécessité devient claire.
Pool de workers simple
On peut utiliser deux channels (jobs / results) pour envoyer des jobs à un pool de workers :
package main
import (
"fmt"
"time"
)
func worker(id int, jobs <-chan int, results chan<- int) {
for j := range jobs {
fmt.Printf("Worker %d traite le job %d\n", id, j)
time.Sleep(time.Second) // simuler le traitement
results <- j * 2
}
}
func main() {
const numJobs = 5
const numWorkers = 3
jobs := make(chan int, numJobs)
results := make(chan int, numJobs)
// Lancement des workers
for w := 1; w <= numWorkers; w++ {
go worker(w, jobs, results)
}
// Envoi des jobs
for j := 1; j <= numJobs; j++ {
jobs <- j
}
close(jobs)
// Récupération des résultats
for a := 1; a <= numJobs; a++ {
fmt.Println("Résultat :", <-results)
}
}
Chaque worker reçoit ses jobs via le channel jobs et renvoie ses résultats sur results. Si on voulait un pool dynamique de workers, on pourrait envoyer les channels des workers eux-mêmes sur un channel de channel pour gérer leur disponibilité.
Quand l'utiliser
Un pool de workers borné est la façon idiomatique de traiter un grand nombre de connexions clientes avec un nombre limité de goroutines actives — utile dès que le traitement de chaque connexion coûte cher (calcul, accès à une base de données, appel à une API externe) et qu'on ne peut pas se permettre une goroutine illimitée par connexion sous forte charge, contrairement au patron « une goroutine par connexion » vu au cours 4.
Quand l'éviter
Le pool simple ne convient pas quand il faut cibler un worker précis plutôt que n'importe lequel de disponible — par exemple pour garder un état propre à une connexion sur le même worker d'un appel à l'autre (« session affinity »), ou pour connaître à tout moment lequel des workers est réellement occupé. Avec jobs := make(chan int, numJobs), les workers sont interchangeables et anonymes : Go choisit lui-même, par compétition sur la réception (for j := range jobs), quel worker traite quel job, sans qu'aucun code n'ait explicitement à faire ce choix.
Le pool simple ne garantit pas non plus que l'ordre des résultats sur results corresponde à l'ordre d'envoi des jobs sur jobs — chaque worker traite à son propre rythme. Si un job et son résultat doivent être associés de façon fiable, il faut inclure cette correspondance dans la donnée elle-même (par exemple un champ id dans le job et dans le résultat), plutôt que de se fier à l'ordre de réception.
Pool simple vs pool avec channel de channel — lequel choisir ?
- Pool simple : suffisant dans la grande majorité des cas — dès que les workers sont interchangeables et qu'on n'a besoin ni de cibler un worker en particulier, ni de savoir explicitement qui est occupé. C'est le cas le plus fréquent en programmation réseau : traiter un flot de requêtes similaires sans affinité de session.
- Pool avec channel de channel : utile seulement quand cette adressabilité individuelle devient réellement nécessaire — routage nominatif vers un worker précis, ou pattern requête/réponse. Voir les deux prochaines sections (« Pool dynamique avec channel de channel » et « Pourquoi un channel de channel plutôt qu'un identifiant de session ? ») pour le mécanisme en détail.
En pratique : commencer par le pool simple, qui couvre la grande majorité des besoins de limitation de concurrence en Go, et ne passer au channel de channel que si ce besoin d'adressabilité se manifeste réellement.
Piège fréquent en programmation réseau : file de jobs sous-dimensionnée qui bloque l'acceptation de nouvelles connexions. Si le channel jobs est non tamponné (ou trop petit) et que tous les workers sont occupés, la boucle qui accepte les connexions se bloque elle-même en essayant d'y déposer un job — le serveur cesse d'accepter de nouveaux clients pendant que d'autres attendent déjà dans la file TCP du système d'exploitation :
jobs := make(chan net.Conn) // BUG : non tamponné, aucune marge
for w := 1; w <= numWorkers; w++ {
go worker(w, jobs)
}
for {
conn, err := listener.Accept()
if err != nil {
continue
}
jobs <- conn // BUG : bloque ICI dès que tous les workers sont occupés,
} // et donc bloque aussi Accept() pour la connexion suivante
La correction consiste à dimensionner le tampon selon la charge attendue (une file d'attente raisonnable, pas illimitée) et à décider explicitement quoi faire quand elle est pleine plutôt que de bloquer silencieusement :
jobs := make(chan net.Conn, 100) // marge pour absorber les pics de connexions
for {
conn, err := listener.Accept()
if err != nil {
continue
}
select {
case jobs <- conn:
// job accepté normalement
default:
conn.Close() // file pleine : refuser plutôt que de bloquer Accept()
}
}
Pool dynamique avec channel de channel
Le pool simple a une limite structurelle : personne ne sait explicitement quel worker est libre à un instant donné — Go arbitre ça implicitement via la compétition sur jobs. Pour lever cette limite, il faut rendre chaque worker adressable individuellement, ce qui demande de faire circuler son channel comme une valeur à part entière. Il est possible de créer un channel de channels, c'est-à-dire un channel dont les éléments sont eux-mêmes des channels :
c1 := make(chan chan int) // c1 transporte des channels de type chan int
Dans un pool dynamique, chaque worker crée son propre channel de jobs et le dépose sur un channel workerPool de ce type dès qu'il est disponible :
package main
import (
"fmt"
"sync"
)
func main() {
var wg sync.WaitGroup
workerPool := make(chan chan int, 3)
// Création de workers
for i := 0; i < 3; i++ {
w := make(chan int)
workerPool <- w
go func(id int, jobs chan int, wg *sync.WaitGroup) {
for job := range jobs {
fmt.Printf("Worker %d traite le job %d\n", id, job)
wg.Done()
}
}(i+1, w, &wg)
}
// Distribution des jobs
for j := 1; j <= 25; j++ {
wg.Add(1)
w := <-workerPool // récupérer un worker disponible
w <- j // envoyer le job
workerPool <- w // remettre le worker dans le pool
}
// Fermer tous les workers
for i := 0; i < 2; i++ {
w := <-workerPool
close(w)
}
wg.Wait()
}
workerPoolest un channel de channel, qui contient les channels de chaque worker disponible.- Le dispatcher récupère un worker libre (
<-workerPool), lui envoie un job (w <- j), puis le remet dans le pool une fois libre (workerPool <- w). - Cela garantit qu'aucun worker ne reçoit deux jobs à la fois, tout en distribuant dynamiquement le travail selon les disponibilités.
Exemple complet : https://gitlab.com/drynish/ecole/-/tree/main/420-M54/Exemples/Cours%208/2.%20Workers
Pourquoi un channel de channel plutôt qu'un identifiant de session ?
Le pool dynamique ci-dessus illustre un premier usage du channel de channel : faire circuler le channel d'un worker à travers workerPool pour signaler sa disponibilité. Il existe un second usage, structurellement différent : envoyer un channel de réponse comme partie d'une requête, pour qu'un résultat revienne directement au bon demandeur — sans registre de correspondance à maintenir. C'est ce second patron qu'on détaille ici.
Une alternative courante pour faire correspondre une réponse à sa requête consiste à attribuer un identifiant de session à chaque requête, puis à répondre sur un channel partagé en y attachant cet identifiant — charge ensuite au demandeur de filtrer les réponses pour retrouver la sienne. Ça fonctionne, mais ça force à maintenir un état partagé (typiquement une map[int]chan Resultat protégée par un sync.Mutex) pour faire le lien entre chaque ID et son demandeur, avec le risque d'oublier de retirer une entrée de la map une fois la réponse livrée (fuite mémoire) ou de générer deux fois le même ID par erreur.
Le channel de channel élimine ce problème : au lieu d'envoyer un identifiant, on envoie directement le channel de réponse avec la requête. La réponse revient alors uniquement à qui l'attend, sans table de correspondance ni verrou à gérer :
type requete struct {
donnees string
reponse chan string // channel de retour privé, propre à CETTE requête
}
func serveur(requetes chan requete) {
for r := range requetes {
r.reponse <- fmt.Sprintf("traité: %s", r.donnees)
}
}
func client(requetes chan requete, donnees string) string {
reponse := make(chan string) // créé pour cette requête seulement
requetes <- requete{donnees: donnees, reponse: reponse}
return <-reponse // reçoit forcément SA réponse, rien à filtrer
}
Chaque appel à client crée son propre channel reponse, l'envoie comme partie de la requête, puis attend dessus. Le serveur n'a besoin d'aucun état pour savoir « à qui répondre » : il répond simplement sur le channel qu'on lui a fourni. Une fois la réponse reçue, ce channel n'est plus référencé nulle part et le garbage collector s'en occupe — contrairement à l'entrée de map qu'il faudrait explicitement supprimer avec l'approche par identifiant.
Exemple complet, avec plusieurs clients concurrents qui reçoivent chacun uniquement leur propre réponse : https://gitlab.com/drynish/ecole/-/tree/main/420-M54/Exemples/Cours%208/7.%20Requete%20Reponse
Limite importante : cette technique ne fonctionne qu'à l'intérieur d'un même processus Go — et pas seulement « pas sur le réseau ». Même deux processus Go distincts tournant côte à côte sur la même machine ne peuvent pas se partager un channel : chacun a son propre espace mémoire et son propre runtime/scheduler Go, complètement indépendants l'un de l'autre.
La raison est plus profonde qu'une simple histoire de distance physique : un channel n'est pas vraiment « de la donnée », contrairement à un int ou une string. C'est une référence vers une structure interne du runtime Go — la file d'attente des goroutines bloquées en envoi/réception, le buffer circulaire, le drapeau « fermé », etc. — qui n'a de sens que pour le scheduler du processus qui l'a créée. Il n'y a donc rien de cohérent à transmettre : même en essayant avec encoding/json ou encoding/gob, l'encodage échouerait, puisque chan n'est simplement pas un type que ces encodeurs savent sérialiser.
Le réseau n'est qu'un cas particulier de cette limite plus générale : un client et un serveur réseau sont presque toujours deux processus séparés — souvent sur deux machines, mais le problème existerait identiquement pour deux processus sur la même machine. Un vrai protocole client-serveur (HTTP, un protocole binaire personnalisé, etc.) doit donc forcément revenir à un identifiant de requête — contrairement au channel, un ID (int, string, UUID) est une vraie donnée sérialisable qu'on peut transmettre sur le fil et retrouver de l'autre côté. Le channel de channel — sous ses deux formes vues dans ce chapitre, disponibilité d'un worker ou adresse de réponse — reste donc une technique strictement interne à un processus Go, utile pour coordonner des goroutines entre elles, jamais pour communiquer entre deux programmes distincts.
Nil channels
Un nil channel est un type particulier de channel qui bloque systématiquement toute opération. Il est souvent utilisé pour désactiver dynamiquement certaines branches d'un select.
Quand l'utiliser
Pour désactiver dynamiquement un canal de contrôle sans complexifier la logique du select avec des if supplémentaires — par exemple, un serveur qui accepte des commandes d'administration sur un channel tant qu'il fonctionne normalement, et qui met ce channel à nil pendant une opération de maintenance pour ignorer temporairement les nouvelles commandes sans fermer la connexion d'administration.
-
Envoi et réception sur un nil channel bloquent indéfiniment :
var ch chan int // nilch <- 1 // bloque indéfiniment<-ch // bloque indéfiniment -
Dans un
select, les cases qui utilisent un nil channel ne peuvent jamais être sélectionnées — ça permet de désactiver temporairement une case sans supprimer le code.
Exemple : désactivation dynamique d'une case dans un select
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int)
var stop chan int // nil channel
go func() {
time.Sleep(2 * time.Second)
fmt.Println("Activation du channel stop")
stop = make(chan int) // channel non-nil
}()
for i := 0; i < 5; i++ {
select {
case ch <- i:
fmt.Println("Envoyé :", i)
case <-stop:
fmt.Println("Stop reçu")
return
default:
fmt.Println("Aucune case disponible")
}
time.Sleep(time.Second)
}
}
Initialement, stop est un nil channel : la case case <-stop: est désactivée. Après 2 secondes, on lui assigne un channel réel : la case devient active. Ça permet de contrôler dynamiquement quelles branches du select sont disponibles.
Points clés : un envoi/réception sur un nil channel bloque toujours ; dans un select, un nil channel n'est jamais choisi ; utile pour désactiver temporairement une communication.
Piège fréquent en programmation réseau : un default qui transforme un envoi voulu bloquant en perte silencieuse de message. Dans l'exemple ci-dessus, case ch <- i: n'envoie i que si un récepteur est immédiatement prêt ; sinon, le default s'exécute et la valeur i est perdue sans avertissement — c'est exactement ce qui se produirait avec de vrais messages réseau si ch représentait une connexion cliente qu'on essaie d'alimenter « seulement si le client est prêt » :
select {
case sortie <- message: // BUG potentiel : si le lecteur de "sortie" (le client) est
// envoyé momentanément occupé, on tombe dans le default ci-dessous
default:
// le message est silencieusement perdu, sans log ni erreur
}
Si la perte de message n'est PAS acceptable — ce qui est le cas pour la plupart des données applicatives, par opposition à un simple signal de contrôle comme dans l'exemple du nil channel ci-dessus — retirez le default pour forcer l'attente, ou journalisez explicitement la perte pour qu'elle reste visible :
select {
case sortie <- message:
case <-time.After(time.Second): // attend au plus 1 seconde avant d'abandonner
log.Printf("message perdu, client trop lent: %v", message)
}
Ordre des goroutines
Même si on ne doit normalement faire aucune hypothèse sur l'ordre d'exécution des goroutines, il est possible d'imposer un ordre déterminé en utilisant des channels comme mécanisme de synchronisation : chaque étape attend un signal (valeur ou fermeture) sur un channel avant de démarrer, puis envoie un signal sur le channel suivant. Cela crée une chaîne de dépendances A → B → C.
package main
import (
"fmt"
"time"
)
func A(start, next chan struct{}) {
<-start
fmt.Println("A()!")
time.Sleep(time.Second)
close(next) // débloque B
}
func B(start, next chan struct{}) {
<-start
fmt.Println("B()!")
time.Sleep(time.Second)
close(next) // débloque C
}
func C(start chan struct{}) {
<-start
fmt.Println("C()!")
}
func main() {
chA := make(chan struct{})
chB := make(chan struct{})
chC := make(chan struct{})
go A(chA, chB)
go B(chB, chC)
go C(chC)
// Lancement initial
close(chA)
time.Sleep(3 * time.Second)
}
A ferme chA → B démarre ; B ferme chB → C démarre : l'ordre d'exécution est forcé par la synchronisation via les channels.
Point à retenir : ne jamais réutiliser ou fermer un channel plus d'une fois (ça cause une panique) — ces channels ne servent ici que de signaux, pas de transport de données.
Fuite de goroutine : channel done
Au cours 4, on a vu qu'une goroutine bloquée pour toujours sur un channel qui ne recevra jamais rien fuit — elle n'est jamais récupérée par le ramasse-miettes, et le programme accumule des goroutines fantômes jusqu'à épuiser la mémoire :
func gererClient(conn net.Conn, notifications chan string) {
go func() {
for msg := range notifications { // BUG : si "notifications" n'est jamais fermé ni
fmt.Fprintln(conn, msg) // alimenté après la fin de gererClient, cette
} // goroutine reste bloquée pour toujours ici
}()
// ... traitement de conn, qui peut se terminer (et fermer conn) sans jamais fermer "notifications" ...
}
Maintenant qu'on a vu select et les channels chan struct{} comme signal (ci-dessus), on peut corriger ça : donner à la goroutine un moyen de se terminer même sans nouveau message, via un channel done dédié vérifié dans un select (ou un context.Context) :
func gererClient(conn net.Conn, notifications chan string, done <-chan struct{}) {
go func() {
for {
select {
case msg, ok := <-notifications:
if !ok {
return
}
fmt.Fprintln(conn, msg)
case <-done: // signalé quand gererClient se termine
return
}
}
}()
}
done ne transporte jamais de valeur — on le ferme (close(done)) pour signaler la fin à toutes les goroutines qui l'écoutent, exactement comme chA/chB/chC plus haut.
Condition de course : régler ça avec un channel plutôt qu'un mutex
Au cours 6, la condition de course était réglée en protégeant la variable partagée avec un Mutex ou sync/atomic. Il existe une alternative idiomatique en Go, souvent résumée par la phrase "ne communiquez pas en partageant la mémoire ; partagez la mémoire en communiquant" : au lieu de laisser plusieurs goroutines accéder directement à une même variable, on confie l'unique accès à une seule goroutine propriétaire, et les autres lui envoient des requêtes via un channel.
type requete struct {
delta int
reponse chan int
}
func gestionnaireCompteur(requetes <-chan requete) {
compteur := 0
for r := range requetes {
compteur += r.delta
r.reponse <- compteur
}
}
func main() {
requetes := make(chan requete)
go gestionnaireCompteur(requetes)
reponse := make(chan int)
requetes <- requete{delta: 1, reponse: reponse}
fmt.Println("Nouveau total :", <-reponse)
}
Aucune donnée n'est jamais partagée entre goroutines — seule gestionnaireCompteur touche compteur, ce qui élimine la condition de course par construction, sans verrou explicite. C'est plus lent qu'un simple mutex pour un cas aussi trivial, mais ce patron devient précieux dès que la logique autour de la donnée partagée se complexifie (validation, agrégation, effets de bord).