zdpgo_pool_goroutine

package module
v0.1.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 11, 2022 License: MIT Imports: 8 Imported by: 1

README

zdpgo_pool_goroutine

go的协程池,基于ants二次开发

版本历史

  • v0.1.0 2022/07/08 新增:基础代码
  • v0.1.1 2022/08/03 新增:批量执行无参任务和带参任务
  • v0.1.2 2022/08/11 优化:带参泛型改为any

使用示例

请查看 examples 目录

Documentation

Index

Constants

View Source
const (
	DefaultCleanIntervalTime = time.Second * 3 // 默认的清除间隔时间
	OPENED                   = iota            // 连接池开启状态
	CLOSED                                     // 连接池关闭状态
)

Variables

View Source
var (
	// ErrLackPoolFunc will be returned when invokers don't provide function for pool.
	ErrLackPoolFunc = errors.New("must provide function for pool")

	// ErrInvalidPoolExpiry will be returned when setting a negative number as the periodic duration to purge goroutines.
	ErrInvalidPoolExpiry = errors.New("invalid expiry for pool")

	// ErrPoolClosed will be returned when submitting task to a closed pool.
	ErrPoolClosed = errors.New("this pool has been closed")

	// ErrPoolOverload will be returned when the pool is full and no workers available.
	ErrPoolOverload = errors.New("too many goroutines blocked on submit or Nonblocking is set")

	// ErrInvalidPreAllocSize will be returned when trying to set up a negative capacity under PreAlloc mode.
	ErrInvalidPreAllocSize = errors.New("can not set up a negative capacity under PreAlloc mode")

	// ErrTimeout will be returned after the operations timed out.
	ErrTimeout = errors.New("operation timed out")
)

Functions

func Cap

func Cap() int

Cap 返回默认连接池的容量

func Free

func Free() int

Free 获取可用的Goroutine数量

func NewSpinLock added in v0.1.1

func NewSpinLock() sync.Locker

NewSpinLock instantiates a spin-lock.

func Reboot

func Reboot()

Reboot 重启默认的连接池

func Release

func Release()

Release 关闭默认的连接池

func RunBatchArgTask added in v0.1.1

func RunBatchArgTask[T any](poolSize int, funcObj func(arg T), args []T)

RunBatchArgTask 批量执行带参数的任务 @param poolSize 协程池的容量大小 @param funcObj 要执行的带参方法 @param args 参数列表,会将每个参数都传给funcObj并发执行

func RunBatchTask added in v0.1.1

func RunBatchTask(funcList []func())

RunBatchTask 批量执行任务

func Running

func Running() int

Running 返回当前整型运行的Goroutine的数量

func Submit

func Submit(task func()) error

Submit 提交任务到连接池

Types

type Logger

type Logger interface {
	// Printf 必须和 log.Printf 具有相同的实现
	Printf(format string, args ...interface{})
}

Logger 用于日志格式化

type Option

type Option func(opts *Options)

Option represents the optional function.

func WithExpiryDuration

func WithExpiryDuration(expiryDuration time.Duration) Option

WithExpiryDuration sets up the interval time of cleaning up goroutines.

func WithLogger

func WithLogger(logger Logger) Option

WithLogger sets up a customized logger.

func WithMaxBlockingTasks

func WithMaxBlockingTasks(maxBlockingTasks int) Option

WithMaxBlockingTasks sets up the maximum number of goroutines that are blocked when it reaches the capacity of pool.

func WithNonblocking

func WithNonblocking(nonblocking bool) Option

WithNonblocking indicates that pool will return nil when there is no available workers.

func WithOptions

func WithOptions(options Options) Option

WithOptions accepts the whole options config.

func WithPanicHandler

func WithPanicHandler(panicHandler func(interface{})) Option

WithPanicHandler sets up panic handler.

func WithPreAlloc

func WithPreAlloc(preAlloc bool) Option

WithPreAlloc indicates whether it should malloc for workers.

type Options

type Options struct {
	// 过期时间。表示 goroutine 空闲多长时间之后会被ants池回收
	ExpiryDuration time.Duration
	// 预分配。调用NewPool()/NewPoolWithFunc()之后预分配worker(管理一个工作 goroutine 的结构体)切片。
	// 而且使用预分配与否会直接影响池中管理worker的结构。
	PreAlloc bool

	// 最大阻塞任务数量。
	// 即池中 goroutine 数量已到池容量,且所有 goroutine 都处理繁忙状态,这时到来的任务会在阻塞列表等待。
	// 这个选项设置的是列表的最大长度。阻塞的任务数量达到这个值后,后续任务提交直接返回失败
	MaxBlockingTasks int

	// When Nonblocking is true, Pool.Submit will never be blocked.
	// ErrPoolOverload will be returned when Pool.Submit cannot be done at once.
	// When Nonblocking is true, MaxBlockingTasks is inoperative.
	// 池是否阻塞,默认阻塞。
	// 提交任务时,如果ants池中 goroutine 已到上限且全部繁忙,阻塞的池会将任务添加的阻塞列表等待(当然受限于阻塞列表长度,见上一个选项)。
	// 非阻塞的池直接返回失败
	Nonblocking bool

	// PanicHandler is used to handle panics from each worker goroutine.
	// if nil, panics will be thrown out again from worker goroutines.
	PanicHandler func(interface{})

	// Logger is the customized logger for logging info, if it is not set,
	// default standard logger from log package is used.
	Logger Logger
}

