Skip to content

Consumer groups

Listing groups and their lag, changing a subscription, seeking a cursor per queue or per partition, and the admin DELETE routes.

Updated View as Markdown

Consumption state on this engine is one cursor per (partition, consumer group): queen.log_consumers.committed, the highest offset that group has completed in that partition. There is no per-message state, so every number on these routes is arithmetic over cursors and partition watermarks rather than a scan.

A pop that carries no consumerGroup uses the reserved group name __QUEUE_MODE__, which is why that name shows up in listings on queues where nobody set a group explicitly. Groups are created implicitly on first contact: the first pop for a (queue, group) writes a marker row in queen.consumer_groups_metadata and seeds a cursor for every existing partition in one statement.

Reads are read-only. subscription and both seek routes are read-write. Both DELETE routes are admin: route_access_level maps every DELETE under /api/v1/consumer-groups/ to admin regardless of what follows.

GET /api/v1/consumer-groups

Every group the tenant has, one entry per (group, queue) pair. The body is a bare JSON array.

[
  {
    "name": "invoicer",
    "topics": ["orders"],
    "queueName": "orders",
    "members": 8,
    "partitionCursors": 8,
    "partitionsWithLag": 2,
    "totalLag": 314,
    "maxTimeLag": 42,
    "state": "Stable",
    "storage": "segments",
    "subscriptionMode": null,
    "subscriptionTimestamp": null,
    "subscriptionCreatedAt": null
  }
]

Read these fields carefully:

  • members counts cursor rows, not consumer processes. One consumer on a 32-partition queue reports members: 32. partitionCursors carries the same number under its true name; members is kept only for wire compatibility.
  • totalLag is the sum over partitions of GREATEST(last_offset - GREATEST(committed, log_start - 1), 0): pure arithmetic, exact, and no segment scan. Frames already evicted by retention are not counted as lag.
  • maxTimeLag is the age in seconds of the oldest unconsumed segment. This is the one value that touches queen.log_segments.
  • state is derived, not stored: Lagging when maxTimeLag > 300, otherwise Stable when the group has ever consumed anything, otherwise Dead.
  • subscriptionMode, subscriptionTimestamp and subscriptionCreatedAt are always null for log-engine queues: the log schema has no per-partition subscription record.

GET /api/v1/consumer-groups/lagging

Query: minLagSeconds (default 3600). Returns a bare array of (group, partition) pairs whose time lag exceeds the threshold. The static lagging path is registered before the :group parameter route, so it cannot be shadowed by a group actually named lagging.

[
  {
    "consumer_group": "invoicer",
    "queue_name": "orders",
    "partition_name": "eu-3",
    "partition_id": "0198...",
    "worker_id": "0198...",
    "offset_lag": 512,
    "time_lag_seconds": 7421,
    "lag_hours": 2.06,
    "oldest_unconsumed_at": "2026-07-30T07:41:02.010Z",
    "last_consumed_at": null
  }
]

Keys are snake_case here: this route predates the camelCase listings and its shape is unchanged. time_lag_seconds is derived from the oldest unconsumed segment’s created_at, so a partition whose entire backlog was evicted drops out of the result rather than reporting a huge lag. last_consumed_at is always null: queen.log_consumers does not carry that column.

GET /api/v1/consumer-groups/:group

One group, broken out per queue and partition. The body is an object keyed by queue name.

{
  "orders": {
    "partitions": [
      {
        "partition": "eu-1",
        "workerId": "0198...",
        "lastConsumedAt": null,
        "totalConsumed": 91204,
        "offsetLag": 0,
        "timeLagSeconds": 0,
        "leaseActive": false
      }
    ]
  }
}

leaseActive is true only while lease_expires_at is in the future. That is the one in-flight leased batch this group holds on the partition. totalConsumed is the group’s lifetime counter from queen.log_consumers.total_consumed. A group with no cursors returns {}.

POST /api/v1/consumer-groups/:group/subscription

Body: subscriptionTimestamp (required, non-empty). A missing or empty value is a 400.

{ "subscriptionTimestamp": "2026-07-30T08:00:00Z" }

Response:

{
  "success": true,
  "consumerGroup": "invoicer",
  "newTimestamp": "2026-07-30T08:00:00Z",
  "rowsUpdated": 3
}

What this actually changes: the stored subscription timestamp on the group’s queen.consumer_groups_metadata rows (for every queue the group has touched, because the update carries no queue predicate), and it deletes the group’s empty-scan watermarks so a backward move re-exposes cold partitions to the next pop’s candidate scan.

Seek

Two routes, one body shape. Send either toEnd or timestamp:

{ "toEnd": true }
{ "timestamp": "2026-07-30T08:00:00Z" }

A body with neither is a 400 (Must specify toEnd=true or a timestamp). The per-partition route additionally accepts an empty body, which it treats as toEnd: true. That is what the dashboard’s per-partition “skip to end” button sends.

Cursor placement:

  • toEnd sets committed = last_offset, so the next wanted offset is last_offset + 1 and the whole existing backlog is skipped.
  • timestamp sets committed to the base offset of the first segment with created_at >= timestamp, minus one. Resolution is therefore segment-granular, not per-message. If no segment is that recent, the cursor lands on last_offset (a future-only cursor).

Either way the seek releases any live lease on the partition and resets the retry/redelivery state for that (partition, group). Because the timestamp walk starts at the retention watermark, seeking to a timestamp older than what is retained lands on the oldest retained segment.

POST /api/v1/consumer-groups/:group/queues/:queue/seek

Seeks every partition of the queue.

{
  "success": true,
  "consumerGroup": "invoicer",
  "queueName": "orders",
  "seekToEnd": true,
  "targetTimestamp": null,
  "partitionsUpdated": 8
}

If no partition matched (the queue does not exist for this tenant, or has no partitions yet), the stored procedure fails loudly rather than reporting a successful no-op:

{
  "success": false,
  "error": "Queue not found or has no partitions",
  "consumerGroup": "invoicer",
  "queueName": "nope",
  "partitionsUpdated": 0
}

That body is served with HTTP 404: the broker maps a stored-procedure error containing “not found” to 404 and anything else to 500.

POST /api/v1/consumer-groups/:group/queues/:queue/partitions/:partition/seek

Seeks one named partition.

{
  "success": true,
  "consumerGroup": "invoicer",
  "queueName": "orders",
  "partitionName": "eu-1",
  "seekToEnd": false,
  "targetTimestamp": "2026-07-30T08:00:00.000Z",
  "updated": true
}

An unknown partition returns {"success": false, "error": "Partition not found"} at 404.

After a successful seek the broker re-seeds its in-memory discovery ring for (queue, group) from committed database state and wakes parked pops, so a backward seek starts redelivering immediately instead of waiting for the periodic reseed. The reseed only re-adds partitions whose last_offset is genuinely ahead of committed, so a seek to end adds nothing.

DELETE routes

Both are admin. Both accept deleteMetadata (default true) and both answer 200 with a JSON body, not 204, because the SDKs read fields off the response.

DELETE /api/v1/consumer-groups/:group

Removes the group’s cursors across every queue of the tenant, plus its empty-scan watermarks, plus (when deleteMetadata) its queen.consumer_groups_metadata rows.

{
  "success": true,
  "consumerGroup": "invoicer",
  "deletedPartitions": 24,
  "metadataDeleted": true
}

deletedPartitions is the number of queen.log_consumers cursor rows removed.

DELETE /api/v1/consumer-groups/:group/queues/:queue

The same operation scoped to one queue: cursors for that queue’s partitions and the (queue, group) watermark row.

{
  "success": true,
  "consumerGroup": "invoicer",
  "queueName": "orders",
  "deletedPartitions": 8,
  "metadataDeleted": true
}

Clearing the watermark is the part that matters: an advanced empty-scan watermark would otherwise fence off every partition and the recreated group would appear to consume nothing.

Access levels

Route Level
GET /api/v1/consumer-groups read-only
GET /api/v1/consumer-groups/lagging read-only
GET /api/v1/consumer-groups/:group read-only
POST /api/v1/consumer-groups/:group/subscription read-write
POST /api/v1/consumer-groups/:group/queues/:queue/seek read-write
POST /api/v1/consumer-groups/:group/queues/:queue/partitions/:partition/seek read-write
DELETE /api/v1/consumer-groups/:group admin
DELETE /api/v1/consumer-groups/:group/queues/:queue admin

Every route here is tenant-scoped: the group name travels into the stored procedure together with the tenant, so two tenants can run a group called workers on their own queues without seeing each other’s cursors. Related pages: queues and resources, messages and DLQ.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close