pond is a minimalistic and high-performance Go library designed to elegantly manage concurrent tasks.
Motivation
This library is meant to provide a simple and idiomatic way to manage concurrency in Go programs. Based on the Worker Pool pattern, it allows running a large number of tasks concurrently while limiting the number of goroutines that are running at the same time. This is useful when you need to limit the number of concurrent operations to avoid resource exhaustion or hitting rate limits.
Some common use cases include: - Processing a large number of tasks concurrently - Limiting the number of concurrent HTTP requests - Limiting the number of concurrent database connections - Sending HTTP requests to a rate-limited API
Features:
- Zero dependencies
- Create pools of goroutines that scale automatically based on the number of tasks submitted
- Limit the number of concurrent tasks running at the same time
- Worker goroutines are only created when needed and immediately removed when idle (scale to zero)
- Minimalistic and fluent APIs for:
- Very high performance and efficient resource usage under heavy workloads, even outperforming unbounded goroutines in some scenarios
- Complete pool metrics such as number of running workers, tasks waiting in the queue and more
- Configurable parent context to stop all workers when it is cancelled
- New features in v2:
Installation
go get -u github.com/alitto/pond/v2
Usage
Submitting tasks to a pool with limited concurrency
`` go
package main
import ( "fmt"
"github.com/alitto/pond/v2" )
func main() {
// Create a pool with limited concurrency pool := pond.NewPool(100)
// Submit 1000 tasks for i := 0; i < 1000; i++ { i := i pool.Submit(func() { fmt.Printf("Running task #%d\n", i) }) }
// Stop the pool and wait for all submitted tasks to complete pool.StopAndWait() }
go // Create a pool with limited concurrency pool := pond.NewPool(100)### Submitting tasks that return errorserrorThis feature allows you to submit tasks that return an
. This is useful when you need to handle errors that occur during the execution of a task.
// Submit a task that returns an error task := pool.SubmitErr(func() error { return errors.New("An error occurred") })
// Wait for the task to complete and get the error err := task.Wait()
### Submitting tasks that return results
This feature allows you to submit tasks that return a value. This is useful when you need to process the result of a task.
go
// Create a pool that accepts tasks that return a string and an error
pool := pond.NewResultPoolstring
// Submit a task that returns a string task := pool.Submit(func() (string) { return "Hello, World!" })
// Wait for the task to complete and get the result result, err := task.Wait() // result = "Hello, World!" and err = nil
go // Create a concurrency limited pool that accepts tasks that return a string pool := pond.NewResultPoolstring### Submitting tasks that return results or errorserrorThis feature allows you to submit tasks that return a value and an
. This is useful when you need to handle errors that occur during the execution of a task.
// Submit a task that returns a string value or an error task := pool.SubmitErr(func() (string, error) { return "Hello, World!", nil })
// Wait for the task to complete and get the result result, err := task.Wait() // result = "Hello, World!" and err = nil
### Submitting tasks associated with a context
If you need to submit a task that is associated with a context, you can pass the context directly to the task function.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a context that can be cancelled ctx, cancel := context.WithCancel(context.Background())
// Submit a task that is associated with a context task := pool.SubmitErr(func() error { return doSomethingWithCtx(ctx) // Pass the context to the task directly })
// Wait for the task to complete and get the error. // If the context is cancelled, the task is stopped and an error is returned. err := task.Wait()
### Submitting a group of related tasks
You can submit a group of tasks that are related to each other. This is useful when you need to execute a group of tasks concurrently and wait for all of them to complete.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a task group group := pool.NewGroup()
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.Submit(func() { fmt.Printf("Running group task #%d\n", i) }) }
// Wait for all tasks in the group to complete err := group.Wait()
### Submitting a group of related tasks associated with a context
You can submit a group of tasks that are linked to a context. This is useful when you need to execute a group of tasks concurrently and stop them when the context is cancelled (e.g. when the parent task is cancelled or times out).
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a context with a 5s timeout timeout, _ := context.WithTimeout(context.Background(), 5*time.Second)
// Create a task group with a context group := pool.NewGroupContext(timeout)
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.Submit(func() { fmt.Printf("Running group task #%d\n", i) }) }
// Wait for all tasks in the group to complete or the timeout to occur, whichever comes first err := group.Wait()
### Submitting a group of related tasks and waiting for the first error
You can submit a group of tasks that are related to each other and wait for the first error to occur. This is useful when you need to execute a group of tasks concurrently and stop the execution if an error occurs.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a task group group := pool.NewGroup()
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.SubmitErr(func() error { return doSomethingThatCanFail() }) }
// Wait for all tasks in the group to complete or the first error to occur err := group.Wait()
go // Create a pool with limited concurrency pool := pond.NewPool(10)#### Cancelling tasks immediatelygroup.Wait()When the first error occurs, tasks that are in the queue will be aborted but any running task will not be disrupted. The call to
will not wait for these "in-flight" tasks to complete but they will continue running until completion nonetheless.group.Context()If you also need to stop these "in-flight" tasks when the first error occurs, you can reference the group's context (accessible via
) from any long-running operation carried out within these tasks. Here's an example:
// Create a task group group := pool.NewGroup()
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.SubmitErr(func() error { // Cancel all in-flight tasks when the first error occurs return doSomethingThatCanFailInContext(group.Context()) }) }
// Wait for all tasks in the group to complete or the first error to occur err := group.Wait()
### Submitting a group of related tasks that return results
You can submit a group of tasks that are related to each other and return results. This is useful when you need to execute a group of tasks concurrently and process the results. Results are returned in the order they were submitted.
go
// Create a pool with limited concurrency
pool := pond.NewResultPoolstring
// Create a task group group := pool.NewGroup()
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.Submit(func() string { return fmt.Sprintf("Task #%d", i) }) }
// Wait for all tasks in the group to complete results, err := group.Wait() // results = ["Task #0", "Task #1", ..., "Task #19"] and err = nil
### Stopping a group of tasks when a context is cancelled
If you need to submit a group of tasks that are associated with a context and stop them when the context is cancelled, you can pass the context directly to the task function.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a context that can be cancelled ctx, cancel := context.WithCancel(context.Background())
// Create a task group group := pool.NewGroupContext(ctx)
// Submit a group of tasks for i := 0; i < 20; i++ { i := i group.SubmitErr(func() error { return doSomethingWithCtx(ctx) // Pass the context to the task directly }) }
// Wait for all tasks in the group to complete. // If the context is cancelled, all tasks are stopped and the first error is returned. err := group.Wait()
go // Create a custom context that can be cancelled customCtx, cancel := context.WithCancel(context.Background())### Using a custom Context at the pool levelcontext.Background()Each pool is associated with a context that is used to stop all workers when the pool is stopped. By default, the context is the background context (
). You can create a custom context and pass it to the pool to stop all workers when the context is cancelled.
// This creates a pool that is stopped when customCtx is cancelled
pool := pond.NewPool(10, pond.WithContext(customCtx))
`
When a pool-level context is canceled, workers stop accepting new work and any queued tasks are drained so the pool can shut down cleanly without deadlocking. Drained tasks are not executed and resolve with a cancellation error.
Stopping a pool
You can stop a pool using the Stop method. This will stop all workers and prevent new tasks from being submitted. You can also wait for all submitted tasks by calling the Wait method.
` go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Submit a task pool.Submit(func() { fmt.Println("Running task") })
// Stop the pool and wait for all submitted tasks to complete pool.Stop().Wait()
go // Create a pool with limited concurrency pool := pond.NewPool(10)A shorthand methodStopAndWaitis also available to stop the pool and wait for all submitted tasks to complete.
// Submit a task pool.Submit(func() { fmt.Println("Running task") })
// Stop the pool and wait for all submitted tasks to complete pool.StopAndWait()
### Recovering from panics
By default, panics that occur during the execution of a task are captured and returned as errors. This allows you to recover from panics and handle them gracefully.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Submit a task that panics task := pool.Submit(func() { panic("A panic occurred") })
// Wait for the task to complete and get the error err := task.Wait()
if err != nil { fmt.Printf("Failed to run task: %v", err) } else { fmt.Println("Task completed successfully") }
If you prefer to keep the default Go runtime behavior (panics crashing the goroutine), you can disable panic interception when creating the pool: go
pool := pond.NewPool(10, pond.WithoutPanicRecovery())
pool.Go(func() { panic("this panic will crash the goroutine") })
### Subpools (v2)
Subpools are pools that can have a fraction of the parent pool's maximum number of workers. This is useful when you need to create a pool of workers that can be used for a specific task or group of tasks.
go
// Create a pool with limited concurrency
pool := pond.NewPool(10)
// Create a subpool with a fraction of the parent pool's maximum number of workers subpool := pool.NewSubpool(5)
// Submit a task to the subpool subpool.Submit(func() { fmt.Println("Running task in subpool") })
// Stop the subpool and wait for all submitted tasks to complete subpool.StopAndWait()
### Default pool (v2)
The default pool is a global pool that is used when no pool is provided. This is useful when you need to submit tasks but don't want to create a pool explicitly.
The default pool does not have a maximum number of workers and scales automatically based on the number of tasks submitted.
go
// Submit a task to the default pool and wait for it to complete
err := pond.SubmitErr(func() error {
fmt.Println("Running task in default pool")
return nil
}).Wait()
if err != nil { fmt.Printf("Failed to run task: %v", err) } else { fmt.Println("Task completed successfully") }
go // Create a pool with a maximum of 10 tasks in the queue pool := pond.NewPool(1, pond.WithQueueSize(10))### Bounded task queues (v2)WithQueueSizeBy default, task queues are unbounded, meaning that tasks are queued indefinitely until the pool is stopped (or the process runs out of memory). You can limit the number of tasks that can be queued by setting a queue size when creating a pool (
option).
go// Create a pool with a maximum of 10 tasks in the queue pool := pond.NewPool(1, pond.WithQueueSize(10))Blocking vs non-blocking task submissionTrySubmitWhen a pool defines a queue size (bounded), you can also specify how to handle tasks submitted when the queue is full. This can be done in two ways:
- Per task submission: You can use
andTrySubmitErrmethods to attempt to submit a task and get a boolean indicating whether the task was submitted successfully.
// Submit a task to the pool task, ok := pool.TrySubmit(func() { // Do some work })
// Check if the task was submitted successfully if !ok { fmt.Println("Task submission failed because the queue is full") }
go // Create a pool with a maximum of 10 tasks in the queue and non-blocking task submission pool := pond.NewPool(1, pond.WithQueueSize(10), pond.WithNonBlocking(true))- Global non-blocking task submission: You can set theNonBlockingoption totruewhen creating a pool to enable non-blocking task submission. If the queue is full and non-blocking task submission is enabled, the task is dropped and an error is returned (ErrQueueFull).
go // Force the queue back to the default unbounded behavior. pool := pond.NewPool(8, pond.WithQueueSize(pond.Unbounded), pond.WithNonBlocking(false))Note: you can technically useTrySubmitandTrySubmitErrmethods orNonBlockingoption for both bounded and unbounded task queues, but they are only useful when the queue is bounded. For unbounded task queues there is always space in the queue and non-blocking submission will always succeed.pond.UnboundedUsing the
constantpond.Unboundedis a convenience constant that resolves to the maximum supported queue size. Passing it toWithQueueSizekeeps the queue effectively infinite while still allowing other queue-related options to be configured explicitly (for example in subpools or helper constructors).
go // Create a pool that never queues tasks: they either run immediately or error. pool := pond.NewPool(4, pond.WithQueueSize(0))#### Disabling the queue altogetherWithQueueSize(0)Passing
disables the queue completely. Only running workers can accept work; if all workers are busy, submissions will either block (default) until a worker becomes available or fail immediately when usingTrySubmit,TrySubmitErr, orWithNonBlocking(true). This is helpful when you need "no backlog" semantics and want to ensure work runs immediately or gets rejected.
if _, ok := pool.TrySubmit(func() { // Do some work that must start right away }); !ok { // All 4 workers were busy so the task was rejected }
go // Create a pool with 5 workers pool := pond.NewPool(5)### Resizing pools (v2)ResizeYou can dynamically change the maximum number of workers in a pool using the
method. This is useful when you need to adjust the pool's capacity based on runtime conditions.
// Submit some tasks for i := 0; i < 20; i++ { pool.Submit(func() { // Do some work }) }
// Increase the pool size to 10 workers pool.Resize(10)
// Submit more tasks that will use the increased capacity for i := 0; i < 20; i++ { pool.Submit(func() { // Do some work }) }
// Decrease the pool size back to 5 workers
pool.Resize(5)
`
When resizing a pool:
- The new maximum concurrency must be greater than or equal to 0 (0 means no limit)
- If you increase the size, new workers will be created as needed up to the new maximum
- If you decrease the size, existing workers will continue running until they complete their current tasks, but no new workers will be created until the number of running workers is below the new maximum
Metrics & monitoring
Each worker pool instance exposes useful metrics that can be queried through the following methods:
- pool.RunningWorkers() int64
: Current number of running workers - pool.SubmittedTasks() uint64
: Total number of tasks submitted since the pool was created and before it was stopped. This includes tasks that were dropped because the queue was full - pool.WaitingTasks() uint64
: Current number of tasks in the queue that are waiting to be executed - pool.SuccessfulTasks() uint64
: Total number of tasks that have successfully completed their execution since the pool was created - pool.FailedTasks() uint64
: Total number of tasks that completed with a non-cancellation error (including panics) since the pool was created - pool.CanceledTasks() uint64
: Total number of tasks accepted by the pool that were canceled before executing user code due to context cancellation - pool.CompletedTasks() uint64
: Total number of tasks that have completed their execution either successfully or with a non-cancellation error since the pool was created - pool.DroppedTasks() uint64
: Total number of tasks that were dropped because the queue was full since the pool was created
Migrating from pond v1 to v2
If you are using pond v1, here are the changes you need to make to migrate to v2:
- Update the import path to github.com/alitto/pond/v2
- Replace pond.New(100, 1000)
withpond.NewPool(100). The second argument is no longer needed since task queues are unbounded by default. - The pool option pond.Context
was renamed topond.WithContext - The following pool options were deprecated:
: This option is no longer needed since workers are created on demand and removed when idle.
- pond.IdleTimeout: This option is no longer needed since workers are immediately removed when idle.
- pond.PanicHandler: Panics are captured and returned as errors. You can handle panics by checking the error returned by the Wait method.
- pond.Strategy: The pool now scales automatically based on the number of tasks submitted.
- The
pool.StopAndWaitFor method was deprecated. Use pool.Stop().Done() channel if you need to wait for the pool to stop in a select statement.
The pool.Group method was renamed to pool.NewGroup.
The pool.GroupContext was renamed to pool.NewGroupWithContext`.
Examples
You can find more examples in the examples directory.
API Reference
Full API reference is available at https://pkg.go.dev/github.com/alitto/pond/v2
Benchmarks
See Benchmarks.
Resources
Here are some of the resources which have served as inspiration when writing this library:
- http://marcio.io/2015/07/handling-1-million-requests-per-minute-with-golang/
- https://brandur.org/go-worker-pool
- https://gobyexample.com/worker-pools
- https://github.com/panjf2000/ants
- https://github.com/gammazero/workerpool
Contribution & Support
Feel free to send a pull request if you consider there's something that should be polished or improved. Also, please open up an issue if you run into a problem when using this library or just have a question about it.