Options contains all options which will be applied when instantiating an ants pool.

type Pool

type Pool struct {
	// contains filtered or unexported fields
}

Pool accepts the tasks from client, it limits the total of goroutines to a given number by recycling goroutines.

func NewPool

func NewPool(size int, options ...Option) (*Pool, error)

NewPool generates an instance of ants pool.

func (*Pool) Cap

func (p *Pool) Cap() int

Cap returns the capacity of this pool.

func (*Pool) Free

func (p *Pool) Free() int

Free returns the number of available goroutines to work, -1 indicates this pool is unlimited.

func (*Pool) IsClosed

func (p *Pool) IsClosed() bool

IsClosed indicates whether the pool is closed.

func (*Pool) Reboot

func (p *Pool) Reboot()

Reboot reboots a closed pool.

func (*Pool) Release

func (p *Pool) Release()

Release closes this pool and releases the worker queue.

func (*Pool) ReleaseTimeout

func (p *Pool) ReleaseTimeout(timeout time.Duration) error

ReleaseTimeout is like Release but with a timeout, it waits all workers to exit before timing out.

func (*Pool) Running

func (p *Pool) Running() int

Running returns the number of workers currently running.

func (*Pool) Submit

func (p *Pool) Submit(task func()) error

Submit submits a task to this pool.

Note that you are allowed to call Pool.Submit() from the current Pool.Submit(), but what calls for special attention is that you will get blocked with the latest Pool.Submit() call once the current Pool runs out of its capacity, and to avoid this, you should instantiate a Pool with ants.WithNonblocking(true).

func (*Pool) Tune

func (p *Pool) Tune(size int)

Tune changes the capacity of this pool, note that it is noneffective to the infinite or pre-allocation pool.

func (*Pool) Waiting

func (p *Pool) Waiting() int

Waiting returns the number of tasks which are waiting be executed.

type PoolWithFunc

type PoolWithFunc struct {
	// contains filtered or unexported fields
}

PoolWithFunc accepts the tasks from client, it limits the total of goroutines to a given number by recycling goroutines.

func NewPoolWithFunc

func NewPoolWithFunc(size int, pf func(interface{}), options ...Option) (*PoolWithFunc, error)

NewPoolWithFunc generates an instance of ants pool with a specific function.

func (*PoolWithFunc) Cap

func (p *PoolWithFunc) Cap() int

Cap returns the capacity of this pool.

func (*PoolWithFunc) Free

func (p *PoolWithFunc) Free() int

Free returns the number of available goroutines to work, -1 indicates this pool is unlimited.

func (*PoolWithFunc) Invoke

func (p *PoolWithFunc) Invoke(args interface{}) error

Invoke submits a task to pool.

Note that you are allowed to call Pool.Invoke() from the current Pool.Invoke(), but what calls for special attention is that you will get blocked with the latest Pool.Invoke() call once the current Pool runs out of its capacity, and to avoid this, you should instantiate a PoolWithFunc with ants.WithNonblocking(true).

func (*PoolWithFunc) IsClosed

func (p *PoolWithFunc) IsClosed() bool

IsClosed indicates whether the pool is closed.

func (*PoolWithFunc) Reboot

func (p *PoolWithFunc) Reboot()

Reboot reboots a closed pool.

func (*PoolWithFunc) Release

func (p *PoolWithFunc) Release()

Release closes this pool and releases the worker queue.

func (*PoolWithFunc) ReleaseTimeout

func (p *PoolWithFunc) ReleaseTimeout(timeout time.Duration) error

ReleaseTimeout is like Release but with a timeout, it waits all workers to exit before timing out.

func (*PoolWithFunc) Running

func (p *PoolWithFunc) Running() int

Running returns the number of workers currently running.

func (*PoolWithFunc) Tune

func (p *PoolWithFunc) Tune(size int)

Tune changes the capacity of this pool, note that it is noneffective to the infinite or pre-allocation pool.

func (*PoolWithFunc) Waiting

func (p *PoolWithFunc) Waiting() int

Waiting returns the number of tasks which are waiting be executed.

Directories

Path Synopsis
examples
ants_with_func command
ants_with_pool command
batch_task_sum command
big_sum command
run_batch_task command

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL