func process(stream chan int) {
p := pool.New().WithMaxGoroutines(10)
for elem := range stream {
elem := elem
p.Go(func() {
handle(elem)
})
}
p.Wait()
I did something similar just with input (optionally output) channel. Close input, goroutines stop, when all of them stop the output is closed [1]. No need to incur function call just to add elements (although I'd imagine go would just inline it so it might not matter either way)This
func mapStream(
in chan int,
out chan int,
f func(int) int,
) {
s := stream.New().WithMaxGoroutines(10)
for elem := range in {
elem := elem
s.Go(func() stream.Callback {
res := f(elem)
return func() { out <- res }
})
}
s.Wait()
}
also seems awfully verbose vs just function (that is now easy and safe thanks to generics) with this signature WorkerPool[T1, T2 any](input chan T1, output chan T2, worker func(T1) T2, concurrency int)
I do like idea of waitgroup on steroids, I might steal it for my generic library.* [1] https://github.com/XANi/goneric/blob/master/worker.go#L92