How to Bind a Queue to a VarMQ Worker: Complete API Guide
Binding a queue to a VarMQ worker is accomplished through the binder API methods such as BindQueue(), BindPriorityQueue(), or WithDistributedQueue(), which instantiate the queue and automatically start the worker goroutine via defer wb.start().
VarMQ is a flexible, open-source job processing framework written in Go that strictly separates workers (the code that executes jobs) from queues (the data structures that store jobs). To establish this connection, you must bind a queue to a VarMQ worker using the binder API defined in worker_binder.go. This single operation creates the queue instance and immediately activates the worker's processing loop.
VarMQ Queue Binding Methods
The IWorkerBinder[T] interface exposed by every worker provides multiple methods for attaching different queue implementations. According to the source code in worker_binder.go, the available binding strategies include:
BindQueue– Creates and binds a standard FIFOQueue[T](lines 74‑78).WithQueue– Binds an existingIQueueimplementation without creating a new one (lines 80‑86).BindPriorityQueue– Instantiates aPriorityQueue[T]that orders jobs by numeric priority (lines 88‑90).WithPriorityQueue– Attaches a customIPriorityQueueimplementation (lines 92‑96).WithPersistentQueue– Connects a durablePersistentQueue[T]for job survival across restarts (lines 98‑101).WithDistributedQueue– Registers a Redis‑backedDistributedQueue[T]or anyIDistributedQueueimplementation (lines 110‑118).
When any binding method executes, the internal workerBinder[T] struct performs three critical actions: it instantiates or receives the concrete queue, invokes defer wb.start() to launch the worker goroutine, and returns a typed queue wrapper exposing Add, AddAll, and NumPending methods. For distributed queues, the binder additionally registers a subscription handler (handleQueueSubscription at lines 63‑71) that listens for "enqueued" events to trigger immediate job retrieval.
Binding Examples by Queue Type
Standard FIFO Queue Binding
The most common scenario involves creating a worker and binding it to a standard in-memory queue.
package main
import (
"fmt"
"github.com/goptics/varmq"
)
func main() {
// Create a worker that processes strings
worker := varmq.NewWorker(func(j varmq.Job[string]) {
fmt.Println("processing:", j.Data())
})
// Bind to a FIFO queue; this starts the worker automatically
queue := worker.BindQueue()
// Enqueue jobs (non-blocking)
queue.Add("hello world")
}
Source reference: worker_binder.go lines 74‑78 implement BindQueue by calling newQueue and starting the worker.
Priority Queue Binding
For jobs requiring urgency levels, bind a priority queue that sorts by numeric priority values.
package main
import (
"fmt"
"github.com/goptics/varmq"
)
func main() {
worker := varmq.NewWorker(func(j varmq.Job[string]) {
fmt.Println("priority job:", j.Data())
})
// Bind to a priority queue
pQueue := worker.BindPriorityQueue()
// Higher numbers indicate higher priority
pQueue.Add("low priority", 1)
pQueue.Add("urgent task", 10)
}
Source reference: worker_binder.go lines 88‑90 define BindPriorityQueue which instantiates a PriorityQueue[T].
Persistent Queue Binding
To ensure jobs survive process crashes, bind a PersistentQueue implementation (typically Redis-backed).
package main
import (
"log"
"github.com/goptics/varmq"
)
func main() {
// Assume redisQueue implements IPersistentQueue[string]
redisQueue := // ... initialize your persistent queue
worker := varmq.NewWorker(func(j varmq.Job[string]) {
log.Println("persisted:", j.Data())
})
// Bind existing persistent queue instance
pQueue := worker.WithPersistentQueue(redisQueue)
pQueue.Add("critical task")
}
Source reference: worker_binder.go lines 98‑101 show WithPersistentQueue accepting an IPersistentQueue[T] and wrapping it with the worker binder.
Distributed Queue Binding
For clustered deployments where multiple worker nodes share a single job source, bind a distributed queue.
package main
import (
"log"
"github.com/goptics/varmq"
)
func main() {
// Assume distQueue implements IDistributedQueue[string]
distQueue := // ... initialize Redis distributed queue
worker := varmq.NewWorker(func(j varmq.Job[string]) {
log.Println("distributed:", j.Data())
})
// Bind distributed queue with subscription handling
dQueue := worker.WithDistributedQueue(distQueue)
dQueue.Add("cluster-wide task")
}
Source reference: worker_binder.go lines 110‑118 implement WithDistributedQueue, which registers handleQueueSubscription to listen for remote enqueue events.
Result Worker Queue Binding
Result workers return values to callers and bind to queues identically.
package main
import (
"fmt"
"github.com/goptics/varmq"
)
func main() {
// Worker returns string length
worker := varmq.NewResultWorker(func(j varmq.Job[string]) (int, error) {
return len(j.Data()), nil
})
q := worker.BindQueue()
if handle, ok := q.Add("hello"); ok {
go func() {
if length, err := handle.Result(); err == nil {
fmt.Println("length =", length) // Output: 5
}
}()
}
}
Source reference: The pattern follows examples/result-worker/main.go and the binding implementation in worker_binder.go.
How Queue Binding Works Under the Hood
The binding mechanism relies on the generic workerBinder[T] struct defined in worker_binder.go. When you invoke any Bind* or With* method, the binder:
- Instantiates the concrete queue type (e.g.,
Queue[T],PriorityQueue[T]) or accepts an existing interface implementation. - Registers the queue with the internal
worker[T,J]instance, which maintains a registry of active queues. - Starts the worker goroutine via
defer wb.start(), ensuring the worker begins polling for jobs immediately upon binding completion. - Returns a typed wrapper (e.g.,
newQueue,newPriorityQueue) that implements the public queue API while exposingqueue.Worker()for introspection.
The internal helper newQueues (and its variants newResultQueues, newErrQueues) creates the workerBinder instance when you call varmq.NewWorker. For distributed scenarios, step 2 includes registering the handleQueueSubscription callback (lines 63‑71) that reacts to "enqueued" actions from remote publishers, triggering immediate job retrieval rather than waiting for the next poll interval.
Summary
Binding a queue to a VarMQ worker unites job storage with job processing through a fluent, type-safe API. Key takeaways include:
- Use
BindQueue()for standard FIFO processing andBindPriorityQueue()for priority-based ordering. - Use
WithQueue()orWithPriorityQueue()to attach custom implementations ofIQueueorIPriorityQueue. - Use
WithPersistentQueue()andWithDistributedQueue()for durable or clustered architectures. - Binding automatically starts the worker via
defer wb.start(); no manual startup is required. - Access the underlying worker through the returned queue's
Worker()method if you need to inspect or manage the processor.
Frequently Asked Questions
What is the difference between BindQueue and WithQueue?
BindQueue creates a new standard Queue[T] instance internally and binds it to your worker, while WithQueue accepts an existing object that implements the IQueue interface. Use WithQueue when you have already instantiated a custom queue or need to share a queue instance across multiple workers.
Does binding a queue immediately start the worker?
Yes. According to the source code in worker_binder.go, every binding method includes defer wb.start(), which launches the worker's internal goroutine to poll for jobs. The worker begins processing as soon as the binding function returns.
Can I bind multiple different queue types to a single worker?
The VarMQ architecture supports registering multiple queues with a single worker instance. You can call multiple binding methods (e.g., BindQueue() for local jobs and WithDistributedQueue() for remote jobs) on the same worker object, and the worker[T,J] will poll all registered queues for available jobs.
How do I bind a queue to a ResultWorker or ErrWorker?
The binding API is identical for all worker types. Whether you create a worker with varmq.NewWorker, varmq.NewResultWorker, or varmq.NewErrWorker, the returned object implements IWorkerBinder[T] and exposes the same BindQueue(), WithQueue(), and other binding methods defined in worker_binder.go.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →