Skip to content

The model

Entity, partition, step and log: how Queen MQ's primitives fit together, with a glossary of every term these docs use.

Updated View as Markdown

Queen has one idea, and every concept page zooms into a part of it. Every entity gets an ordered partition; a worker handles one event of that partition as a step; and the step commits as one entry of the broker’s log. Read this page first, and the others will each fill in one piece.

A step on the wire

This is one step, as the broker receives it. A worker in the billing group took a payment event from the partition order-9137, and now acks it, pushes a receipt, records the order as paid, makes sure the payment is handled once, and asks for a reminder tomorrow:

{
  "operations": [
    { "type": "ack", "transactionId": "payment-9137-succeeded", "partitionId": "5",
      "consumerGroup": "billing", "leaseId": "01a0fcd7-e9bc-7001-8866-1685bf9d4343",
      "status": "completed" },
    { "type": "push", "items": [{ "queue": "receipts", "partition": "customer-42",
      "transactionId": "receipt-9137", "payload": { "orderId": 9137 } }] }
  ],
  "kv": [
    { "op": "put", "ns": "orders", "key": "9137", "value": { "status": "paid" },
      "ttlSeconds": 2592000 },
    { "op": "putIfAbsent", "ns": "paid", "key": "payment-9137-succeeded", "value": true,
      "ttlSeconds": 86400, "required": true }
  ],
  "timers": [
    { "op": "schedule", "queue": "reminders", "timerKey": "9137", "delayMs": 86400000,
      "txn": "remind-9137", "payload": "eyJvcmRlcklkIjo5MTM3fQ==" }
  ]
}

POST /api/v1/transaction with this body answers "success": true and one result per part, and the broker has written all of it as one log entry. Sent a second time, it answers "success": false with "reason": "rejected_ack", because the message is already acked, and writes nothing. kv and timers are top-level arrays beside operations, never elements of it. The SDKs build this body for you (queen.transaction() in JavaScript); see transactions.

Entity, partition, step, log

An entity is what you think in: a customer, an order, a conversation. Its partition is its ordered input. You name it on push ("partition": "order-9137"), and the first push creates it, so partitions are as many as your entities and cost the broker a few rows each.

A step is what a worker does with one event of a partition. A consumer group leases the partition’s next batch to one worker at a time, which keeps the entity’s events in order while other workers handle other entities, and the worker’s step ends in one commit.

The log is where every write lands. On a single node, the node fsyncs an entry before it answers; in a cluster, raft copies the entry to a majority of the nodes, each of which fsyncs it, before you get an answer. One step is one entry, so a step is never half written.

Producers push events for order-9137 into that order's own partition, which keeps them in order. A worker of the consumer group leases the partition's next batch, and its step ends in one commit: one log entry holding the ack, the state change, the next events and any timer. Other entities have partitions of their own, leased to other workers in parallel, and their steps commit into the same log.producerspush to order-9137order-9137one entity, in orderother entitiesa partition eachworkerholds the leaseother workersin parallelone log entryack, state, events, timerpushnext batchone committheirs
An entity's events wait in its own partition, one worker at a time takes the next batch, and the step it runs ends in a single entry of the log.

The five roles of a step

What a worker does with an event falls into five roles, and each is a primitive of the same log:

Role What the worker does Primitive Page
Events emits the next events queues, one partition per entity partitions
Progress takes the event and acks it consumer groups, leases, acks, retries, the DLQ, replay consuming
State changes the entity’s state KV KV
Time schedules what comes later timers, delayed delivery timers
Identity makes a retry harmless transactionId dedup, once dedup and once

In most stacks those are four systems, a broker, a database, a scheduler and an idempotency table, with an outbox to carry events from the database to the broker. What makes event-driven code hard is keeping their separate commits consistent when something crashes between two of them, so we put all five in one log, where a step has a single commit to get right.

Two ways to run one step. On the left, the usual stack: the worker acks and publishes the next events on a broker, writes the state and an outbox row in a database, schedules a reminder in a scheduler, and records the event id in an idempotency table, four systems with four separate commits. On the right, Queen: the worker sends one transaction, and the ack, the next events, the state, the timer and the identity check land in one log entry.workerbrokerack, next eventsdatabasestate, outboxschedulerreminderidempotencytableworkerone log entryack, events, state,timer, identityfour systems, four commitsQueen, one commit
Four commits can disagree when something crashes between two of them. One entry cannot: it is written whole or not at all.

Stateless workers, a state machine per entity

Per-entity order and a one-entry commit together let you run every entity as a small state machine. Its partition is the input, in order; its KV entry is the state; a transaction is one transition; and a timer is a transition that time triggers. Workers keep nothing between events, so any worker can take any partition, and you scale by adding workers until every partition with work has one. See one state machine per entity.

Exactly-once, precisely

Delivery is at-least-once: a lease that expires, a nack or a crash delivers a message again. What Queen makes exactly-once is the effect of a step inside the broker. The ack and the step’s writes are one entry, fenced by the lease the ack carries, and when the step carries a gate (a push with a deterministic transactionId, or a once marker), a retry of a step that already committed is refused whole. A call your worker makes to another system (a payment provider, an email API) is outside any commit, and needs that system’s own idempotency key. See guarantees.

Glossary

Term Meaning
Queue A named set of partitions with one configuration (lease time, retries, dedup window, retention). Created by the first push that names it, or by POST /api/v1/configure.
Partition One entity’s ordered stream inside a queue. Offsets start at 0 and have no gaps. A push with no partition goes to Default.
Consumer group A name that reads a queue at its own pace, with one cursor per partition. A pop without a group uses the shared group __QUEUE_MODE__.
Cursor, offset The offset is a message’s position in its partition. The cursor is the last offset a group acked there; the group’s next message is the one after it.
Lease A pop’s claim on one partition for one group: a leaseId, the end of the batch and an expiry. One lease per partition and group at a time.
Ack, nack An ack (completed) moves the cursor past a message. A nack (failed) releases the lease and spends one retry. retry and dlq are the other two statuses.
Dead-letter queue Where a message goes when its retries are spent, or when a worker acks it with dlq. It keeps the payload and the error until you replay or delete it.
Transaction One POST /api/v1/transaction: acks, pushes, KV writes and timer operations written as one log entry, all or nothing.
transactionId A message’s idempotency key. A second push of the same id to the same partition, inside the dedup window, writes nothing.
KV namespace A named set of keys inside a tenant, such as orders.
Timer A message promised for later, keyed by queue and timerKey, pushed into its queue through the log when it fires.
Tenant An isolated set of queues, KV namespaces and timers. With QUEEN_TENANCY_HEADER off, every request uses the default tenant.
Raft group One replicated log with its own leader. A node runs QUEEN_RAFT_GROUPS of them (default 1), and a tenant lives in exactly one.

Limits

  • Order holds inside a partition, and two partitions have no order between them.
  • A consumer group leases one batch per partition at a time, so a hot partition is sequential and parallelism comes from many partitions.
  • A transaction is one call inside one tenant. There is no interactive BEGIN and COMMIT.

Next

Navigation

Type to search…

↑↓ navigate↵ selectEsc close