Skip to content

Queue

Durable FIFO buffer with acknowledgments. Essential for load leveling and reliable background processing — video transcoding, email sending, order processing.

Basic Usage

typescript
// Create queue
const mailQ = await client.queue<MailJob>("emails").create();

// Push message
await mailQ.push({ to: "test@test.com" });

// Subscribe (auto-ACK on success)
await mailQ.subscribe((msg) => console.log(msg));

// Delete queue
await mailQ.delete();
python
# Create queue
mail_q: NexoQueue[MailJob] = await client.queue("emails").create()

# Push message
await mail_q.push({"to": "test@test.com"})

# Subscribe (auto-ACK on success)
async def handle_email(msg: MailJob) -> None:
    print(msg)

await mail_q.subscribe(handle_email)

# Delete queue
await mail_q.delete()

Persistence

All queues are persisted to disk by default using a Write-Ahead Log (WAL) backed by SQLite. To maximize throughput and performance, Nexo uses an asynchronous flush strategy for all queues. Writes are buffered in memory and flushed to disk periodically.

By default, the server flushes data to disk every 100ms. This interval is globally configurable via the QUEUE_DEFAULT_FLUSH_MS environment variable (see Configuration below).

Advanced Creation

Configure reliability and timeout settings:

typescript
const criticalQueue = await client.queue<CriticalTask>('critical-tasks').create({
  // RELIABILITY
  visibilityTimeoutMs: 10000,  // Retry if not ACKed within 10s (default: 30s)
  maxDeliveries: 5,            // Move to DLQ after 5 failed deliveries (default: 5)
});
python
critical_queue: NexoQueue[CriticalTask] = await client.queue("critical-tasks").create({
    # RELIABILITY
    "visibility_timeout_ms": 10000,  # Retry if not ACKed within 10s (default: 30s)
    "max_deliveries": 5,             # Move to DLQ after 5 failed deliveries (default: 5)
})

Priority

typescript
// PRIORITY: Higher value = delivered first (0-255)
await criticalQueue.push({ type: 'urgent' }, { priority: 255 });
python
# PRIORITY: Higher value = delivered first (0-255)
await critical_queue.push({"type": "urgent"}, {"priority": 255})

Batch Push

Push multiple messages in a single network request. Reduces round-trip overhead and improves throughput when producing bursts of messages.

typescript
await mailQ.pushBatch([
  { data: { to: 'user1@example.com' } },
  { data: { to: 'user2@example.com' } },
  { data: { to: 'user3@example.com' }, options: { priority: 10 } },
]);
python
await mail_q.push_batch([
    {"data": {"to": "user1@example.com"}},
    {"data": {"to": "user2@example.com"}},
    {"data": {"to": "user3@example.com"}, "options": {"priority": 10}},
])

Each item can have its own priority. The server processes all items atomically under a single lock, then notifies consumers once.

Consumer Tuning

Queues are pull-based: the SDK continuously polls the server for new messages in a loop, processes them, and polls again. The server never pushes messages to the client. Three parameters control this behavior:

batchSize (default: 50)

How many messages the SDK fetches from the server in a single network request. Higher values reduce round-trips but increase memory usage per cycle.

Choosing batchSize vs concurrency

The SDK fetches batchSize messages per request and processes them with concurrency parallel workers. If batchSize is much larger than concurrency, some messages may sit in the client's local buffer waiting for a worker. The server marks each delivered message with a visibility timeout (set at queue creation via visibilityTimeoutMs). If a message waits too long in the buffer, its timeout expires and the server redelivers it to another consumer — resulting in duplicate processing.

This is safe (Nexo guarantees at-least-once delivery), but wasteful. As a rule of thumb: if your callbacks are fast, batchSize can be much larger than concurrency. If your callbacks are slow or your visibility timeout is short, keep batchSize close to concurrency.

waitMs (default: 20000)

When the queue is empty, the server holds the connection open for up to waitMs milliseconds before responding with an empty result (long-polling). This avoids the client hammering the server with tight empty loops. If a message arrives during the wait, the server responds immediately.

concurrency (default: 5)

How many messages are processed in parallel within a single batch. This is useful when your callback involves I/O (HTTP calls, DB writes) — Node.js is single-threaded for CPU, but can run multiple async I/O operations concurrently.

FIFO Ordering

