Skip to content
Get started
Get started
QueueGroup: Namespace Related Queues

import { Tabs, TabItem } from ‘@astrojs/starlight/components’;

guide · queue-group

Many queues, one namespace.

QueueGroup prefixes a set of queues with a shared name, so "invoices" inside the "billing" group becomes "billing:invoices". Handy for multi-tenant apps and keeping domains apart.

A QueueGroup is a thin organizer: it creates normal Queue and Worker instances whose names carry the group prefix, and it can pause, resume, or clear all of them at once.

import { QueueGroup } from 'bunqueue/client';
const billing = new QueueGroup('billing');
// Queues are automatically prefixed
const invoices = billing.getQueue('invoices', { embedded: true }); // "billing:invoices"
const payments = billing.getQueue('payments', { embedded: true }); // "billing:payments"
await invoices.add('create', { amount: 100 });
await payments.add('process', { orderId: '123' });
// Workers use the same prefixed names
const invoiceWorker = billing.getWorker('invoices', async (job) => {
console.log('Processing invoice:', job.data);
return { processed: true };
}, { embedded: true });
import { Queue, Worker } from 'bunqueue-client';
const invoices = new Queue('billing:invoices');
const payments = new Queue('billing:payments');
await invoices.add('create', { amount: 100 });
await payments.add('process', { orderId: '123' });
const invoiceWorker = new Worker('billing:invoices', async (job) => ({
processed: true,
invoice: job.data,
}));
from bunqueue import Queue, Worker
invoices = Queue("billing:invoices")
payments = Queue("billing:payments")
invoices.add("create", {"amount": 100})
payments.add("process", {"order_id": "123"})
invoice_worker = Worker(
"billing:invoices",
lambda job: {"processed": True, "invoice": job.data},
)
invoice_worker.run() # blocking loop
use Bunqueue\Queue;
use Bunqueue\Worker;
$invoices = new Queue('billing:invoices');
$payments = new Queue('billing:payments');
$invoices->add('create', ['amount' => 100]);
$payments->add('process', ['orderId' => '123']);
$invoiceWorker = new Worker('billing:invoices',
fn (Bunqueue\Job $job) => ['processed' => true, 'invoice' => $job->data()]
);
$invoiceWorker->run(); // blocking loop
invoices := bunqueue.NewQueue("billing:invoices", bunqueue.Options{})
payments := bunqueue.NewQueue("billing:payments", bunqueue.Options{})
invoices.Add("create", map[string]any{"amount": 100}, nil)
payments.Add("process", map[string]any{"orderId": "123"}, nil)
invoiceWorker := bunqueue.NewWorker("billing:invoices", processor,
bunqueue.WorkerOptions{})
invoiceWorker.Run() // blocking loop
use bunqueue_client::{ConnectionOptions, JobOptions, Queue, Value, Worker, WorkerOptions};
let options = ConnectionOptions::default();
let invoices = Queue::new("billing:invoices", options.clone());
let payments = Queue::new("billing:payments", options);
invoices.add("create", Value::from(100), JobOptions::default())?;
payments.add("process", Value::from("123"), JobOptions::default())?;
let invoice_worker = Worker::new(
"billing:invoices", processor, WorkerOptions::default(),
);
invoice_worker.run()?; // blocking loop
invoices = Bunqueue.queue("billing:invoices")
payments = Bunqueue.queue("billing:payments")
{:ok, _job} = Bunqueue.Queue.add(invoices, "create", %{amount: 100})
{:ok, _job} = Bunqueue.Queue.add(payments, "process", %{order_id: "123"})
invoice_worker =
Bunqueue.Worker.new("billing:invoices", fn job ->
{:ok, %{processed: true, invoice: job.data}}
end)
Bunqueue.Worker.run(invoice_worker) # blocking loop

In Bun, getQueue and getWorker accept the same options as Queue and Worker, including embedded mode. The network SDK examples create normal queues and workers with the same prefixed broker names, which is the portable equivalent.

