The use cases teams run.
Six patterns you can copy into a real app: emails, webhooks, images, payments, cron and multi-step flows. Each one moves slow work out of your API request and into a background job.
This page shows the most common things people build with bunqueue, each with a small working example. If a term is new to you, it gets a plain-words explanation the first time it appears.
The core pattern
Section titled “The core pattern”Every use case below is a variation of the same idea: your API handler adds a job (a unit of work saved as data) to a queue and returns immediately. A worker (a function that pulls jobs and runs them) does the slow part in the background.
import { Queue, Worker } from 'bunqueue/client';
// Embedded mode: the queue runs inside your process, no separate server.// dataPath persists jobs to a SQLite file so they survive restarts.const queue = new Queue('emails', { embedded: true, dataPath: './data/app.db' });
new Worker('emails', async (job) => { await sendEmail(job.data); return { sent: true };}, { embedded: true, concurrency: 10 }); // up to 10 jobs in parallel
// In your API handler: this returns in microsecondsawait queue.add('welcome', { to: 'user@example.com' });import { Queue, Worker } from 'bunqueue-client';
// Connects to a bunqueue server on localhost:6789const queue = new Queue('emails');
new Worker('emails', async (job) => { await sendEmail(job.data); return { sent: true };}, { concurrency: 10 }); // up to 10 jobs in parallel
// In your API handler: this returns as soon as the job is queuedawait queue.add('welcome', { to: 'user@example.com' });from bunqueue import Queue, Worker
# Connects to a bunqueue server on localhost:6789queue = Queue("emails")
def process(job): send_email(job.data) return {"sent": True}
Worker("emails", process, concurrency=10) # up to 10 jobs in parallel
# In your API handler: this returns as soon as the job is queuedqueue.add("welcome", {"to": "user@example.com"})use Bunqueue\Queue;use Bunqueue\Worker;
// Connects to a bunqueue server on localhost:6789$queue = new Queue('emails');
// In your API handler: this returns as soon as the job is queued$queue->add('welcome', ['to' => 'user@example.com']);
// worker.php, a separate long-running process$worker = new Worker('emails', function (Bunqueue\Job $job) { sendEmail($job->data()); return ['sent' => true];});$worker->run();// Connects to a bunqueue server on localhost:6789queue := bunqueue.NewQueue("emails", bunqueue.Options{})defer queue.Close()
// In your API handler: this returns as soon as the job is queuedqueue.Add("welcome", map[string]any{"to": "user@example.com"}, nil)
// worker, a separate long-running processworker := bunqueue.NewWorker("emails", func(job *bunqueue.Job) (any, error) { return sendEmail(job.Data())}, bunqueue.WorkerOptions{Concurrency: 10}) // up to 10 jobs in parallelworker.Run()use bunqueue_client::{ConnectionOptions, JobOptions, ProcessError, Queue, Value, Worker, WorkerOptions};
// Connects to a bunqueue server on localhost:6789let queue = Queue::new("emails", ConnectionOptions::default());
// In your API handler: this returns as soon as the job is queuedlet data = Value::Map(vec![(Value::from("to"), Value::from("user@example.com"))]);queue.add("welcome", data, JobOptions::default())?;
// worker, a separate long-running processlet worker = Worker::new( "emails", |job| { deliver(job.data()) .map(|_| Value::from(true)) .map_err(|error| ProcessError::retryable(error.to_string())) }, WorkerOptions { concurrency: 10, ..Default::default() }, // up to 10 jobs in parallel);worker.run()?;# Connects to a bunqueue server on localhost:6789queue = Bunqueue.queue("emails")
# In your API handler: this returns as soon as the job is queued{:ok, _job} = Bunqueue.Queue.add(queue, "welcome", %{to: "user@example.com"})
# worker, a separate long-running processworker = Bunqueue.Worker.new("emails", fn job -> send_email(job.data) {:ok, %{sent: true}} end, concurrency: 10)
Bunqueue.Worker.run(worker)That is the whole model. The sections below add the options that make each use case reliable. New to bunqueue? Start with the quickstart.
Email delivery
Section titled “Email delivery”Sending email inside an API request is slow and fragile: the provider can be down or rate limited. Queue it instead, and let bunqueue retry on failure.
import { Queue, Worker } from 'bunqueue/client';
interface EmailJob { to: string; template: string; data: Record<string, unknown> }
const emails = new Queue<EmailJob>('emails', { embedded: true, dataPath: './data/app.db', defaultJobOptions: { attempts: 5, // try up to 5 times backoff: 2000, // wait 2s, then 4s, 8s... between tries (exponential backoff) removeOnComplete: true, // drop finished jobs to keep the queue lean },});
new Worker<EmailJob>('emails', async (job) => { const result = await sendEmail(job.data); return { messageId: result.messageId };}, { embedded: true, concurrency: 10 });
// One emailawait emails.add('welcome', { to: 'user@example.com', template: 'welcome', data: { name: 'John' } });
// Bulk newsletter, batched in one callawait emails.addBulk( subscribers.map((s) => ({ name: 'newsletter', data: { to: s.email, template: 'news', data: {} } })));import { Queue, Worker } from 'bunqueue-client';
interface EmailJob { to: string; template: string; data: Record<string, unknown> }
const emails = new Queue<EmailJob>('emails');
// Retry options are passed per add (defaultJobOptions is a Bun bunqueue feature)const retryOpts = { attempts: 5, // try up to 5 times backoff: 2000, // wait 2s, then 4s, 8s... between tries (exponential backoff) removeOnComplete: true, // drop finished jobs to keep the queue lean};
new Worker<EmailJob>('emails', async (job) => { const result = await sendEmail(job.data); return { messageId: result.messageId };}, { concurrency: 10 });
// One emailawait emails.add('welcome', { to: 'user@example.com', template: 'welcome', data: { name: 'John' } }, retryOpts);
// Bulk newsletter, batched in one callawait emails.addBulk( subscribers.map((s) => ({ name: 'newsletter', data: { to: s.email, template: 'news', data: {} }, opts: retryOpts })));from bunqueue import Queue, Worker
emails = Queue("emails")
# Retry options are passed per addretry_opts = { "attempts": 5, # try up to 5 times "backoff": 2000, # wait 2s, then 4s, 8s... between tries (exponential backoff) "remove_on_complete": True, # drop finished jobs to keep the queue lean}
def process(job): result = send_email(job.data) return {"message_id": result["message_id"]}
Worker("emails", process, concurrency=10)
# One emailemails.add("welcome", {"to": "user@example.com", "template": "welcome", "data": {"name": "John"}}, **retry_opts)
# Bulk newsletter, batched in one callemails.add_bulk([ {"name": "newsletter", "data": {"to": s["email"], "template": "news", "data": {}}, **retry_opts} for s in subscribers])emails := bunqueue.NewQueue("emails", bunqueue.Options{})defer emails.Close()
// Retry options are passed per addretryOpts := bunqueue.JobOptions{ "attempts": 5, // try up to 5 times "backoff": 2000, // wait 2s, then 4s, 8s... between tries (exponential backoff) "removeOnComplete": true, // drop finished jobs to keep the queue lean}
worker := bunqueue.NewWorker("emails", func(job *bunqueue.Job) (any, error) { result, err := sendEmail(job.Data()) if err != nil { return nil, err // retried automatically } return map[string]any{"messageId": result.MessageID}, nil}, bunqueue.WorkerOptions{Concurrency: 10})go worker.Run()
// One emailemails.Add("welcome", map[string]any{"to": "user@example.com", "template": "welcome", "data": map[string]any{"name": "John"}}, retryOpts)
// Bulk newsletter, batched in one callentries := make([]bunqueue.BulkEntry, 0, len(subscribers))for _, s := range subscribers { entries = append(entries, bunqueue.BulkEntry{ Name: "newsletter", Data: map[string]any{"to": s.Email, "template": "news"}, Opts: retryOpts, })}emails.AddBulk(entries)PHP, Rust and Elixir follow the same shape, add / addBulk on the producer and a worker with retries handled server-side, see the SDK guide.
If all 5 attempts fail, the job lands in the dead letter queue (DLQ), a holding area for jobs that ran out of retries. Inspect it with queue.getDlq() and retry with queue.retryDlq(). See the DLQ guide.
Webhook delivery
Section titled “Webhook delivery”Partner endpoints go down and return 5xx errors. Treat every delivery as a job with retries, and let the DLQ auto-retry the stubborn ones on a schedule.
const webhooks = new Queue('webhooks', { embedded: true, dataPath: './data/app.db', defaultJobOptions: { attempts: 8, backoff: 5000 }, // 5s, 10s, 20s, 40s...});
// Jobs that exhaust all 8 attempts go to the DLQ.// Auto-retry the DLQ every hour, up to 3 times, then keep entries 7 days.webhooks.setDlqConfig({ autoRetry: true, autoRetryInterval: 3_600_000, maxAutoRetries: 3, maxAge: 604_800_000,});
new Worker('webhooks', async (job) => { const { endpoint, event, payload } = job.data; const res = await fetch(endpoint, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Webhook-Event': event }, body: JSON.stringify(payload), signal: AbortSignal.timeout(30_000), // never hang on a dead endpoint }); if (!res.ok) throw new Error(`HTTP ${res.status}`); // throwing triggers a retry return { status: res.status };}, { embedded: true, concurrency: 20 });
await webhooks.add('order.created', { endpoint: 'https://partner.com/webhooks', event: 'order.created', payload: { orderId: 'ORD-123' },});import { Queue, Worker } from 'bunqueue-client';
const webhooks = new Queue('webhooks');
// Jobs that exhaust all 8 attempts go to the DLQ.// Auto-retry the DLQ every hour, up to 3 times, then keep entries 7 days.await webhooks.setDlqConfig({ autoRetry: true, autoRetryInterval: 3_600_000, maxAutoRetries: 3, maxAge: 604_800_000,});
new Worker('webhooks', async (job) => { const { endpoint, event, payload } = job.data; const res = await fetch(endpoint, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Webhook-Event': event }, body: JSON.stringify(payload), signal: AbortSignal.timeout(30_000), // never hang on a dead endpoint }); if (!res.ok) throw new Error(`HTTP ${res.status}`); // throwing triggers a retry return { status: res.status };}, { concurrency: 20 });
await webhooks.add('order.created', { endpoint: 'https://partner.com/webhooks', event: 'order.created', payload: { orderId: 'ORD-123' },}, { attempts: 8, backoff: 5000 }); // 5s, 10s, 20s, 40s...import requestsfrom bunqueue import Queue, Worker
webhooks = Queue("webhooks")
# Jobs that exhaust all 8 attempts go to the DLQ.# Auto-retry the DLQ every hour, up to 3 times, then keep entries 7 days.webhooks.set_dlq_config({ "autoRetry": True, "autoRetryInterval": 3_600_000, "maxAutoRetries": 3, "maxAge": 604_800_000,})
def deliver(job): res = requests.post( job.data["endpoint"], json=job.data["payload"], headers={"X-Webhook-Event": job.data["event"]}, timeout=30, # never hang on a dead endpoint ) res.raise_for_status() # raising triggers a retry return {"status": res.status_code}
Worker("webhooks", deliver, concurrency=20)
webhooks.add("order.created", { "endpoint": "https://partner.com/webhooks", "event": "order.created", "payload": {"orderId": "ORD-123"},}, attempts=8, backoff=5000) # 5s, 10s, 20s, 40s...webhooks := bunqueue.NewQueue("webhooks", bunqueue.Options{})defer webhooks.Close()
// DLQ auto-retry configuration (setDlqConfig) is available from the// Bun, TypeScript and Python clients.
worker := bunqueue.NewWorker("webhooks", func(job *bunqueue.Job) (any, error) { data := job.Data() payload, _ := json.Marshal(data["payload"]) ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() // never hang on a dead endpoint
req, _ := http.NewRequestWithContext(ctx, "POST", data["endpoint"].(string), bytes.NewReader(payload)) req.Header.Set("Content-Type", "application/json") req.Header.Set("X-Webhook-Event", data["event"].(string)) res, err := http.DefaultClient.Do(req) if err != nil { return nil, err // returning an error triggers a retry } defer res.Body.Close() if res.StatusCode >= 400 { return nil, fmt.Errorf("HTTP %d", res.StatusCode) } return map[string]any{"status": res.StatusCode}, nil}, bunqueue.WorkerOptions{Concurrency: 20})go worker.Run()
webhooks.Add("order.created", map[string]any{ "endpoint": "https://partner.com/webhooks", "event": "order.created", "payload": map[string]any{"orderId": "ORD-123"},}, bunqueue.JobOptions{"attempts": 8, "backoff": 5000}) // 5s, 10s, 20s, 40s...PHP, Rust and Elixir deliver webhooks with the same worker shape (see the SDK guide). DLQ auto-retry configuration (setDlqConfig) is available from the Bun, TypeScript and Python clients.
Image processing
Section titled “Image processing”Generating thumbnails and variants during an upload request makes the upload slow. Queue one job per image, report progress as each variant finishes, and cap the runtime with a timeout.
const images = new Queue('images', { embedded: true, dataPath: './data/app.db', defaultJobOptions: { attempts: 3, timeout: 120_000 }, // kill stuck jobs after 2 minutes});
new Worker('images', async (job) => { const { sourceUrl, variants } = job.data; const source = await downloadImage(sourceUrl); const urls: Record<string, string> = {};
for (let i = 0; i < variants.length; i++) { const v = variants[i]; // updateProgress(percent, message): your frontend can poll or subscribe to this await job.updateProgress(Math.round((i / variants.length) * 100), `Processing ${v.name}`); const out = await sharp(source).resize(v.width, v.height).webp().toBuffer(); urls[v.name] = await uploadToCDN(out, `${job.id}/${v.name}.webp`); }
await job.updateProgress(100, 'Done'); return { urls };}, { embedded: true, concurrency: 5 });
await images.add('product-image', { sourceUrl: 'https://uploads.example.com/raw/product-123.jpg', variants: [ { name: 'thumb', width: 150, height: 150 }, { name: 'full', width: 1200, height: 900 }, ],});import { Queue, Worker } from 'bunqueue-client';
const images = new Queue('images');
new Worker('images', async (job) => { const { sourceUrl, variants } = job.data; const source = await downloadImage(sourceUrl); const urls: Record<string, string> = {};
for (let i = 0; i < variants.length; i++) { const v = variants[i]; // updateProgress(percent, message): your frontend can poll this await job.updateProgress(Math.round((i / variants.length) * 100), `Processing ${v.name}`); const out = await sharp(source).resize(v.width, v.height).webp().toBuffer(); urls[v.name] = await uploadToCDN(out, `${job.id}/${v.name}.webp`); }
await job.updateProgress(100, 'Done'); return { urls };}, { concurrency: 5 });
await images.add('product-image', { sourceUrl: 'https://uploads.example.com/raw/product-123.jpg', variants: [ { name: 'thumb', width: 150, height: 150 }, { name: 'full', width: 1200, height: 900 }, ],}, { attempts: 3, timeout: 120_000 }); // kill stuck jobs after 2 minutesfrom bunqueue import Queue, Worker
images = Queue("images")
def process(job): variants = job.data["variants"] source = download_image(job.data["sourceUrl"]) urls = {}
for i, v in enumerate(variants): # update_progress(percent, message): your frontend can poll this job.update_progress(round(i / len(variants) * 100), f"Processing {v['name']}") out = resize_image(source, v["width"], v["height"]) # e.g. Pillow urls[v["name"]] = upload_to_cdn(out, f"{job.id}/{v['name']}.webp")
job.update_progress(100, "Done") return {"urls": urls}
Worker("images", process, concurrency=5)
images.add("product-image", { "sourceUrl": "https://uploads.example.com/raw/product-123.jpg", "variants": [ {"name": "thumb", "width": 150, "height": 150}, {"name": "full", "width": 1200, "height": 900}, ],}, attempts=3, timeout=120_000) # kill stuck jobs after 2 minutesimages := bunqueue.NewQueue("images", bunqueue.Options{})defer images.Close()
worker := bunqueue.NewWorker("images", func(job *bunqueue.Job) (any, error) { data := job.Data() variants := data["variants"].([]any) source, err := downloadImage(data["sourceUrl"].(string)) if err != nil { return nil, err } urls := map[string]string{}
for i, raw := range variants { v := raw.(map[string]any) name := v["name"].(string) // UpdateProgress(percent, message): your frontend can poll this job.UpdateProgress(float64(i)/float64(len(variants))*100, "Processing "+name) out := resizeVariant(source, v) // your image library of choice urls[name] = uploadToCDN(out, job.ID()+"/"+name+".webp") }
job.UpdateProgress(100, "Done") return map[string]any{"urls": urls}, nil}, bunqueue.WorkerOptions{Concurrency: 5})go worker.Run()
images.Add("product-image", map[string]any{ "sourceUrl": "https://uploads.example.com/raw/product-123.jpg", "variants": []any{ map[string]any{"name": "thumb", "width": 150, "height": 150}, map[string]any{"name": "full", "width": 1200, "height": 900}, },}, bunqueue.JobOptions{"attempts": 3, "timeout": 120000}) // kill stuck jobs after 2 minutesPHP, Rust and Elixir report progress the same way (updateProgress / update_progress), see the Worker guide.
The same shape works for video transcoding or any long CPU-bound task, just raise the timeout. For heavy CPU work see CPU-intensive workers.
Payments and other critical jobs
Section titled “Payments and other critical jobs”A payment job must never be lost and never run twice. Two options handle this:
durable: trueskips the write buffer (the small in-memory batch bunqueue normally flushes to disk every 10ms) and writes the job to SQLite beforeadd()returns. A crash cannot lose it.jobIdmakes the add idempotent: adding the samejobIdtwice returns the existing job instead of creating a duplicate.
const payments = new Queue('payments', { embedded: true, dataPath: './data/payments.db', // durable writes need a SQLite file defaultJobOptions: { attempts: 3, backoff: 5000, timeout: 60_000 },});
// Failed payments need a human, never an automatic re-chargepayments.setDlqConfig({ autoRetry: false, maxAge: 2_592_000_000 }); // keep 30 days
new Worker('payments', async (job) => { const { orderId, amount, idempotencyKey } = job.data; const intent = await stripe.paymentIntents.create( { amount, currency: 'usd', confirm: true }, { idempotencyKey } // provider-side guard against double charges ); if (intent.status !== 'succeeded') throw new Error(`Payment failed: ${intent.status}`); await recordTransaction(orderId, intent.id); return { paymentIntentId: intent.id };}, { embedded: true, concurrency: 5 });
await payments.add( 'charge', { orderId: 'ORD-123', amount: 9999, idempotencyKey: 'order-ORD-123' }, { jobId: 'charge-ORD-123', durable: true } // no duplicate, no loss on crash);import { Queue, Worker } from 'bunqueue-client';
const payments = new Queue('payments');
// Failed payments need a human, never an automatic re-chargeawait payments.setDlqConfig({ autoRetry: false, maxAge: 2_592_000_000 }); // keep 30 days
new Worker('payments', async (job) => { const { orderId, amount, idempotencyKey } = job.data; const intent = await stripe.paymentIntents.create( { amount, currency: 'usd', confirm: true }, { idempotencyKey } // provider-side guard against double charges ); if (intent.status !== 'succeeded') throw new Error(`Payment failed: ${intent.status}`); await recordTransaction(orderId, intent.id); return { paymentIntentId: intent.id };}, { concurrency: 5 });
await payments.add( 'charge', { orderId: 'ORD-123', amount: 9999, idempotencyKey: 'order-ORD-123' }, { jobId: 'charge-ORD-123', durable: true, // no duplicate, no loss on crash attempts: 3, backoff: 5000, timeout: 60_000, });from bunqueue import Queue, Worker
payments = Queue("payments")
# Failed payments need a human, never an automatic re-chargepayments.set_dlq_config({"autoRetry": False, "maxAge": 2_592_000_000}) # keep 30 days
def charge(job): intent = stripe.PaymentIntent.create( amount=job.data["amount"], currency="usd", confirm=True, idempotency_key=job.data["idempotencyKey"], # guard against double charges ) if intent.status != "succeeded": raise Exception(f"Payment failed: {intent.status}") record_transaction(job.data["orderId"], intent.id) return {"payment_intent_id": intent.id}
Worker("payments", charge, concurrency=5)
payments.add( "charge", {"orderId": "ORD-123", "amount": 9999, "idempotencyKey": "order-ORD-123"}, job_id="charge-ORD-123", durable=True, # no duplicate, no loss on crash attempts=3, backoff=5000, timeout=60_000,)payments := bunqueue.NewQueue("payments", bunqueue.Options{})defer payments.Close()
worker := bunqueue.NewWorker("payments", func(job *bunqueue.Job) (any, error) { data := job.Data() // Pass data["idempotencyKey"] to your provider as its idempotency key: // a provider-side guard against double charges. intent, err := chargeProvider(data) if err != nil { return nil, err } if err := recordTransaction(data["orderId"].(string), intent.ID); err != nil { return nil, err } return map[string]any{"paymentIntentId": intent.ID}, nil}, bunqueue.WorkerOptions{Concurrency: 5})go worker.Run()
payments.Add("charge", map[string]any{"orderId": "ORD-123", "amount": 9999, "idempotencyKey": "order-ORD-123"}, bunqueue.JobOptions{ "jobId": "charge-ORD-123", // no duplicate "durable": true, // no loss on crash "attempts": 3, "backoff": 5000, "timeout": 60000, })PHP, Rust and Elixir support jobId and durable identically, see the SDK guide. DLQ auto-retry configuration (setDlqConfig) is available from the Bun, TypeScript and Python clients.
Scheduled tasks (cron)
Section titled “Scheduled tasks (cron)”Recurring jobs use upsertJobScheduler(). Schedules are stored in SQLite and survive restarts, in both embedded and server mode. A cron expression like 0 3 * * * means “every day at 3 AM”.
const scheduled = new Queue('scheduled', { embedded: true, dataPath: './data/app.db' });
// Daily cleanup at 3 AMawait scheduled.upsertJobScheduler( 'daily-cleanup', { pattern: '0 3 * * *' }, { data: { task: 'cleanup' } });
// Health check every 5 minutesawait scheduled.upsertJobScheduler( 'health-check', { every: 300_000 }, { data: { task: 'health-check' } });
new Worker('scheduled', async (job) => { if (job.data.task === 'cleanup') return { deleted: await cleanupOldRecords() }; return await checkSystemHealth();}, { embedded: true });import { Queue, Worker } from 'bunqueue-client';
const scheduled = new Queue('scheduled');
// Daily cleanup at 3 AMawait scheduled.upsertJobScheduler( 'daily-cleanup', { pattern: '0 3 * * *' }, { data: { task: 'cleanup' } });
// Health check every 5 minutesawait scheduled.upsertJobScheduler( 'health-check', { every: 300_000 }, { data: { task: 'health-check' } });
new Worker('scheduled', async (job) => { if (job.data.task === 'cleanup') return { deleted: await cleanupOldRecords() }; return await checkSystemHealth();});from bunqueue import Queue, Worker
scheduled = Queue("scheduled")
# Daily cleanup at 3 AMscheduled.upsert_job_scheduler("daily-cleanup", {"pattern": "0 3 * * *"}, {"data": {"task": "cleanup"}})
# Health check every 5 minutesscheduled.upsert_job_scheduler("health-check", {"every": 300_000}, {"data": {"task": "health-check"}})
def process(job): if job.data["task"] == "cleanup": return {"deleted": cleanup_old_records()} return check_system_health()
Worker("scheduled", process)$scheduled = new Bunqueue\Queue('scheduled');
// Daily cleanup at 3 AM$scheduled->upsertJobScheduler('daily-cleanup', ['pattern' => '0 3 * * *'], ['data' => ['task' => 'cleanup']],);
// Health check every 5 minutes$scheduled->upsertJobScheduler('health-check', ['every' => 300000], ['data' => ['task' => 'health-check']],);
$worker = new Bunqueue\Worker('scheduled', function (Bunqueue\Job $job) { if ($job->data()['task'] === 'cleanup') { return ['deleted' => cleanupOldRecords()]; } return checkSystemHealth();});$worker->run();scheduled := bunqueue.NewQueue("scheduled", bunqueue.Options{})defer scheduled.Close()
// Daily cleanup at 3 AMscheduled.UpsertJobScheduler("daily-cleanup", bunqueue.SchedulerRepeat{Pattern: "0 3 * * *"}, bunqueue.SchedulerTemplate{Data: map[string]any{"task": "cleanup"}},)
// Health check every 5 minutesscheduled.UpsertJobScheduler("health-check", bunqueue.SchedulerRepeat{EveryMs: 300000}, bunqueue.SchedulerTemplate{Data: map[string]any{"task": "health-check"}},)
worker := bunqueue.NewWorker("scheduled", func(job *bunqueue.Job) (any, error) { if job.Data()["task"] == "cleanup" { return map[string]any{"deleted": cleanupOldRecords()}, nil } return checkSystemHealth()}, bunqueue.WorkerOptions{})worker.Run()use bunqueue_client::{ConnectionOptions, Queue, SchedulerRepeat, SchedulerTemplate, Value};
let scheduled = Queue::new("scheduled", ConnectionOptions::default());
// Daily cleanup at 3 AMscheduled.upsert_job_scheduler( "daily-cleanup", SchedulerRepeat { pattern: Some("0 3 * * *".into()), ..Default::default() }, SchedulerTemplate { data: Value::Map(vec![(Value::from("task"), Value::from("cleanup"))]), ..Default::default() },)?;
// Health check every 5 minutesscheduled.upsert_job_scheduler( "health-check", SchedulerRepeat { every_ms: Some(300_000), ..Default::default() }, SchedulerTemplate { data: Value::Map(vec![(Value::from("task"), Value::from("health-check"))]), ..Default::default() },)?;scheduled = Bunqueue.queue("scheduled")
# Daily cleanup at 3 AM:ok = Bunqueue.Queue.upsert_scheduler(scheduled, "daily-cleanup", %{pattern: "0 3 * * *"}, %{data: %{task: "cleanup"}} )
# Health check every 5 minutes:ok = Bunqueue.Queue.upsert_scheduler(scheduled, "health-check", %{every: 300_000}, %{data: %{task: "health-check"}} )
worker = Bunqueue.Worker.new("scheduled", fn job -> case job.data["task"] do "cleanup" -> {:ok, %{deleted: cleanup_old_records()}} _ -> {:ok, check_system_health()} end end)
Bunqueue.Worker.run(worker)Timezones, one-off delayed jobs and the repeat shorthand on queue.add() are covered in the cron guide.
Multi-step flows
Section titled “Multi-step flows”Some work has dependencies: an order ships only after inventory and payment both check out. FlowProducer runs child jobs first, in parallel, then runs the parent with access to every child result.
import { FlowProducer, Worker } from 'bunqueue/client';
type OrderData = { orderId: string };const flow = new FlowProducer({ embedded: true });
const checks = new Worker<OrderData>('checks', async (job) => ({ check: job.name, approved: true,}), { embedded: true });
const orders = new Worker<OrderData>('orders', async (job) => { // Children finished first; read what each one returned const results = await job.getChildrenValues(); return { orderId: job.data.orderId, shipped: true, checks: results };}, { embedded: true });
const node = await flow.add<OrderData>({ name: 'fulfill-order', queueName: 'orders', data: { orderId: 'ORD-123' }, children: [ { name: 'check-inventory', queueName: 'checks', data: { orderId: 'ORD-123' } }, { name: 'check-payment', queueName: 'checks', data: { orderId: 'ORD-123' } }, ],});
const result = await node.job.waitUntilFinished(null, 10_000);console.log(result);
await checks.close();await orders.close();await flow.close();import { FlowProducer, Worker } from 'bunqueue-client';
const flow = new FlowProducer();
await flow.add({ name: 'fulfill-order', queueName: 'orders', data: { orderId: 'ORD-123' }, children: [ { name: 'check-inventory', queueName: 'checks', data: { orderId: 'ORD-123' } }, { name: 'check-payment', queueName: 'checks', data: { orderId: 'ORD-123' } }, ],});
new Worker('orders', async (job) => { // Children finished first; read what each one returned const results = await job.getChildrenValues(); return shipOrder(job.data.orderId, results);});from bunqueue import FlowProducer, Worker
flow = FlowProducer()
flow.add({ "name": "fulfill-order", "queueName": "orders", "data": {"orderId": "ORD-123"}, "children": [ {"name": "check-inventory", "queueName": "checks", "data": {"orderId": "ORD-123"}}, {"name": "check-payment", "queueName": "checks", "data": {"orderId": "ORD-123"}}, ],})
def fulfill(job): # Children finished first; read what each one returned results = job.get_children_values() return ship_order(job.data["orderId"], results)
Worker("orders", fulfill)use Bunqueue\FlowProducer;use Bunqueue\Queue;use Bunqueue\Worker;
$flow = new FlowProducer();
$flow->add([ 'name' => 'fulfill-order', 'queueName' => 'orders', 'data' => ['orderId' => 'ORD-123'], 'children' => [ ['name' => 'check-inventory', 'queueName' => 'checks', 'data' => ['orderId' => 'ORD-123']], ['name' => 'check-payment', 'queueName' => 'checks', 'data' => ['orderId' => 'ORD-123']], ],]);
$orders = new Queue('orders');$worker = new Worker('orders', function (Bunqueue\Job $job) use ($orders) { // Children finished first; read what each one returned $results = $orders->getChildrenValues($job->id()); return shipOrder($job->data()['orderId'], $results);});$worker->run();flow := bunqueue.NewFlowProducer(bunqueue.Options{})defer flow.Close()
flow.Add(bunqueue.FlowJob{ Name: "fulfill-order", QueueName: "orders", Data: map[string]any{"orderId": "ORD-123"}, Children: []bunqueue.FlowJob{ {Name: "check-inventory", QueueName: "checks", Data: map[string]any{"orderId": "ORD-123"}}, {Name: "check-payment", QueueName: "checks", Data: map[string]any{"orderId": "ORD-123"}}, },})
orders := bunqueue.NewQueue("orders", bunqueue.Options{})worker := bunqueue.NewWorker("orders", func(job *bunqueue.Job) (any, error) { // Children finished first; read what each one returned results, err := orders.GetChildrenValues(job.ID()) if err != nil { return nil, err } return shipOrder(job.Data()["orderId"].(string), results)}, bunqueue.WorkerOptions{})worker.Run()Rust and Elixir support the same flow trees (FlowProducer::add / Bunqueue.FlowProducer.add) and sequential chains, see the SDK guide.
flow.addChain([...]) runs jobs one after another, and flow.addBulkThen(jobs, finalJob) fans out in parallel and merges at the end (fan-in is available in the TypeScript and Python clients). See the flow guide. For richer orchestration (branching, rollback on failure, human approval steps) use the workflow engine.
Job options cheat sheet
Section titled “Job options cheat sheet”The options you saw above, all set per job or via defaultJobOptions:
| Option | What it does | Default |
|---|---|---|
attempts | Max tries before the job goes to the DLQ | 3 |
backoff | Base delay between retries, doubles each time | 1000 ms |
timeout | Max processing time before the job is failed | none |
priority | Higher numbers run sooner | 0 |
delay | Wait this many ms before the job is runnable | 0 |
jobId | Custom ID, adding the same ID twice returns the existing job | auto |
durable | Write to disk immediately instead of the 10ms buffer | false |
removeOnComplete | Delete the job once it succeeds | false |
Full list in the queue guide.
Gotchas
Section titled “Gotchas”- No
dataPathmeans no persistence. An embedded queue withoutdataPath(or aDATA_PATHenv var) keeps everything in memory and loses it on restart. - The default write buffer trades 10ms for speed. Jobs are flushed to SQLite every 10ms. A hard crash can lose jobs accepted in that window. Use
durable: truewhere that matters. - Throwing is how you retry. A worker that catches every error and returns normally marks the job completed. Let errors propagate when you want a retry.
- Do not auto-retry money. Set
setDlqConfig({ autoRetry: false })on payment-like queues so failed charges wait for review. - Close workers on shutdown.
await worker.close()waits for active jobs to finish;worker.close(true)forces a stop. See production for the full shutdown pattern.
More patterns
Section titled “More patterns”- AI agents via MCP, let Claude or any MCP client schedule and monitor jobs
- Multi-tenant isolation, one namespaced queue set per tenant
- Rate-limited API calls, token buckets and worker limiters
- Edge and IoT forwarding, queue locally, drain to a central server
- Copy-paste examples, shorter recipes for common tasks