| name | cloudflare-queues |
| description | Complete knowledge domain for Cloudflare Queues - flexible message queue for asynchronous processing
and background tasks on Cloudflare Workers.
Use when: creating message queues, async processing, background jobs, batch processing, handling retries,
configuring dead letter queues, implementing consumer concurrency, or encountering "queue timeout",
"batch retry", "message lost", "throughput exceeded", "consumer not scaling" errors.
Keywords: cloudflare queues, queues workers, message queue, queue bindings, async processing,
background jobs, queue consumer, queue producer, batch processing, dead letter queue, dlq,
message retry, queue ack, consumer concurrency, queue backlog, wrangler queues
|
| license | MIT |
Cloudflare Queues
Status: Production Ready ✅
Last Updated: 2025-10-21
Dependencies: cloudflare-worker-base (for Worker setup)
Latest Versions: wrangler@4.43.0, @cloudflare/workers-types@4.20251014.0
Quick Start (10 Minutes)
1. Create a Queue
npx wrangler queues create my-queue
npx wrangler queues list
npx wrangler queues info my-queue
2. Set Up Producer (Send Messages)
wrangler.jsonc:
{
"name": "my-producer",
"main": "src/index.ts",
"compatibility_date": "2025-10-11",
"queues": {
"producers": [
{
"binding": "MY_QUEUE",
"queue": "my-queue"
}
]
}
}
src/index.ts (Producer):
import { Hono } from 'hono';
type Bindings = {
MY_QUEUE: Queue;
};
const app = new Hono<{ Bindings: Bindings }>();
app.post('/send', async (c) => {
const body = await c.req.json();
await c.env.MY_QUEUE.send({
userId: body.userId,
action: 'process-order',
timestamp: Date.now(),
});
return c.json({ status: 'queued' });
});
app.post('/send-batch', async (c) => {
const items = await c.req.json();
await c.env.MY_QUEUE.sendBatch(
items.map((item) => ({
body: { userId: item.userId, action: item.action },
}))
);
return c.json({ status: 'queued', count: items.length });
});
export default app;
3. Set Up Consumer (Process Messages)
Create consumer Worker:
npm create cloudflare@latest my-consumer -- --type hello-world --ts
cd my-consumer
wrangler.jsonc:
{
"name": "my-consumer",
"main": "src/index.ts",
"compatibility_date": "2025-10-11",
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_batch_size": 10,
"max_batch_timeout": 5
}
]
}
}
src/index.ts (Consumer):
export default {
async queue(
batch: MessageBatch,
env: Env,
ctx: ExecutionContext
): Promise<void> {
console.log(`Processing batch of ${batch.messages.length} messages`);
for (const message of batch.messages) {
console.log('Message:', message.id, message.body, `Attempt: ${message.attempts}`);
await processMessage(message.body);
}
},
};
async function processMessage(body: any) {
console.log('Processing:', body);
}
4. Deploy and Test
cd my-producer
npm run deploy
cd my-consumer
npm run deploy
curl -X POST https://my-producer.<your-subdomain>.workers.dev/send \
-H "Content-Type: application/json" \
-d '{"userId": "123", "action": "welcome-email"}'
npx wrangler tail my-consumer
Complete Producer API
send() - Send Single Message
interface QueueSendOptions {
delaySeconds?: number;
}
await env.MY_QUEUE.send(body: any, options?: QueueSendOptions);
Examples:
await env.MY_QUEUE.send({ userId: '123', action: 'send-email' });
await env.MY_QUEUE.send(
{ userId: '123', action: 'reminder' },
{ delaySeconds: 600 }
);
await env.MY_QUEUE.send({
type: 'order-confirmation',
orderId: 'ORD-123',
email: 'user@example.com',
items: [{ sku: 'ITEM-1', quantity: 2 }],
total: 49.99,
timestamp: Date.now(),
});
CRITICAL:
- Message body must be JSON serializable (structured clone algorithm)
- Maximum message size: 128 KB (including ~100 bytes internal metadata)
- Messages >128 KB will fail - split them or store in R2 and send reference
sendBatch() - Send Multiple Messages
interface MessageSendRequest<Body = any> {
body: Body;
delaySeconds?: number;
}
interface QueueSendBatchOptions {
delaySeconds?: number;
}
await env.MY_QUEUE.sendBatch(
messages: Iterable<MessageSendRequest>,
options?: QueueSendBatchOptions
);
Examples:
await env.MY_QUEUE.sendBatch([
{ body: { userId: '1', action: 'email' } },
{ body: { userId: '2', action: 'email' } },
{ body: { userId: '3', action: 'email' } },
]);
await env.MY_QUEUE.sendBatch([
{ body: { task: 'task1' }, delaySeconds: 60 },
{ body: { task: 'task2' }, delaySeconds: 300 },
{ body: { task: 'task3' }, delaySeconds: 600 },
]);
await env.MY_QUEUE.sendBatch(
[
{ body: { task: 'task1' } },
{ body: { task: 'task2' } },
],
{ delaySeconds: 3600 }
);
const tasks = await getTasks();
await env.MY_QUEUE.sendBatch(
tasks.map((task) => ({
body: {
taskId: task.id,
userId: task.userId,
priority: task.priority,
},
}))
);
Limits:
- Maximum 100 messages per batch
- Maximum 256 KB total batch size
- Each message still limited to 128 KB individually
Complete Consumer API
Queue Handler Function
export default {
async queue(
batch: MessageBatch,
env: Env,
ctx: ExecutionContext
): Promise<void> {
},
};
Parameters:
batch - MessageBatch object containing messages
env - Environment bindings (KV, D1, R2, etc.)
ctx - Execution context for waitUntil(), passThroughOnException()
MessageBatch Interface
interface MessageBatch<Body = unknown> {
readonly queue: string;
readonly messages: Message<Body>[];
ackAll(): void;
retryAll(options?: QueueRetryOptions): void;
}
Properties:
Methods:
Message Interface
interface Message<Body = unknown> {
readonly id: string;
readonly timestamp: Date;
readonly body: Body;
readonly attempts: number;
ack(): void;
retry(options?: QueueRetryOptions): void;
}
Properties:
id - System-generated unique ID (UUID)
timestamp - Date object when message was sent to queue
body - Your message content (any JSON serializable type)
attempts - Number of times consumer has processed this message
- Starts at 1 on first delivery
- Increments on each retry
- Use for exponential backoff:
delaySeconds: 60 * message.attempts
Methods:
QueueRetryOptions
interface QueueRetryOptions {
delaySeconds?: number;
}
Example:
message.retry();
message.retry({ delaySeconds: 300 });
message.retry({
delaySeconds: Math.min(60 * Math.pow(2, message.attempts - 1), 3600),
});
Consumer Patterns
1. Basic Consumer (Implicit Acknowledgement)
Best for: Idempotent operations where retries are safe
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const message of batch.messages) {
await sendEmail(message.body.email, message.body.content);
}
},
};
Behavior:
- If function returns successfully → all messages acknowledged
- If function throws error → all messages retried
- Simple but can cause duplicate processing on partial failures
2. Explicit Acknowledgement (Non-Idempotent Operations)
Best for: Database writes, API calls, financial transactions
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const message of batch.messages) {
try {
await env.DB.prepare(
'INSERT INTO orders (id, user_id, amount) VALUES (?, ?, ?)'
).bind(message.body.orderId, message.body.userId, message.body.amount).run();
message.ack();
} catch (error) {
console.error(`Failed to process ${message.id}:`, error);
}
}
},
};
Why explicit ack?
- Prevents duplicate writes if one message in batch fails
- Only successfully processed messages are acknowledged
- Failed messages retry independently
3. Retry with Exponential Backoff
Best for: Rate-limited APIs, temporary failures
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const message of batch.messages) {
try {
await fetch('https://api.example.com/process', {
method: 'POST',
body: JSON.stringify(message.body),
});
message.ack();
} catch (error) {
if (error.status === 429) {
const delaySeconds = Math.min(
60 * Math.pow(2, message.attempts - 1),
3600
);
console.log(`Rate limited. Retrying in ${delaySeconds}s (attempt ${message.attempts})`);
message.retry({ delaySeconds });
} else {
message.retry();
}
}
}
},
};
4. Dead Letter Queue (DLQ) Pattern
Best for: Handling permanently failed messages
Setup DLQ:
npx wrangler queues create my-dlq
wrangler.jsonc:
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_batch_size": 10,
"max_retries": 3,
"dead_letter_queue": "my-dlq"
}
]
}
}
DLQ Consumer:
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const message of batch.messages) {
console.error('PERMANENTLY FAILED MESSAGE:', {
id: message.id,
attempts: message.attempts,
body: message.body,
timestamp: message.timestamp,
});
await env.DB.prepare(
'INSERT INTO failed_messages (id, body, attempts, failed_at) VALUES (?, ?, ?, ?)'
).bind(
message.id,
JSON.stringify(message.body),
message.attempts,
new Date().toISOString()
).run();
await sendAlert({
type: 'queue-dlq',
messageId: message.id,
queue: batch.queue,
});
message.ack();
}
},
};
How it works:
- Message fails in main queue
- Retries up to
max_retries (default 3)
- After max retries, sent to DLQ
- DLQ consumer processes failed messages
- Without DLQ, messages are deleted permanently
5. Multiple Queues, Single Consumer
Best for: Centralized processing logic
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
switch (batch.queue) {
case 'high-priority-queue':
await processHighPriority(batch.messages, env);
break;
case 'low-priority-queue':
await processLowPriority(batch.messages, env);
break;
case 'email-queue':
await processEmails(batch.messages, env);
break;
default:
console.warn(`Unknown queue: ${batch.queue}`);
}
},
};
async function processHighPriority(messages: Message[], env: Env) {
for (const message of messages) {
await fastProcess(message.body);
message.ack();
}
}
async function processLowPriority(messages: Message[], env: Env) {
for (const message of messages) {
await slowProcess(message.body);
message.ack();
}
}
wrangler.jsonc:
{
"queues": {
"consumers": [
{ "queue": "high-priority-queue" },
{ "queue": "low-priority-queue" },
{ "queue": "email-queue" }
]
}
}
Consumer Configuration
Batch Settings
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_batch_size": 100,
"max_batch_timeout": 30
}
]
}
}
How batching works:
- Consumer called when either condition met:
max_batch_size messages accumulated
max_batch_timeout seconds elapsed (whichever comes first)
Example:
max_batch_size: 100, max_batch_timeout: 10
- If 100 messages arrive in 3 seconds → batch delivered immediately
- If only 50 messages arrive in 10 seconds → batch of 50 delivered
Tuning guidelines:
- High volume, low latency →
max_batch_size: 100, max_batch_timeout: 1
- Low volume, batch writes →
max_batch_size: 50, max_batch_timeout: 30
- Cost optimization → Larger batches = fewer invocations
Retry Settings
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_retries": 5,
"retry_delay": 300
}
]
}
}
max_retries:
- Number of times to retry failed message
- After max retries:
- With DLQ → message sent to DLQ
- Without DLQ → message deleted permanently
retry_delay:
- Default delay for all retried messages (seconds)
- Can be overridden with
message.retry({ delaySeconds })
- Maximum delay: 43200 seconds (12 hours)
Concurrency Settings
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_concurrency": 10
}
]
}
}
How concurrency works:
- Queues auto-scales consumers based on backlog
- Default: Scales up to 250 concurrent invocations
- Setting
max_concurrency limits scaling
When to set max_concurrency:
- ✅ Upstream API has rate limits
- ✅ Database connection limits
- ✅ Want to control costs
- ❌ Most cases - leave unset for best performance
Auto-scaling triggers:
- Growing backlog (messages accumulating)
- High error rate
- Processing speed vs. incoming rate
Dead Letter Queue
{
"queues": {
"consumers": [
{
"queue": "my-queue",
"max_retries": 3,
"dead_letter_queue": "my-dlq"
}
]
}
}
CRITICAL:
- DLQ must be created separately:
npx wrangler queues create my-dlq
- Without DLQ, failed messages are deleted permanently
- Messages in DLQ persist for 4 days without consumer
- Always configure DLQ for production queues
Wrangler Commands
Create Queue
npx wrangler queues create my-queue
npx wrangler queues create my-queue --message-retention-period-secs 1209600
npx wrangler queues create my-queue --delivery-delay-secs 60
List Queues
npx wrangler queues list
Get Queue Info
npx wrangler queues info my-queue
Update Queue
npx wrangler queues update my-queue --message-retention-period-secs 604800
npx wrangler queues update my-queue --delivery-delay-secs 3600
Delete Queue
npx wrangler queues delete my-queue
Consumer Management
npx wrangler queues consumer add my-queue my-consumer-worker \
--batch-size 50 \
--batch-timeout 10 \
--message-retries 5 \
--max-concurrency 20 \
--retry-delay-secs 300
npx wrangler queues consumer remove my-queue my-consumer-worker
Purge Queue
npx wrangler queues purge my-queue
Pause/Resume Delivery
npx wrangler queues pause-delivery my-queue
npx wrangler queues resume-delivery my-queue
Use cases:
- Maintenance on consumer Workers
- Temporarily stop processing
- Debug issues without losing messages
Limits & Quotas
| Feature | Limit |
|---|
| Queues per account | 10,000 |
| Message size | 128 KB (includes ~100 bytes metadata) |
| Message retries | 100 max |
| Batch size | 1-100 messages |
| Batch timeout | 0-60 seconds |
| Messages per sendBatch | 100 (or 256 KB total) |
| Queue throughput | 5,000 messages/second per queue |
| Message retention | 4 days (default), 14 days (max) |
| Queue backlog size | 25 GB per queue |
| Concurrent consumers | 250 (push-based, auto-scale) |
| Consumer duration | 15 minutes (wall clock) |
| Consumer CPU time | 30 seconds (default), 5 minutes (max) |
| Visibility timeout | 12 hours (pull consumers) |
| Message delay | 12 hours (max) |
| API rate limit | 1200 requests / 5 minutes |
Pricing
Requires Workers Paid plan ($5/month)
Operations Pricing:
- First 1,000,000 operations/month: FREE
- After that: $0.40 per million operations
What counts as an operation:
- Each 64 KB chunk written, read, or deleted
- Messages >64 KB count as multiple operations:
- 65 KB message = 2 operations
- 127 KB message = 2 operations
- 128 KB message = 2 operations
Typical message lifecycle:
- 1 write + 1 read + 1 delete = 3 operations
Retries:
- Each retry = additional read operation
- Message retried 3 times = 1 write + 4 reads + 1 delete = 6 operations
Dead Letter Queue:
- Writing to DLQ = additional write operation
Cost examples:
- 1M messages/month (no retries): ((1M × 3) - 1M) / 1M × $0.40 = $0.80
- 10M messages/month: ((10M × 3) - 1M) / 1M × $0.40 = $11.60
- 100M messages/month: ((100M × 3) - 1M) / 1M × $0.40 = $119.60
TypeScript Types
interface Queue<Body = any> {
send(body: Body, options?: QueueSendOptions): Promise<void>;
sendBatch(
messages: Iterable<MessageSendRequest<Body>>,
options?: QueueSendBatchOptions
): Promise<void>;
}
interface QueueSendOptions {
delaySeconds?: number;
}
interface MessageSendRequest<Body = any> {
body: Body;
delaySeconds?: number;
}
interface QueueSendBatchOptions {
delaySeconds?: number;
}
export default {
queue(
batch: MessageBatch,
env: Env,
ctx: ExecutionContext
): Promise<void>;
}
interface MessageBatch<Body = unknown> {
readonly queue: string;
readonly messages: Message<Body>[];
ackAll(): void;
retryAll(options?: QueueRetryOptions): void;
}
interface Message<Body = unknown> {
readonly id: string;
readonly timestamp: Date;
readonly body: Body;
readonly attempts: number;
ack(): void;
retry(options?: QueueRetryOptions): void;
}
interface QueueRetryOptions {
delaySeconds?: number;
}
Error Handling
Common Errors
1. Message Too Large
await env.MY_QUEUE.send({
data: largeArray,
});
const message = { data: largeArray };
const size = new TextEncoder().encode(JSON.stringify(message)).length;
if (size > 128000) {
const key = `messages/${crypto.randomUUID()}.json`;
await env.MY_BUCKET.put(key, JSON.stringify(message));
await env.MY_QUEUE.send({ type: 'large-message', r2Key: key });
} else {
await env.MY_QUEUE.send(message);
}
2. Throughput Exceeded
for (let i = 0; i < 10000; i++) {
await env.MY_QUEUE.send({ id: i });
}
const messages = Array.from({ length: 10000 }, (_, i) => ({
body: { id: i },
}));
for (let i = 0; i < messages.length; i += 100) {
await env.MY_QUEUE.sendBatch(messages.slice(i, i + 100));
}
for (let i = 0; i < messages.length; i += 100) {
await env.MY_QUEUE.sendBatch(messages.slice(i, i + 100));
if (i + 100 < messages.length) {
await new Promise(resolve => setTimeout(resolve, 100));
}
}
3. Consumer Timeout
export default {
async queue(batch: MessageBatch): Promise<void> {
for (const message of batch.messages) {
await processForMinutes(message.body);
}
},
};
wrangler.jsonc:
{
"limits": {
"cpu_ms": 300000
}
}
4. Backlog Growing
{
"queues": {
"consumers": [{
"queue": "my-queue",
"max_batch_size": 100
}]
}
}
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
await Promise.all(
batch.messages.map(async (message) => {
await process(message.body);
message.ack();
})
);
},
};
Always Do ✅
- Configure Dead Letter Queue for production queues
- Use explicit ack() for non-idempotent operations (DB writes, API calls)
- Validate message size before sending (<128 KB)
- Use sendBatch() for multiple messages (more efficient)
- Implement exponential backoff for retries
- Set appropriate batch settings based on workload
- Monitor queue backlog and consumer errors
- Use ctx.waitUntil() for async cleanup in consumers
- Handle errors gracefully - log, alert, retry
- Let concurrency auto-scale (don't set max_concurrency unless needed)
Never Do ❌
- Never assume message ordering - not guaranteed FIFO
- Never rely on implicit ack for non-idempotent ops - use explicit ack()
- Never send messages >128 KB - will fail
- Never delete queues with active messages - data loss
- Never skip DLQ configuration in production
- Never exceed 5000 msg/s per queue - rate limit error
- Never process messages synchronously in loop - use Promise.all()
- Never ignore message.attempts - use for backoff logic
- Never set max_concurrency=1 unless you have a very specific reason
- Never forget to ack() in explicit acknowledgement patterns
Troubleshooting
Issue: Messages not being delivered to consumer
Possible causes:
- Consumer not deployed
- Wrong queue name in wrangler.jsonc
- Delivery paused
- Consumer throwing errors
Solution:
npx wrangler queues info my-queue
npx wrangler queues resume-delivery my-queue
npx wrangler tail my-consumer
Issue: Entire batch retried when one message fails
Cause: Using implicit acknowledgement with non-idempotent operations
Solution: Use explicit ack()
for (const message of batch.messages) {
try {
await dbWrite(message.body);
message.ack();
} catch (error) {
console.error(`Failed: ${message.id}`);
}
}
Issue: Messages deleted without processing
Cause: No Dead Letter Queue configured
Solution:
npx wrangler queues create my-dlq
{
"queues": {
"consumers": [{
"queue": "my-queue",
"dead_letter_queue": "my-dlq"
}]
}
}
Issue: Consumer not auto-scaling
Possible causes:
max_concurrency set to 1
- Consumer returning errors (not processing)
- Batch processing too fast (no backlog)
Solution:
{
"queues": {
"consumers": [{
"queue": "my-queue",
"max_batch_size": 50
}]
}
}
Production Checklist
Before deploying to production:
Related Documentation
Last Updated: 2025-10-21
Version: 1.0.0
Maintainer: Jeremy Dawes | jeremy@jezweb.net