You need to process 10 million records from a file with bounded memory and bounded parallelism. How do you design it?
Question 245HardGo 1.22 to 1.25
The core ideas: stream the input instead of loading it, bound the number of goroutines, and apply backpressure through channels.
func Process(ctx context.Context, r io.Reader, workers int) error {
g, ctx := errgroup.WithContext(ctx)
lines := make(chan string, workers*2) // small buffer = backpressure
g.Go(func() error { // single producer
defer close(lines)
sc := bufio.NewScanner(r)
for sc.Scan() {
select {
case lines <- sc.Text():
case <-ctx.Done():
return ctx.Err()
}
}
return sc.Err()
})
for range workers { // fixed consumers
g.Go(func() error {
for line := range lines {
if err := handle(ctx, line); err != nil {
return err // cancels ctx; producer stops
}
}
return nil
})
}
return g.Wait()
}
Talking points:
- Memory is O(workers + buffer), not O(N).
- Pick the worker count by workload: about
runtime.GOMAXPROCS(0)for CPU-bound work, and higher for I/O-bound work, capped by downstream limits such as the DB pool size or rate limits. - Batch writes to cut per-record overhead.
bufio.Scannerfails withbufio.ErrTooLongon lines over 64 KB by default. Callsc.Buffer(make([]byte, 0, 1<<20), 16<<20)if records can be long.- If a worker exits on error while the producer is blocked, the ctx cancellation unblocks the producer. If
handleignores ctx, the remaining workers still drainlines. - Spawning a goroutine per record (10M goroutines) "works" but wastes memory and hammers dependencies.
More on Concurrency Patterns & sync
- Q243How do you implement a semaphore in Go? Compare a buffered channel with golang.org/x/sync/semaphore.
- Q244Explain errgroup: WithContext, SetLimit, TryGo. What are its semantics and gotchas?
- Q246How do you implement rate limiting in Go? Compare time.Ticker with golang.org/x/time/rate.
- Q247What problem does singleflight solve, and what are its gotchas?
- Q248Implement graceful shutdown for an HTTP server with background workers.
- Q249Which newer context features matter for concurrent code: WithCancelCause, AfterFunc, WithoutCancel?