Documentation
¶
Overview ¶
Package executors provides batching executors that coalesce sporadic tasks into fewer, larger executions: periodical (time-driven), bulk (count-driven), delay (debounced), chunk (byte-size-driven) and less (rate-limited).
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type BulkExecutor ¶
type BulkExecutor struct {
// contains filtered or unexported fields
}
BulkExecutor executes tasks when either the batch size is reached or the flush interval elapses.
func NewBulkExecutor ¶
func NewBulkExecutor(execute Execute, opts ...BulkOption) *BulkExecutor
NewBulkExecutor returns a BulkExecutor that executes batches with execute.
func (*BulkExecutor) Add ¶
func (be *BulkExecutor) Add(task any) error
Add adds a task to the batch.
func (*BulkExecutor) Wait ¶
func (be *BulkExecutor) Wait()
Wait flushes and waits for execution to finish.
type BulkOption ¶
type BulkOption func(*bulkOptions)
BulkOption customizes a BulkExecutor.
func WithBulkInterval ¶
func WithBulkInterval(duration time.Duration) BulkOption
WithBulkInterval sets the flush interval.
func WithBulkTasks ¶
func WithBulkTasks(tasks int) BulkOption
WithBulkTasks sets the batch size that triggers a flush.
type ChunkExecutor ¶
type ChunkExecutor struct {
// contains filtered or unexported fields
}
ChunkExecutor executes tasks when either the accumulated chunk size is reached or the flush interval elapses.
func NewChunkExecutor ¶
func NewChunkExecutor(execute Execute, opts ...ChunkOption) *ChunkExecutor
NewChunkExecutor returns a ChunkExecutor that executes chunks with execute.
func (*ChunkExecutor) Add ¶
func (ce *ChunkExecutor) Add(task any, size int) error
Add adds a task with the given byte size to the chunk.
func (*ChunkExecutor) Wait ¶
func (ce *ChunkExecutor) Wait()
Wait flushes and waits for execution to finish.
type ChunkOption ¶
type ChunkOption func(*chunkOptions)
ChunkOption customizes a ChunkExecutor.
func WithChunkBytes ¶
func WithChunkBytes(size int) ChunkOption
WithChunkBytes sets the chunk size that triggers a flush.
func WithFlushInterval ¶
func WithFlushInterval(duration time.Duration) ChunkOption
WithFlushInterval sets the flush interval.
type DelayExecutor ¶
type DelayExecutor struct {
// contains filtered or unexported fields
}
DelayExecutor runs a task after a debounce delay, coalescing repeated triggers within the delay window.
func NewDelayExecutor ¶
func NewDelayExecutor(fn func(), delay time.Duration) *DelayExecutor
NewDelayExecutor returns a DelayExecutor running fn after delay.
func (*DelayExecutor) Trigger ¶
func (de *DelayExecutor) Trigger()
Trigger schedules fn to run after the delay. It is safe to call repeatedly; subsequent triggers within the delay window are coalesced.
type LessExecutor ¶
type LessExecutor struct {
// contains filtered or unexported fields
}
LessExecutor runs a task at most once per threshold interval.
func NewLessExecutor ¶
func NewLessExecutor(threshold time.Duration) *LessExecutor
NewLessExecutor returns a LessExecutor with the given threshold interval.
func (*LessExecutor) DoOrDiscard ¶
func (le *LessExecutor) DoOrDiscard(fn func()) bool
DoOrDiscard runs fn if threshold has elapsed since the last execution, and discards it otherwise. It reports whether fn ran. The last-execution time is updated with a compare-and-swap so concurrent callers within a threshold window still execute at most once.
type PeriodicalExecutor ¶
type PeriodicalExecutor struct {
// contains filtered or unexported fields
}
PeriodicalExecutor coalesces tasks and executes them periodically, or immediately when the container signals fullness.
func NewPeriodicalExecutor ¶
func NewPeriodicalExecutor(interval time.Duration, container TaskContainer) *PeriodicalExecutor
NewPeriodicalExecutor returns a PeriodicalExecutor with the given interval and container. Call Flush (or Wait) before shutdown to drain pending tasks.
func (*PeriodicalExecutor) Add ¶
func (pe *PeriodicalExecutor) Add(task any)
Add adds a task to be executed.
func (*PeriodicalExecutor) Flush ¶
func (pe *PeriodicalExecutor) Flush() bool
Flush forces pe to execute all pending tasks.
func (*PeriodicalExecutor) Sync ¶
func (pe *PeriodicalExecutor) Sync(fn func())
Sync runs fn while holding the executor's container lock.
func (*PeriodicalExecutor) Wait ¶
func (pe *PeriodicalExecutor) Wait()
Wait flushes and waits for all executions to finish.
type TaskContainer ¶
type TaskContainer interface {
// AddTask adds task, returning true if the container needs flushing now.
AddTask(task any) bool
// Execute handles the collected tasks.
Execute(tasks any)
// RemoveAll removes and returns all collected tasks.
RemoveAll() any
}
TaskContainer accumulates tasks for a PeriodicalExecutor.