Queues
Workflow queues allow you to ensure that workflow functions will be run, without starting them immediately. Queues are useful for controlling the number of workflows run in parallel, or the rate at which they are started.
Queue configuration is persisted to the system database, so any DBOS process connected to the same system database can register, retrieve, and reconfigure queues. If multiple applications share a system database, each queue is owned by the application that registers it, and only that application dequeues workflows from it.
Queue Management
RegisterQueue
func RegisterQueue(ctx Client, name string, options ...QueueOption) (Queue, error)
Register a queue and persist its configuration to the system database, returning a Queue.
If a queue with the same name already exists in the database, the WithQueueOnConflict option controls whether its configuration is overwritten.
Queues may be registered at any time, including after Launch(); live workers periodically reload queue configuration, so changes take effect without a restart.
You can enqueue a workflow using the WithQueue parameter of RunWorkflow.
Parameters:
- ctx: The DBOS client or context.
- name: The name of the queue. Must be unique among all queues in the application.
- options: Functional options for the queue, documented below.
Example Syntax:
queue, err := dbos.RegisterQueue(ctx, "email-queue",
dbos.WithWorkerConcurrency(5),
dbos.WithRateLimiter(&dbos.RateLimiter{
Limit: 100,
Period: 60 * time.Second, // 100 workflows per minute
}),
dbos.WithPriorityEnabled(),
)
// Enqueue workflows to this queue by passing its handle to WithQueue:
handle, err := dbos.RunWorkflow(ctx, SendEmailWorkflow, emailData, dbos.WithQueue(queue))
WithWorkerConcurrency
func WithWorkerConcurrency(concurrency int) QueueOption
Set the maximum number of workflows from this queue that may run concurrently within a single DBOS process.
WithGlobalConcurrency
func WithGlobalConcurrency(concurrency int) QueueOption
Set the maximum number of workflows from this queue that may run concurrently. Unset by default (no limit). This concurrency limit is global across all DBOS processes using this queue.
WithPriorityEnabled
func WithPriorityEnabled() QueueOption
Enable setting priority for workflows on this queue.
WithRateLimiter
func WithRateLimiter(limiter *RateLimiter) QueueOption
type RateLimiter struct {
Limit int // Maximum number of workflows to start within the period
Period time.Duration // Time period for the rate limit
}
A limit on the maximum number of functions which may be started in a given period.
WithPartitionConcurrency
func WithPartitionConcurrency(concurrency int) QueueOption
Set the maximum number of workflows from any one partition of this queue that may run concurrently across all DBOS processes. Must be at least 1 and less than or equal to the queue's global concurrency.
Setting any partition limit (WithPartitionConcurrency, WithPartitionWorkerConcurrency, or WithPartitionRateLimiter) makes the queue partitioned: every workflow enqueued on it must supply a partition key with WithQueuePartitionKey, and the queue dequeues from each partition separately.
A partitioned queue enforces its partition limits and its queue-wide limits (WithGlobalConcurrency, WithWorkerConcurrency, WithRateLimiter) at the same time.
WithPartitionWorkerConcurrency
func WithPartitionWorkerConcurrency(concurrency int) QueueOption
Set the maximum number of workflows from any one partition of this queue that may run concurrently within a single DBOS process. Must be at least 1 and less than or equal to the queue's partition concurrency, worker concurrency, and global concurrency. Setting this limit makes the queue partitioned.
WithPartitionRateLimiter
func WithPartitionRateLimiter(limiter *RateLimiter) QueueOption
A limit on the maximum number of workflows which may be started from any one partition of this queue in a given period. The limit is applied to each partition separately. Setting this limit makes the queue partitioned.
Example Syntax:
// Create a partitioned queue with a per-partition concurrency limit of 1
partitionedQueue, err := dbos.RegisterQueue(ctx, "user-tasks",
dbos.WithPartitionConcurrency(1),
)
// Enqueue workflows with different partition keys
// At most one workflow per user can run at once, but workflows from different users can run concurrently
handle1, _ := dbos.RunWorkflow(ctx, ProcessTask, task1,
dbos.WithQueue(partitionedQueue),
dbos.WithQueuePartitionKey("user-123"),
)
handle2, _ := dbos.RunWorkflow(ctx, ProcessTask, task2,
dbos.WithQueue(partitionedQueue),
dbos.WithQueuePartitionKey("user-456"),
)
WithPartitionQueue
func WithPartitionQueue() QueueOption
WithPartitionQueue is deprecated.
Partition a queue by setting a partition limit (WithPartitionConcurrency, WithPartitionWorkerConcurrency, or WithPartitionRateLimiter) instead.
Enable the legacy partitioned queue mode, under which the queue's global concurrency, worker concurrency, and rate limit each apply to individual partitions instead of the queue as a whole.
For example, a queue registered with WithPartitionQueue() and WithGlobalConcurrency(1) runs at most one workflow from each partition at a time.
The equivalent under the partition limits is WithPartitionConcurrency(1), which additionally lets you keep queue-wide limits (see Combining Queue-Wide and Per-Partition Limits).
WithPartitionQueue cannot be combined with a partition limit in the same RegisterQueue call.
A queue registered with it rejects the Set* methods for its limits with an error matching dbos.ErrInvalidOption: re-register the queue with the partition limits instead.
Its Get* methods report each limit at the scope it is enforced, so for example GetPartitionConcurrency returns the value passed to WithGlobalConcurrency and GetGlobalConcurrency returns nil.
WithQueueBasePollingInterval
func WithQueueBasePollingInterval(interval time.Duration) QueueOption
Set the base polling interval for this queue. This also acts as the minimum (fastest) interval. Polling intervals are subject to base 2 exponential backoff.
Example Syntax:
queue, err := dbos.RegisterQueue(ctx, "email-queue", dbos.WithQueueBasePollingInterval(100*time.Millisecond))
WithQueueApplicationName
func WithQueueApplicationName(name string) QueueOption
Set the application that owns the queue and dequeues workflows from it.
Defaults to the registering context's own application (for a standalone client, its AppName).
Registering a queue already owned by a different application returns an error.
WithQueueOnConflict
func WithQueueOnConflict(policy QueueConflictResolution) QueueOption
type QueueConflictResolution string
const (
QueueConflictUpdateIfLatestVersion QueueConflictResolution = "update_if_latest_version"
QueueConflictAlwaysUpdate QueueConflictResolution = "always_update"
QueueConflictNeverUpdate QueueConflictResolution = "never_update"
)
Set how RegisterQueue behaves when a queue with the same name already exists in the system database:
- QueueConflictUpdateIfLatestVersion (default): overwrite the existing configuration only if the running application is the latest registered application version. This prevents older versions in a rolling deploy from overwriting a newer configuration.
- QueueConflictAlwaysUpdate: always overwrite the existing configuration.
- QueueConflictNeverUpdate: leave the existing configuration unchanged. The returned queue reflects the persisted configuration, not the supplied options.
RetrieveQueue
func RetrieveQueue(ctx Client, name string) (Queue, error)
Retrieve a queue by name from the system database. If no queue with that name has been registered, returns an error matching dbos.ErrQueueNotFound:
q, err := dbos.RetrieveQueue(ctx, name)
if errors.Is(err, dbos.ErrQueueNotFound) { /* absent */ }
Example Syntax:
queue, err := dbos.RetrieveQueue(ctx, "email-queue")
if err != nil {
return err
}
fmt.Println("Priority enabled:", queue.GetPriorityEnabled())
ListQueues
func ListQueues(ctx Client, opts ...ListQueuesOption) ([]Queue, error)
Return queues registered in the system database.
By default, only queues owned by the calling context's application (plus queues owned by no application) are returned; a standalone client with no AppName returns every application's queues.
WithListQueuesApplicationNames
func WithListQueuesApplicationNames(names ...string) ListQueuesOption
List queues owned by these applications instead (queues owned by no application are always included).
DeleteQueue
func DeleteQueue(ctx Client, name string) error
Delete a queue from the system database. No-op if no queue with that name exists.
Workflows already enqueued on a deleted queue can no longer be dequeued, executed, or recovered. However, if a queue with the same name is later registered, it will dequeue the leftover workflows. Do not rely on this: stale workflows unexpectedly resuming on a future queue is rarely the intended behavior. Instead, cancel or drain pending workflows on the queue before deleting it.
Queue Interface
A Queue is returned from RegisterQueue, RetrieveQueue, and ListQueues.
Its Get* methods reflect the queue's configuration as of the most recent read from the database; the Set* methods update the configuration in the database.
GetPartitionQueue reports whether the queue is partitioned, whether by a partition limit or by the deprecated WithPartitionQueue option.
Unlike the other properties, ownership cannot be reconfigured: there is no SetApplicationName. Ownership is only transferred by RenameApplication.
type Queue interface {
GetName() string
GetGlobalConcurrency() *int
GetWorkerConcurrency() *int
GetRateLimit() *RateLimiter
GetPartitionConcurrency() *int
GetPartitionWorkerConcurrency() *int
GetPartitionRateLimit() *RateLimiter
GetPriorityEnabled() bool
GetPartitionQueue() bool
GetPollingInterval() time.Duration
GetApplicationName() string
SetGlobalConcurrency(ctx Client, value *int) error
SetWorkerConcurrency(ctx Client, value *int) error
SetRateLimit(ctx Client, value *RateLimiter) error
SetPartitionConcurrency(ctx Client, value *int) error
SetPartitionWorkerConcurrency(ctx Client, value *int) error
SetPartitionRateLimit(ctx Client, value *RateLimiter) error
SetPriorityEnabled(ctx Client, value bool) error
SetPartitionQueue(ctx Client, value bool) error
SetPollingInterval(ctx Client, value time.Duration) error
}
Reconfiguring Queues
Because queue configuration lives in the system database, you can change a queue's configuration at runtime without redeploying or restarting your workers.
Workers pick up the new configuration on their next polling iteration.
For the concurrency and rate limit setters, pass nil to clear the limit.
Each change is validated against the queue's latest persisted configuration: a concurrency limit must be greater than or equal to its partition counterpart, a worker limit must be less than or equal to its global counterpart, and partition limits must be at least 1.
Setting any partition limit (SetPartitionConcurrency, SetPartitionWorkerConcurrency, or SetPartitionRateLimit) makes the queue partitioned; clearing all of them makes it unpartitioned again.
SetPartitionQueue toggles the deprecated WithPartitionQueue mode, and cannot be used on a queue partitioned by its partition limits.
Take care when partitioning a queue at runtime: workflows already enqueued on it have no partition key and will not be dequeued until the queue is unpartitioned.
queue, err := dbos.RetrieveQueue(ctx, "email-queue")
if err != nil {
return err // dbos.ErrQueueNotFound if the queue does not exist
}
concurrency := 50
if err := queue.SetGlobalConcurrency(ctx, &concurrency); err != nil {
return err
}
if err := queue.SetRateLimit(ctx, &dbos.RateLimiter{Limit: 500, Period: 60 * time.Second}); err != nil {
return err
}
If your application calls RegisterQueue on startup, the next process to start can overwrite settings you applied at runtime via Set* methods.
Either update the RegisterQueue call to match the new configuration, or pass WithQueueOnConflict(dbos.QueueConflictNeverUpdate) to preserve the runtime changes.
ListenQueues
func ListenQueues(ctx Context, names ...string)
Configure which queues the current DBOS process should listen to for workflow execution.
By default, all registered queues are listened to.
When ListenQueues is called, only the specified queues (and the internal DBOS queue) will be processed by the queue runner.
This allows multiple DBOS processes to share the same queues but listen to different subsets.
A queue is identified by name, so a queue can be listened to even before it exists in the database; names are resolved against the database on each polling iteration.
dbos.RegisterQueue(ctx, "queue-1")
dbos.RegisterQueue(ctx, "queue-2")
// Only listen to queue-1 and queue-2.
dbos.ListenQueues(ctx, "queue-1", "queue-2")
Each call to ListenQueues replaces the whole listen set; passing an empty set listens to every queue. The listen set may be changed at any time, including after Launch(). Use ListenedQueues(ctx) to retrieve the current set.
Debouncer
A debouncer delays workflow execution until a configurable delay has elapsed since the last invocation. Each subsequent call pushes back the start time by the delay amount. This is useful when you want to coalesce rapid successive triggers (e.g., text field edits, sensor data) into a single workflow execution.
See the debouncing tutorial for usage examples.
NewDebouncer
func NewDebouncer[R any, P any](ctx Context, workflow Workflow[P, R], opts ...DebouncerOption) (*Debouncer[R, P], error)
Create a new debouncer for the specified workflow.
The workflow must be registered before creating the debouncer.
Debouncers can be created at any time, including after Launch().
Multiple debouncers can be created for the same workflow.
Both type parameters are inferred from the workflow function.
Parameters:
- ctx: The Context.
- workflow: The workflow function to debounce (must be registered).
- opts: Optional configuration, documented below.
WithDebouncerTimeout
func WithDebouncerTimeout(timeout time.Duration) DebouncerOption
Set the maximum time before starting the workflow, measured from the first debounce call for a given key. If the timeout is zero (the default), there is no maximum time limit and calling the workflow can be pushed back indefinitely.
WithDebouncerQueue
func WithDebouncerQueue(queueName string) DebouncerOption
Run the debounced workflow on the named queue instead of the DBOS internal queue.
Debounce keys are scoped to the queue.
The queue is fixed per debouncer and must be registered (see RegisterQueue); Debounce calls cannot override it.
NewDebouncer validates at creation time that the queue is registered; NewDebouncerClient does not.
WithDebouncerInstance
func WithDebouncerInstance(instance ConfiguredInstance) DebouncerOption
Target the workflow registration bound to the given configured instance (see WithInstance).
Required when the debounced workflow is a method of a configured instance.
debouncer, err := dbos.NewDebouncer(ctx, slack.Send, dbos.WithDebouncerInstance(slack))
Debouncer.Debounce
func (d *Debouncer[R, P]) Debounce(ctx Context, key string, delay time.Duration, input P, opts ...WorkflowOption) (WorkflowHandle[R], error)
Debounce a workflow invocation. If no debouncer is active for the given key, one is started with the specified delay. If a debouncer is already active for the key, the delay is pushed back and the input is updated. When the delay expires or the debouncer preconfigured timeout is reached, the target workflow is executed with the most recent input.
Parameters:
- ctx: The Context.
- key: A unique key to group debounce calls. Calls with the same key are debounced together.
- delay: Time by which to delay workflow execution from this call.
- input: Input parameters to pass to the workflow.
- opts: Optional workflow options (e.g.,
WithWorkflowID). Options the debounce owns or cannot support (WithQueue,WithDeduplicationID,WithDelay,WithPriority,WithQueuePartitionKey,WithDeduplicationPolicy) are rejected with an error matchingdbos.ErrInvalidOption.
Returns:
- A WorkflowHandle for the target workflow.
You can also create a debouncer from outside a DBOS application using a DebouncerClient.
NewDebouncerClient
func NewDebouncerClient[R any, P any](workflowName string, client Client, opts ...DebouncerOption) *DebouncerClient[R, P]
R is the workflow's result type and P its input type; neither can be inferred, so both must be named explicitly.
Create a new debouncer for use from outside a DBOS application.
Similar to NewDebouncer but uses a standalone client instead of a Context and takes a workflow name string instead of a function reference.
Parameters:
- workflowName: The name of the workflow to debounce.
- client: The DBOS client to use for operations.
- opts: Optional configuration — the same
DebouncerOptions asNewDebouncer, plus the client-specific options below.
WithDebouncerClassName
func WithDebouncerClassName(className string) DebouncerOption
Set the class/namespace name recorded for the debounced workflow.
Use with NewDebouncerClient when the target workflow is registered under a class name — for example by another language's runtime, which may resolve dequeued workflows by class name.
WithDebouncerConfigName
func WithDebouncerConfigName(configName string) DebouncerOption
Target the workflow registration bound to the configured instance with the given config name (see WithInstance).
Required when the debounced workflow is a method of a configured instance.
Use with NewDebouncerClient, where the instance object itself is not available (from a Context, use WithDebouncerInstance instead).
dc := dbos.NewDebouncerClient[string, string]("Send", client,
dbos.WithDebouncerConfigName("slack"))
DebouncerClient.Debounce
func (dc *DebouncerClient[R, P]) Debounce(key string, delay time.Duration, input P, opts ...WorkflowOption) (WorkflowHandle[R], error)
Debounce a workflow invocation from outside a DBOS application.
Behaves the same as Debouncer.Debounce but does not require a Context.
Parameters:
- key: A unique key to group debounce calls.
- delay: Time by which to delay workflow execution.
- input: Input parameters to pass to the workflow.
- opts: Optional workflow options, with the same restrictions as
Debouncer.Debounce.