NOTE

2.15 Goroutine Pool

What a goroutine pool is, why it is usually unnecessary, and two example designs for reusing goroutines.

GoCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is a Goroutine Pool

A pool that reuses goroutines.

2. Why a Goroutine Pool Is Needed

Unlike How to Design Pooling Technology, goroutines have low creation and destruction overhead: creation and destruction happen in user space, the number of goroutines is theoretically unlimited, and one goroutine occupies less than 2 KB of memory. So in 99% of cases, a goroutine pool is unnecessary. But in extreme cases, reasonable reuse is always useful, such as ultra-high-concurrency, low-latency gateway scenarios.

3. How to Design a Goroutine Pool

Goroutine Pool

package main

import (
	"fmt"
	"time"
)

type IRunnable interface {
	// Execute the business method
	Run()
}

/* Define a task type Task */
type Task struct {
	f func() error // A Task contains one concrete business operation
}

func (this *Task) Run() {
	this.f()
}

func NewTask(f func() error) *Task {
	return &Task{f: f}
}

type IPool interface {
	// Submit a task
	Submit(IRunnable)
}

/* Pool goroutine pool */
type Pool struct {
	// This Channel is accessed externally
	EntryChannel chan IRunnable
	// Internal Task queue
	JobsChannel chan IRunnable
	// Maximum number of workers
	workerNum int
}

// Start all worker goroutines
func (this *Pool) startWorkers() {
	for i := 0; i < this.workerNum; i++ {
		go this.startOneWorker(i)
	}

	go this.tansferTaskToQueue()
}

// Start one worker goroutine
func (this *Pool) startOneWorker(workerId int) {
	for task := range this.JobsChannel {
		task.Run()
		fmt.Println("workerId: ", workerId, " finished one task")
	}
}

func (this *Pool) Submit(runnable IRunnable) {
	this.EntryChannel <- runnable
}

// Transfer tasks from the external Channel (EntryChannel) to the internal Channel (JobsChannel)
func (this *Pool) tansferTaskToQueue() {
	for task := range this.EntryChannel {
		this.JobsChannel <- task
	}
}

func NewPool(workerNum int) *Pool {
	pool := &Pool{
		EntryChannel: make(chan IRunnable),
		JobsChannel:  make(chan IRunnable),
		workerNum:    workerNum}
	pool.startWorkers()
	return pool
}

func main() {
	// Create a task
	task := NewTask(func() error {
		fmt.Println(time.Now())
		return nil
	})
	// Create a goroutine pool with 4 goroutines
	pool := NewPool(4)
	// Continuously submit tasks to the goroutine pool
	go func() {
		for true {
			pool.Submit(task)
		}
	}()
	//pool.startWorkers()
	select {}
}
  • Simplified version
package taskqueue

import (
	log "<INTERNAL_MODULE>/log"
	metrics "<INTERNAL_MODULE>/metrics"
)

// ITask ...
type ITask interface {
	// Execute the business method
	Run()
}

// Task ...
type Task struct {
	desc string
	f    func() error
}

// String ...
func (t *Task) String() string {
	return t.desc
}

// NewTask ...
func NewTask(desc string, f func() error) *Task {
	return &Task{
		desc: desc,
		f:    f,
	}
}

// Run ...
func (t *Task) Run() {
	err := t.f()
	if err != nil {
		metrics.Counter("task-execution-error").Incr()
		log.Errorf("Task Run: executing task error. err=%v", err)
	}
}

package taskqueue

import (
	trpc "<INTERNAL_MODULE>"
	log "<INTERNAL_MODULE>/log"
	metrics "<INTERNAL_MODULE>/metrics"
)

const (
	// DefaultWorkerNum ...
	DefaultWorkerNum = 10
	// DefaultQueueSize ...
	DefaultQueueSize = 10000
)

// GlobalTaskQueue ...
var GlobalTaskQueue = NewTaskQueue(DefaultWorkerNum, DefaultQueueSize)

// ITaskQueue ...
type ITaskQueue interface {
	// Submit a task to the queue
	Submit(ITask)
}

// TaskQueue ...
type TaskQueue struct {
	// Internal Task queue
	TaskQueue chan ITask
	// Maximum number of workers
	workerNum int
}

// NewTaskQueue ...
func NewTaskQueue(workerNum int, queueSize int) *TaskQueue {
	t := &TaskQueue{
		TaskQueue: make(chan ITask, queueSize),
		workerNum: workerNum,
	}
	t.startWorkers()
	return t
}

// Start all worker goroutines
func (p *TaskQueue) startWorkers() {
	for i := 0; i < p.workerNum; i++ {
		go p.startOneWorker(i)
	}
}

// startOneWorker
func (p *TaskQueue) startOneWorker(workerId int) {
	for task := range p.TaskQueue {
		task.Run()
		log.InfoContextf(trpc.BackgroundContext(),
			"TaskQueue: workerId %v execute task: %v", workerId, task)
	}
}

// Submit ...
func (p *TaskQueue) Submit(task ITask) {
	select {
	case p.TaskQueue <- task:
	default:
		metrics.Counter("task-queue-full-discard-task").Incr()
		log.ErrorContextf(trpc.BackgroundContext(),
			"TaskQueue Submit: queue is full, discard task %v", task)
	}
}

4. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub