This means that the queue could grow unboundedly, but that is ok, as this type of structure is not meant for constant request handling but rather for bursts of requests that can happen asynchronously (and in parallel), and allows us to limit the maximum number that can run while also never blocking callers if there are no available go routines to process their request when submitted.
Here is a minimially-viable example of how I do it:
https://play.golang.org/p/mD_jpMdoY_g
Contents here:
package main
import ( "fmt" "time" "sync" )
func main() { fmt.Println("Hello, playground") wg := &sync.WaitGroup{} wg.Add(1) numRequests := 100
inCh := make(chan int) go func() { resCh := make(chan int) queue := []int{} inFlight := 0 max := 20 completed := 0
for {
select{
case m := <- inCh:
queue = append(queue, m)
fmt.Println(fmt.Sprintf("Enqueing: %d", m))
case <- resCh:
inFlight -= 1
completed += 1
}
if len(queue) > 0 && inFlight < max {
v := queue[0]
queue = queue[1:]
fmt.Println(fmt.Sprintf("Submitting: %d", v))
inFlight += 1
go func(_v int) {
time.Sleep(1 * time.Second)
fmt.Println(fmt.Sprintf("run: %d", _v))
resCh <- v
}(v)
}
if completed == numRequests {
wg.Done()
}
}
}()
for i := 0; i < numRequests; i++ {
inCh <- i
}
fmt.Println("Waiting for completion...")
wg.Wait()
fmt.Println("All processes completed. Exiting.")
}