With concurrency: 1, messages are processed strictly in order (true FIFO). With concurrency > 1, messages are still fetched in FIFO order, but since each callback may take a different amount of time, the completion order is not guaranteed. Use concurrency: 1 when ordering matters.

typescript
await criticalQueue.subscribe(
  async (task) => { await processTask(task); },
  {
    batchSize: 100,    // Request up to 100 messages per network request
    concurrency: 10,   // Process 10 messages concurrently (I/O-bound tasks)
    waitMs: 5000       // If empty, wait 5s (server-side) before responding
  }
);
python
async def handle_task(task: CriticalTask) -> None:
    print(task)

await critical_queue.subscribe(
    handle_task,
    {
        "batch_size": 100,    # Request up to 100 messages per network request
        "concurrency": 10,    # Process 10 messages concurrently (I/O-bound tasks)
        "wait_ms": 5000       # If empty, wait 5s (server-side) before responding
    }
)

Dead Letter Queue (DLQ)

Every queue automatically has a dedicated DLQ. When a message exceeds maxDeliveries (default: 5), it's moved to the DLQ automatically — no setup needed.

Since DLQs are created alongside their parent queue, you can inspect failed messages at any time via queue.dlq.

Inspect Failed Messages

typescript
const failedMessages = await criticalQueue.dlq.peek(10);
console.log(`Found ${failedMessages.total} failed messages`);

for (const msg of failedMessages.items) {
  console.log(`Message ${msg.id}: attempts=${msg.attempts}, reason=${msg.failureReason}`);
  console.log(`Payload:`, msg.data);
}
python
failed_messages = await critical_queue.dlq.peek(10)
print(f"Found {failed_messages['total']} failed messages")

for msg in failed_messages["items"]:
    print(f"Message {msg['id']}: attempts={msg['attempts']}, reason={msg['failure_reason']}")
    print(f"Payload: {msg['data']}")

Replay or Discard

typescript
// Replay: move back to main queue (resets attempts to 0)
const moved = await criticalQueue.dlq.moveToQueue(msg.id);

// Discard: permanently delete from DLQ
const deleted = await criticalQueue.dlq.delete(msg.id);

// Purge: clear all DLQ messages
const purgedCount = await criticalQueue.dlq.purge();
python
# Replay: move back to main queue (resets attempts to 0)
moved = await critical_queue.dlq.move_to_queue(msg["id"])

# Discard: permanently delete from DLQ
deleted = await critical_queue.dlq.delete(msg["id"])

# Purge: clear all DLQ messages
purged_count = await critical_queue.dlq.purge()

API Reference

MethodDescriptionReturns
peek(limit, offset)Inspect messages without removing them{ total, items[] }
moveToQueue(messageId)Replay message to main queue (resets attempts)boolean
delete(messageId)Permanently remove a single messageboolean
purge()Remove all messages from DLQnumber (count)

Configuration

How it works

  1. Server starts → reads env vars (global defaults)
  2. Queue created → server snapshots defaults into config.json (per-queue)
  3. SDK can override visibilityTimeoutMs and maxDeliveries at creation — everything else uses system defaults
  4. On restart → each queue reads its own config.json (ignores current env vars)

Existing queues are not affected by env var changes. Only new queues pick up new defaults.

Environment Variables

Global, set at server startup.

VariableDefaultDescription
QUEUE_ROOT_PERSISTENCE_PATH./data/queuesBase directory for all queue SQLite DBs
QUEUE_VISIBILITY_MS30000 (30s)Default visibility timeout — how long before an unacked message is redelivered
QUEUE_MAX_DELIVERIES5Default max delivery attempts before moving to DLQ
QUEUE_DEFAULT_BATCH_SIZE10Default batch size for server-side consume
QUEUE_DEFAULT_WAIT_MS0Default long-polling wait (ms) when queue is empty
QUEUE_DEFAULT_FLUSH_MS100Max durability window (ms) — how often writes are flushed to disk
QUEUE_WRITER_BATCH_SIZE50000SQLite writer batch size (internal tuning)

SDK Overrides

Fields settable at create() time. If omitted, system defaults apply.

FieldSDK optionSystem default (env var)
Visibility timeoutvisibilityTimeoutMs30000 (QUEUE_VISIBILITY_MS)
Max deliveriesmaxDeliveries5 (QUEUE_MAX_DELIVERIES)