Count word frequencies in a large set of files concurrently and return the top K words. How do you merge per-worker maps efficiently?
Use a fixed pool of workers that read file paths from a channel. Each worker counts into its own private map, so the hot loop needs no locks and has no contention. After Wait, merge the maps and pick the top K.
type WordCount struct {
Word string
Count int
}
func TopK(ctx context.Context, paths []string, k, workers int) ([]WordCount, error) {
g, ctx := errgroup.WithContext(ctx)
jobs := make(chan string)
partial := make([]map[string]int, workers) // one slot per worker, no mutex
g.Go(func() error {
defer close(jobs)
for _, p := range paths {
select {
case jobs <- p:
case <-ctx.Done():
return ctx.Err()
}
}
return nil
})
for w := range workers {
g.Go(func() error {
local := make(map[string]int, 1<<12)
partial[w] = local
for p := range jobs {
if err := countFile(p, local); err != nil {
return err // cancels ctx; producer stops, jobs closes
}
}
return nil
})
}
if err := g.Wait(); err != nil { // Wait gives happens-before for partial
return nil, err
}
total := partial[0]
for _, m := range partial[1:] {
for w, c := range m {
total[w] += c
}
}
all := make([]WordCount, 0, len(total))
for w, c := range total {
all = append(all, WordCount{w, c})
}
slices.SortFunc(all, func(a, b WordCount) int {
return cmp.Or(cmp.Compare(b.Count, a.Count), strings.Compare(a.Word, b.Word))
})
return all[:min(k, len(all))], nil
}
func countFile(path string, counts map[string]int) error {
f, err := os.Open(path)
if err != nil {
return err
}
defer f.Close()
sc := bufio.NewScanner(f)
sc.Split(bufio.ScanWords)
for sc.Scan() {
counts[strings.ToLower(sc.Text())]++
}
return sc.Err()
}
Merging efficiently: merge into the largest map instead of allocating a new one. For huge vocabularies, shard by hash: each worker splits its counts into S sub-maps by hash(word)%S, then S goroutines each merge one shard across all workers in parallel. For top K, use a size-K min-heap (container/heap), which is O(U log K) instead of sorting all U words.
Anti-patterns interviewers watch for:
- One global map behind a mutex, locked per word, which gives massive contention.
- Sending every word over a channel. That costs roughly 100ns per word, which is far more than the counting itself.
sync.Mapwith read-modify-write increments, which is racy unless you use an atomic counter as the value.
Mention that bufio.Scanner has a 64KB default token limit (see sc.Buffer), and that the per-worker map is written by one goroutine and read only after Wait, so it is race-free.
More on Classic Concurrency Coding Problems
- Q524Implement a concurrent prime sieve with a daisy chain of goroutines. How many goroutines does it create, and how do you stop it without leaks?
- Q525Concurrently crawl URLs (the Tour of Go web crawler) with depth limits, deduplication and bounded parallelism.
- Q527Implement a debouncer and a throttler for a stream of events using time.Timer and channels.
- Q528Implement a sharded concurrent map (N shards with separate mutexes). When does it beat sync.Map and a single RWMutex?
- Q529Implement a readers-writer lock that prefers writers using only channels or sync.Mutex. Why is it hard to get right?
- Q530Implement a retry-with-timeout helper that runs a function in a goroutine and returns early when the per-attempt or overall deadline expires, without leaking goroutines.