पाठ 18 / 25

Worker Pools and Bounded Concurrency

Limiting how much runs at once.

Three ways to put a ceiling on concurrency

Starting one goroutine per item is fine for ten items and dangerous for a million: memory, open connections and downstream load all grow without bound. Three common bounds: (1) a fixed pool of N workers reading from a jobs channel, as in the fan-out topic; (2) a semaphore built from a buffered channel of capacity N: acquire by sending a token before starting work and release by receiving it afterwards, so at most N goroutines are active; (3) errgroup.Group.SetLimit(n) from golang.org/x/sync/errgroup, which makes g.Go block until a slot is free and also collects errors. There is also a weighted semaphore in golang.org/x/sync/semaphore for jobs of different sizes. Choose N from the bottleneck: CPU cores for computation, the connection-pool size or API quota for I/O.

A buffered-channel semaphore and errgroup.SetLimit

Both keep at most 4 downloads in flight.

package main

import (
	"context"
	"fmt"
	"sync"

	"golang.org/x/sync/errgroup"
)

func download(ctx context.Context, url string) error {
	fmt.Println("downloading", url)
	return nil
}

func withSemaphore(ctx context.Context, urls []string) {
	sem := make(chan struct{}, 4) // capacity = max concurrency
	var wg sync.WaitGroup
	for _, u := range urls {
		sem <- struct{}{} // acquire: blocks while 4 are running
		wg.Add(1)
		go func() {
			defer wg.Done()
			defer func() { <-sem }() // release
			_ = download(ctx, u)
		}()
	}
	wg.Wait()
}

func withErrgroup(ctx context.Context, urls []string) error {
	g, ctx := errgroup.WithContext(ctx)
	g.SetLimit(4) // g.Go blocks until fewer than 4 are active
	for _, u := range urls {
		g.Go(func() error { return download(ctx, u) })
	}
	return g.Wait()
}

func main() {
	urls := []string{"a", "b", "c", "d", "e", "f"}
	withSemaphore(context.Background(), urls)
	if err := withErrgroup(context.Background(), urls); err != nil {
		fmt.Println("error:", err)
	}
}

Bound at the source

Acquire the semaphore before starting the goroutine, not inside it. Acquiring inside still creates one parked goroutine per item, which defeats the point for very large inputs.

त्वरित जाँच: How does a buffered channel act as a semaphore?

  • Its capacity limits how many tokens can be held, so sends block once N workers are active
  • Buffered channels automatically limit the number of goroutines in the program
  • Closing the channel releases all workers at once
  • The runtime schedules only cap(ch) goroutines
Answer

Its capacity limits how many tokens can be held, so sends block once N workers are active — Send to acquire, receive to release; capacity is the limit.