billing.listQueues(); // ['invoices', 'payments'] (names without prefix)
billing.pauseAll(); // pause every queue in the group
billing.resumeAll(); // resume them
billing.drainAll(); // remove all waiting jobs
billing.obliterateAll(); // remove ALL data from every queue
// Awaitable forms are authoritative in embedded and TCP modes
await billing.pauseAllAsync();
await billing.resumeAllAsync();
const removed = await billing.drainAllAsync();
await billing.obliterateAllAsync();
const queues = [invoices, payments];
await Promise.all(queues.map((queue) => queue.pause()));
await Promise.all(queues.map((queue) => queue.resume()));
const removed = await Promise.all(queues.map((queue) => queue.drain()));
await Promise.all(queues.map((queue) => queue.obliterate()));
queues = [invoices, payments]
for queue in queues:
queue.pause()
for queue in queues:
queue.resume()
removed = [queue.drain() for queue in queues]
$queues = [$invoices, $payments];
foreach ($queues as $queue) {
$queue->pause();
}
foreach ($queues as $queue) {
$queue->resume();
}
$removed = array_map(fn ($queue) => $queue->drain(), $queues);
queues := []*bunqueue.Queue{invoices, payments}
for _, queue := range queues {
if err := queue.Pause(); err != nil { return err }
}
for _, queue := range queues {
if err := queue.Resume(); err != nil { return err }
}
let queues = [&invoices, &payments];
for queue in queues {
queue.pause()?;
}
for queue in queues {
queue.resume()?;
}
queues = [invoices, payments]
Enum.each(queues, fn queue -> :ok = Bunqueue.Queue.pause(queue) end)
Enum.each(queues, fn queue -> :ok = Bunqueue.Queue.resume(queue) end)
{:ok, removed} =
Enum.reduce_while(queues, {:ok, []}, fn queue, {:ok, counts} ->
case Bunqueue.Queue.drain(queue) do
{:ok, count} -> {:cont, {:ok, [count | counts]}}
{:error, error} -> {:halt, {:error, error}}
end
end)
const tenantA = new QueueGroup('tenant-a');
const tenantB = new QueueGroup('tenant-b');
const tasksA = tenantA.getQueue('tasks', { embedded: true }); // "tenant-a:tasks"
const tasksB = tenantB.getQueue('tasks', { embedded: true }); // "tenant-b:tasks"
const tasksA = new Queue('tenant-a:tasks');
const tasksB = new Queue('tenant-b:tasks');
tasks_a = Queue("tenant-a:tasks")
tasks_b = Queue("tenant-b:tasks")
$tasksA = new Queue('tenant-a:tasks');
$tasksB = new Queue('tenant-b:tasks');
tasksA := bunqueue.NewQueue("tenant-a:tasks", bunqueue.Options{})
tasksB := bunqueue.NewQueue("tenant-b:tasks", bunqueue.Options{})
let tasks_a = Queue::new("tenant-a:tasks", ConnectionOptions::default());
let tasks_b = Queue::new("tenant-b:tasks", ConnectionOptions::default());
tasks_a = Bunqueue.queue("tenant-a:tasks")
tasks_b = Bunqueue.queue("tenant-b:tasks")
const env = process.env.NODE_ENV || 'development';
const group = new QueueGroup(`${env}-tasks`);
const queue = group.getQueue('jobs', { embedded: true });
// "development-tasks:jobs" or "production-tasks:jobs"
const env = process.env.NODE_ENV || 'development';
const queue = new Queue(`${env}-tasks:jobs`);
env = os.getenv("APP_ENV", "development")
queue = Queue(f"{env}-tasks:jobs")
$env = getenv('APP_ENV') ?: 'development';
$queue = new Queue("{$env}-tasks:jobs");
env := os.Getenv("APP_ENV")
if env == "" { env = "development" }
queue := bunqueue.NewQueue(env+"-tasks:jobs", bunqueue.Options{})
let env = std::env::var("APP_ENV").unwrap_or_else(|_| "development".into());
let queue = Queue::new(format!("{env}-tasks:jobs"), ConnectionOptions::default());
env = System.get_env("APP_ENV", "development")
queue = Bunqueue.queue("#{env}-tasks:jobs")
MethodDescription
getQueue(name, opts?)Get a queue within the group (embedded or TCP)
getWorker(name, processor, opts?)Create a worker for a queue in the group (embedded or TCP)
listQueues()List queue names in the group, without prefix (embedded only)
pauseAll()Pause all queues in the group (embedded only)
resumeAll()Resume all queues in the group (embedded only)
drainAll()Remove waiting jobs from all queues (embedded only)
obliterateAll()Remove all data from all queues (embedded only)
listQueuesAsync()List tracked group queues in either runtime
pauseAllAsync() / resumeAllAsync()Await group control in either runtime
drainAllAsync()Drain all tracked queues and return the aggregate count
obliterateAllAsync()Await removal of all tracked queue data