3. Implémentation du Worker Pool
Modèle canonique en Go pour réguler un traitement par lots :
1func StartWorkerPool(numWorkers int, jobs <-chan Job, results chan<- Result) {
2 for w := 0; w < numWorkers; w++ {
3 go func() {
4 for job := range jobs {
5 res, err := ProcessJob(job)
6 results <- Result{JobID: job.ID, Output: res, Err: err}
7 }
8 }()
9 }
10}
- Chaque travailleur boucle sur le canal
jobset s'éteint proprement dès sa fermeture.
4. Dimensionnement Idéal des Workers
Comment calibrer le nombre optimal de goroutines travailleuses ?
Calculs purs, cryptographie, compression.
Formule : Nombre de Workers = Cœurs Physiques.
Ajouter des workers au-delà crée de la contention inutile.
Requêtes HTTP, lectures SQL, accès disque.
Formule : Nombre de Workers = Cœurs × 2 à 10.
Permet de masquer la latence d'attente réseau.
5. Patterns Fan-Out & Fan-In
Distribution parallèle d'un flux de données volumineux vers des travailleurs indépendants et regroupement des résultats :
6. Pipelines Étagés de Traitement
Découpage d'un algorithme en plusieurs étapes séquentielles indépendantes fonctionnant simultanément :
7. Canaux Tampons & Contre-Pression
Amortissement des à-coups de charge avec le mécanisme de contre-pression (Backpressure) :
L'afflux temporaire est absorbé par le tampon mémoire du canal (make(chan Task, 1000)).
Si le tampon se remplit, l'émetteur est ralenti à la vitesse exacte des consommateurs aval.
Empêche l'allocation infinie et protège les composants aval d'un écroulement mémoire.
8. Annulation Précoce avec Context
Interrompre immédiatement les opérations dès qu'un client abandonne pour éliminer tout calcul fantôme :
1func HandleRequest(w http.ResponseWriter, r *http.Request) {
2 // 1. Récupérer le context lié au socket client et fixer un timeout de 2s
3 ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second)
4 defer cancel()
5
6 // 2. Transmettre le context aux fonctions aval (SQL, HTTP)
7 res, err := QueryDatabase(ctx)
8 if errors.Is(err, context.Canceled) {
9 return // Requête abandonnée par le client : aucun calcul inutile !
10 }
11}
- Si le client coupe la connexion, le runtime Go déclenche
ctx.Done(), annulant instantanément la requête PostgreSQL en cours.
9. Délais d'Expiration & Boucles Réactives
Règle absolue : toujours écouter ctx.Done() dans les boucles de traitement intensives :
1func LongProcessing(ctx context.Context, items []Item) error {
2 for _, item := range items {
3 select {
4 case <-ctx.Done():
5 // Sortie immédiate si le délai est dépassé ou la requête annulée
6 return ctx.Err()
7 default:
8 // Poursuite du calcul si le temps est valide
9 processItem(item)
10 }
11 }
12 return nil
13}
10. Synthèse des Patterns Asynchrones
Les 3 piliers des architectures hautement scalables :
Utiliser systématiquement des Worker Pools pour maîtriser l'empreinte mémoire.
Zéro explosion de goroutines.Isoler les étapes séquentielles avec des canaux tamponnés pour amortir les pics.
Contre-pression naturelle.Transmettre context.Context à toutes les I/O pour interrompre les calculs inutiles.
Sobriété CPU garantie.TP Fil Rouge (Séance 6) : Contention de Verrous & sync.Pool
Mission : Diagnostiquer l'effondrement causé par les verrous partagés sur le compteur d'essais, passer aux opérations atomiques et recycler les buffers d'état de hachage.
1. Lock Contention sous pprof
Observer la dégradation causée par un sync.Mutex partagé lorsque tous les workers incrémentent un compteur centralisé. Utiliser pprof -mutex pour mesurer le temps d'attente sur le verrou.
2. Compteurs Atomiques
Remplacer le Mutex par des opérations atomiques lock-free matérielles (sync/atomic.AddUint64) pour éliminer les blocages d'ordonnancement.
3. Recyclage via sync.Pool
Instancier un sync.Pool de buffers d'états de hachage (sha256.New()) afin que chaque worker réutilise son instance de travail sans allouer sur le Tas.
4. Validation Débit Max
Vérifier au benchmark que le débit de calcul reste maximal et stable même avec une saturation complète de tous les cœurs CPU.