| name | cloudflare-queues |
| description | Build async message queues with Cloudflare Queues for background processing. Use when: handling async tasks, batch processing, implementing retries, configuring dead letter queues, managing consumer concurrency, or troubleshooting queue timeout, batch retry, message loss, or throughput exceeded.
|
| 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