---
title: "S3 sink"
description: "Mirror queues into an S3-compatible bucket as JSONL or Parquet from inside the broker: a MinIO setup to try, how one writer per queue lands every record exactly once, the object layout, a bucket per tenant through the control plane, every setting, and the lag to watch."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# S3 sink

The broker can mirror its queues into an S3-compatible bucket, as JSONL or Parquet files in a Hive
layout that DuckDB, Spark or ClickHouse read with nothing in front of them. The sink runs inside the
broker on every node of a cluster, with no process or port of its own, and it is exactly once: each
record lands in one object, and a crash, a failover or a retried upload rewrites the same bytes under
the same key instead of adding a second copy.

## Try it

A MinIO bucket and a broker on one Docker network. MinIO needs a moment to start before the bucket
can be created:

```bash
docker network create queen-s3
docker run -d --name minio --network queen-s3 \
  -e MINIO_ROOT_USER=queen -e MINIO_ROOT_PASSWORD=queen-secret-key \
  minio/minio server /data
sleep 3
docker run --rm --network queen-s3 -e MC_HOST_m=http://queen:queen-secret-key@minio:9000 \
  minio/mc mb m/lake
```

Then the broker, with the sink pointed at that bucket:

```bash
docker run -d --name queen --network queen-s3 --platform linux/amd64 -p 6632:6632 \
  -v queen-data:/var/lib/queen/raft \
  -e QUEEN_S3_EMBEDDED=true \
  -e QUEEN_S3_QUEUES=orders \
  -e QUEEN_S3_ENDPOINT=http://minio:9000 \
  -e QUEEN_S3_REGION=us-east-1 \
  -e QUEEN_S3_PATH_STYLE=true \
  -e QUEEN_S3_BUCKET=lake \
  -e QUEEN_S3_ACCESS_KEY=queen \
  -e QUEEN_S3_SECRET_KEY=queen-secret-key \
  -e QUEEN_S3_START=earliest \
  -e QUEEN_S3_MAX_WINDOW_MS=10000 \
  ghcr.io/queen-mq/queen:latest
```

`QUEEN_S3_EMBEDDED` starts the sink, and `QUEEN_S3_QUEUES` names the queues it mirrors (`*` takes
every queue). The next six say where the bucket is: MinIO wants the bucket in the URL path, and
`us-east-1` is the region it answers to. The last two are for this demo. The queue `orders` does not
exist yet, and with the default start, `latest`, the sink would begin at the moment it first finds
the queue, after the records that create it; `earliest` reads it from the first record. And a window
closes after ten seconds here, where the default is five minutes.

Push three messages to two partitions:

```bash
curl -s -X POST localhost:6632/api/v1/push -H 'content-type: application/json' -d '{"items":[
  {"queue":"orders","partition":"cust-1","payload":{"order":1,"status":"new","total":42.00}},
  {"queue":"orders","partition":"cust-2","payload":{"order":2,"status":"new","total":17.50}},
  {"queue":"orders","partition":"cust-1","payload":{"order":1,"status":"paid","total":42.00}}
]}' > /dev/null
```

Until that push, the sink's row for `orders` said `missing`, and the sink looked for the queue again
every five seconds. In our run the window was committed seven seconds after the push:

```bash
curl -s localhost:6632/status | jq -c '.s3.sinks[0].queues[0] | {state, k, windowsCommitted, records, lagSeconds}'
```

```text
{"state":"filling","k":1,"windowsCommitted":1,"records":3,"lagSeconds":6.366447}
```

Window 1 holds the three records, and the sink is filling window 2. List the bucket:

```bash
docker run --rm --network queen-s3 -e MC_HOST_m=http://queen:queen-secret-key@minio:9000 \
  minio/mc ls --recursive m/lake
```

It holds two objects, the window's data and its manifest. Their keys, with the timestamps of our
run:

```text
queen/tenant=00000000-0000-0000-0000-000000000001/queue=orders/dt=2026-10-03/hour=06/w-0000000001-1791007200000000-1791009458377855.jsonl.zst
queen/_queen/tenant=00000000-0000-0000-0000-000000000001/queue=orders/windows/0000000001.json
```

Read the data object, with the key from your own listing:

```bash
docker run --rm --network queen-s3 -e MC_HOST_m=http://queen:queen-secret-key@minio:9000 \
  minio/mc cat 'm/lake/queen/tenant=00000000-0000-0000-0000-000000000001/queue=orders/dt=2026-10-03/hour=06/w-0000000001-1791007200000000-1791009458377855.jsonl.zst' \
  | zstd -dc
```

```text
{"partition":"cust-1","offset":0,"transactionId":"01a1007b-855d-7001-8107-adfdabd3e51c","ts":"2026-10-03T06:37:37.501942Z","payload":{"order":1,"status":"new","total":42.00}}
{"partition":"cust-1","offset":1,"transactionId":"01a1007b-855d-7003-ab83-e480f30b0bb6","ts":"2026-10-03T06:37:37.501942Z","payload":{"order":1,"status":"paid","total":42.00}}
{"partition":"cust-2","offset":0,"transactionId":"01a1007b-855d-7002-aa32-19bf977724d0","ts":"2026-10-03T06:37:37.501942Z","payload":{"order":2,"status":"new","total":17.50}}
```

Each record is one line, sorted by partition and offset. `ts` is the stamp the broker gave the push,
the same for every record of one push. The payload is the producer's own text, byte for byte:
`42.00` is still `42.00`, because the sink never parses a payload and prints it again.

## How it works

The sink is a library linked into the broker. The image carries it, and `QUEEN_S3_EMBEDDED=true` is
the only thing that starts it. It runs on threads of its own, named `queen-s3`, and reaches the
broker in memory: it reads each queue's log from this node's applied state, through the same calls
as [`POST /api/v1/fetch` and `POST /api/v1/partitions/changed`](/reference/http/#reads-by-offset),
and keeps three small documents per queue in the tenant's KV store. Everything else lives in the
bucket. Nothing connects to the sink, and it opens no port.

**One writer per queue.** Every node runs the same configuration, so every node has a task for every
queue, but a node writes a queue only while it holds the queue's lease. The lease is a KV row,
`s3:<sink>:<queue>:lease` in the namespace `queen-s3`, taken with `putIfAbsent` and a TTL of
`QUEEN_S3_LEASE_TTL_MS` (30 seconds), and refreshed every third of the TTL. The other nodes show the
queue as `held`, name the holder in `heldBy`, and look again later. A node stopped with SIGTERM gives
its leases back as it goes, and another node takes each queue at its next look. A node that dies
stops refreshing, and its queues move once the TTL runs out. A node restarted within the TTL finds
rows under its own name and takes them back at once.

The queues spread over the nodes by themselves. A node waits one second for every queue it already
runs before it claims a free one, so the least loaded node claims first. Each node also keeps a
presence row, which tells every node how many run the sink, and a node's fair share is the number of
queues divided by the number of nodes, rounded up. A node that has held more than its share for a
whole TTL gives back one queue at a time, the one with the least buffered, after draining it as a
SIGTERM would. In a steady state nothing moves, and a node that joins or leaves moves only what the
shares require.

**The fence.** A window's intent and its commit pointer are each written in one KV batch whose first
operation rewrites the lease row and must find there the version this node last wrote. A node whose
lease was taken over while it stalled, or while it was cut off from the leader, cannot move a pointer
when it comes back: the whole batch rolls back, `queen_s3_commit_precondition_lost_total` counts it,
and the queue stops on that node. Two nodes can never commit different windows of one queue.

**Windows on the broker's clock.** Window `k` of a queue holds the records whose `ts` is at or after
the end of window `k-1` and before its own end, and that end is never later than the node's
`safeTime` minus `QUEEN_S3_SAFE_GUARD_MS`. `safeTime` is the greatest stamp the node has applied, and
every record it has not applied yet will be stamped above it, so once a window closes no record can
join it. Rebuilding window `k`, on the same node or on the next owner, gives the same records in the
same order, hence the same bytes, and the key depends on `k` and the window's bounds alone. That is
the whole exactly-once argument, and we are fond of it because it asks nothing of the object store:
no conditional PUT, no LIST to find out what exists, no read-after-write. Any service that speaks the
S3 API can hold the lake.

A window commits in three steps. The intent fixes the window's end in KV before anything is
uploaded; then the objects go up, followed by the manifest; then the commit pointer moves, and only
that makes the window real. A node killed anywhere in between leaves the next owner either an intent
to redo, which writes the same objects under the same keys once that owner has applied the whole
window, or a finished upload whose manifest names that intent, which it commits as it is. A follower
can own a queue too: it reads its own applied log, and its KV writes reach the leader the way any
client's writes on that node do. Its `safeTime` trails the leader's, so its windows close that much
later.

A window closes when it holds `QUEEN_S3_TARGET_MB` of uncompressed records, when it is
`QUEEN_S3_MAX_WINDOW_MS` old on the broker's clock, when it reaches the next hour (`QUEEN_S3_ALIGN`),
or when the node's buffers pass `QUEEN_S3_MEMORY_MB` and it is the queue holding the most. A stretch
of time with no records is skipped with no object. On a broker with no other traffic, the sink's own
lease refreshes are log entries, so `safeTime` keeps moving and the last window before a quiet spell
still reaches its age.

Two consequences surprise people. The lake mirrors the log, whatever consumers did with it: the
sink's reads take no lease and move no cursor, so a message a consumer group nacked or dead-lettered
is in the lake like any other. And the lake is plaintext: the node decrypts the payloads of an
[encrypted queue](/operate/security/#payload-encryption-at-rest) before writing them, so give the
bucket encryption of its own (`QUEEN_S3_SSE`, or the bucket's default encryption).

## The objects

```text
<prefix>/tenant=<tenant>/queue=<queue>/dt=<YYYY-MM-DD>/hour=<HH>/w-<k>-<tStart>-<tEnd>.<ext>
<prefix>/_queen/tenant=<tenant>/queue=<queue>/windows/<k>.json
<prefix>/_queen/tenant=<tenant>/queue=<queue>/checkpoint/<k>.json.zst
```

The prefix is `queen` unless `QUEEN_S3_PREFIX` says otherwise. The tenant comes first, so two tenants
that share a bucket never share a key; a broker without tenancy writes the default tenant,
`00000000-0000-0000-0000-000000000001`. Tenant, queue and partition names are percent-encoded outside
`A-Z a-z 0-9 . _ -`, so no name can climb out of its prefix. `dt=` and `hour=` come from the window's
start, and they are exact for every record in the object because a window never crosses the
alignment boundary (`QUEEN_S3_ALIGN=day` drops `hour=`; `none` keeps both, from the window's start).
`k` is padded to ten digits, so a listing returns the windows in commit order, and `tStart` and
`tEnd` are microseconds since the epoch: the first window of the run above starts at 06:00:00, the
hour its first record falls in. With `QUEEN_S3_LAYOUT=per-partition` a window writes one object per
partition, its key ending in `-p-<partition>-<first offset>-<last offset>`. The sidecars live under
`_queen/`, which is not a Hive partition, so a reader that globs `tenant=*/queue=*/dt=*/hour=*/*`
never sees them.

The manifest of that window, the one file in the bucket with a wall-clock time in it
(`committedAt`):

```json
{
  "sink": "default",
  "tenant": "00000000-0000-0000-0000-000000000001",
  "queue": "orders",
  "k": 1,
  "tStart": 1791007200000000,
  "tEnd": 1791009458377855,
  "format": "jsonl",
  "compression": "zstd",
  "layout": "merged",
  "objects": [
    {
      "key": "queen/tenant=00000000-0000-0000-0000-000000000001/queue=orders/dt=2026-10-03/hour=06/w-0000000001-1791007200000000-1791009458377855.jsonl.zst",
      "bytes": 207,
      "records": 3,
      "sha256": "5f2ac4baa20940c6c33a5e3c7207ef7d1087db6538fea0cce88cb16ff43def78"
    }
  ],
  "records": 3,
  "bytes": 207,
  "partitions": 2,
  "minTs": 1791009457501942,
  "maxTs": 1791009457501942,
  "lost": [],
  "writer": "queen-s3/2.0.0 jsonl+zstd",
  "committedAt": "2026-10-03T06:37:44.741000Z"
}
```

`lost` is empty unless retention deleted records before the sink read them, the one way the lake can
miss data ([see the lag](#lag-and-the-retention-hold)). The checkpoint, written every
`QUEEN_S3_CHECKPOINT_EVERY` windows, holds the next offset of each partition. It only shortens the
re-read after a restart, and losing it costs time and nothing else.

**JSONL**, the default, writes one record per line with five fields in a fixed order: `partition`,
`offset`, `transactionId`, `ts` (microseconds, UTC) and `payload`. The payload is the stored text,
`null` when it was empty, with one change: a raw line break inside a pretty-printed payload is
written as a space, so a record stays on one line. A payload the node cannot decrypt is written as a
JSON string of its stored bytes. `QUEEN_S3_COMPRESSION` picks `zstd` (the default), `gzip` or `none`.

**Parquet** has the same five columns:

```text
message queen_record {
  required binary partition (STRING);
  required int64 offset;
  required binary transaction_id (STRING);
  required int64 ts (TIMESTAMP(MICROS,true));
  optional binary payload (STRING);
}
```

`payload` is the JSON text, for `json_extract` and its cousins. Row groups close every 100,000
records, the codec is `zstd` or `snappy` (`QUEEN_S3_PARQUET_CODEC`), and the footer names the queue
and the tenant (`queen.queue`, `queen.tenant`), so a lone file still says where it came from.

The tenant and the queue live only in the path, which is where Hive, Spark and Athena expect a
partition key, and a reader with Hive partitioning hands them back as columns, with `dt` and `hour`
to prune on. In DuckDB, with its S3 settings pointed at the bucket:

```sql
SELECT partition, "offset", ts, payload
FROM read_parquet('s3://lake/queen/tenant=*/queue=orders/dt=*/hour=*/*.parquet', hive_partitioning = true)
WHERE dt = '2026-10-03'
ORDER BY partition, ts, "offset";
```

The connector's compatibility lane read these five-field records with DuckDB 1.5.5, ClickHouse
26.8, Spark 4.0.1, Polars, PyArrow and pandas on 2026-09-04, before keys had the `tenant=` level
(`connectors/queen-s3/compat/MATRIX.md`; the query above is its DuckDB call with that level added).
Spark and pandas could not open zstd JSONL out of the box. That is why `gzip` exists, and why
Parquet, which names its codec inside the file, is the safe choice for a lake several tools read.

## A bucket per tenant

On a cell that runs the [embedded proxy](/operate/tenants/), each cluster can mirror its own broker
tenant into its own bucket, with its own credentials. The control plane sets it up, and the
`QUEEN_S3_*` environment stays the default tenant's. Every node needs `QUEEN_S3_EMBEDDED=true`
(without it the routes answer `404 s3_unavailable`), `QUEEN_PROXY_CP_TOKEN`, and the same
`QUEEN_ENCRYPTION_KEY`, because the cluster's S3 secret is stored only sealed with that key, in the
replicated KV under the proxy's own tenant:

```bash
curl -s -X PUT localhost:6711/api/cp/clusters/acme/s3 -H "x-queen-cp-token: $CP_TOKEN" \
  -H 'content-type: application/json' -d '{
  "endpoint": "https://s3.eu-central-1.amazonaws.com",
  "region": "eu-central-1",
  "bucket": "acme-lake",
  "accessKey": "<access key id>",
  "secretKey": "<secret access key>",
  "queues": "*",
  "format": "parquet"
}'
```

The body is the per-tenant settings of the [table below](#reference), by their field names, plus two
more. `secretKey` is required the first time (4 KiB at most) and can be left out later to keep the
stored one; no answer or log line ever carries it. `enabled` is `true` unless you send `false`,
which stops the sink and keeps its row. The broker checks the settings with the sink's own rules
before storing anything: a node-wide setting (`memoryMb`, `leaseTtlMs` and the rest), an unknown
field, or another field named like a secret is refused, and the settings fit in 32 KiB. The answer,
and `GET` on the same path, is the stored sink without its secret: `cluster`, `tenant`, `enabled`,
`config`, `secretKeySet` and `updatedAt`. `DELETE` answers `{"cluster": ..., "removed": true}`,
or `false` when there was nothing to remove.

Every node reads these rows again every five seconds and runs one sink per tenant to match. A sink
runs while its row is enabled and its cluster is `active` or `push_blocked` (the worse of the
cluster's and its tenant's status, as for the data plane), so a cluster that may not push still
ships what it holds, and a `suspended` one ships nothing. Every write to the row moves `updatedAt`,
and every node then rebuilds that tenant's sink from the new row; a `PUT` that carries `secretKey`
always writes, which is how a rotated secret reaches the sink. Purging or deleting a tenant removes
its sinks too. Nothing ever deletes objects from a bucket.

| Answer | When |
|---|---|
| `404 s3_unavailable` | The node that answered does not run the sink |
| `404 cluster_unknown` | No cluster has that slug |
| `404 s3_unset` | `GET` of a cluster with no sink |
| `400 invalid` | A body that is not an object, a missing first `secretKey`, or a setting the sink refuses, with the sink's own sentence |
| `409 system_tenant` | The cluster is the broker's default tenant or the proxy's own, or `QUEEN_PROXY_TENANT_HEADER` is off |
| `409 deleting` | The cluster or its tenant is being deleted |
| `409 encryption_required` | The cell has no `QUEEN_ENCRYPTION_KEY` to seal the secret with |

A node whose key differs from the one that sealed a secret cannot open it. That tenant's sink shows
phase `error` on that node, which runs none of its queues, while the other nodes carry on.

## Reference

Every variable is read once, at boot, and every node should get the same values. A blank value
counts as unset, enumerated values ignore case, and a required variable that is missing, or a value
outside its range, stops the boot with one line that names it:

```text
ERROR boot: FATAL: QUEEN_S3_EMBEDDED=true: QUEEN_S3_MAX_WINDOW_MS=50 is outside 100..=86400000. It is how long a window may stay open; the lag SLO is about this plus the safe lag
```

**The broker's own.**

| Variable | Default | |
|---|---|---|
| `QUEEN_S3_EMBEDDED` | `false` | `true` starts the sink on this node. Off, nothing of the sink runs and its other settings do not matter |
| `QUEEN_S3_THREADS` | a quarter of the cores, 1 or 2 | 1 to 64 threads, shared by the sinks of every tenant |
| `QUEEN_S3_SHUTDOWN_GRACE_MS` | `30000` | 100 to 3,600,000: how long the drain may take, counted from SIGTERM |

**Per sink**, as an environment variable for the default tenant and as a document field for a
tenant the control plane configures, with the same defaults and ranges:

| Variable | Field | Default | |
|---|---|---|---|
| `QUEEN_S3_QUEUES` | `queues` | required | Comma-separated names, or `*` for every queue of the tenant, listed again every ten discovery intervals. A document may also send a JSON list. Setting it is what turns the default tenant's sink on |
| `QUEEN_S3_ENDPOINT` | `endpoint` | required | `http://` or `https://`, a host and an optional port, no path |
| `QUEEN_S3_REGION` | `region` | required | The region the requests are signed for; gateways commonly accept `us-east-1` |
| `QUEEN_S3_BUCKET` | `bucket` | required | One name, no slash or space |
| `QUEEN_S3_ACCESS_KEY` | `accessKey` | required | |
| `QUEEN_S3_SECRET_KEY` | `secretKey`, beside the document | required | Never logged or returned |
| `QUEEN_S3_PREFIX` | `prefix` | `queen` | The root of every key, sidecars included. Outer slashes are dropped; an empty, `.` or `..` segment is refused |
| `QUEEN_S3_PATH_STYLE` | `pathStyle` | `false` | `true` puts the bucket in the URL path, as MinIO and most self-hosted gateways want |
| `QUEEN_S3_SSE` | `sse` | unset | `AES256` or `aws:kms`, sent with every upload. Unset sends no header, and the bucket's default applies |
| `QUEEN_S3_SSE_KMS_KEY_ID` | `sseKmsKeyId` | unset | With `aws:kms` only |
| `QUEEN_S3_FORMAT` | `format` | `jsonl` | `jsonl` or `parquet` |
| `QUEEN_S3_COMPRESSION` | `compression` | `zstd` | JSONL only: `zstd`, `gzip` or `none` |
| `QUEEN_S3_PARQUET_CODEC` | `parquetCodec` | `zstd` | Parquet only: `zstd` or `snappy` |
| `QUEEN_S3_LAYOUT` | `layout` | `merged` | `merged`, one object per window, or `per-partition`, one per partition per window |
| `QUEEN_S3_ALIGN` | `align` | `hour` | `hour`, `day` or `none`: the boundary no window crosses |
| `QUEEN_S3_START` | `start` | `latest` | Where a queue with no commit pointer starts: `latest`, when the sink first finds it, or `earliest`, everything retention still holds |
| `QUEEN_S3_TARGET_MB` | `targetMb` | `128` | 1 to 5,120 MB of uncompressed records close a window |
| `QUEEN_S3_MAX_WINDOW_MS` | `maxWindowMs` | `300000` | 100 to 86,400,000: the most a window stays open, on the broker's clock |
| `QUEEN_S3_SINK` | `sink` | `default` | 1 to 64 characters of `A-Z a-z 0-9 . _ -`. Names the KV documents and is what `retentionSinkHold` refers to. It is in no object key, so two sinks mirroring one queue into one bucket need different prefixes |

**Node-wide**, environment only, shared by every sink on the node. A tenant document that names one
is refused:

| Variable | Default | |
|---|---|---|
| `QUEEN_S3_MEMORY_MB` | `512` | 1 to 1,048,576: one buffer budget for every queue of every sink on the node |
| `QUEEN_S3_FETCH_CONCURRENCY` | `4` | 1 to 256 reads in flight per queue, the throttle for a backfill |
| `QUEEN_S3_DISCOVERY_INTERVAL_MS` | `2000` | 10 to 3,600,000: how often an idle queue asks which partitions moved |
| `QUEEN_S3_SAFE_GUARD_MS` | `5000` | 0 to 3,600,000, subtracted from `safeTime` before a window may close |
| `QUEEN_S3_LEASE_TTL_MS` | `30000` | 1,000 to 3,600,000 |
| `QUEEN_S3_MULTIPART_THRESHOLD_MB` | `64` | 5 to 5,120: larger objects go up as a multipart upload in 16 MiB parts |
| `QUEEN_S3_CHECKPOINT_EVERY` | `20` | 1 to 100,000 windows between checkpoints |
| `QUEEN_S3_INSTANCE` | `node-<id>@<host>` | This node's name in lease rows, built from `QUEEN_RAFT_NODE_ID` and `HOSTNAME`. A value you set gets `/node-<id>` appended unless it already ends with it, so two nodes never share a name |
| `QUEEN_S3_CRASH_AT` | `never` | For the crash tests only: it aborts the broker at a point of the commit sequence |

`QUEEN_S3_PARTITIONS` (a static partition list) stops the boot, because a window can only close
against what discovery reports. `QUEEN_S3_LISTEN` and `QUEEN_S3_LOG_FORMAT`, left over from the
standalone 1.5.0 sink, are named in one warning and ignored. The sink logs under the target
`queen-s3`, and one line per sink at start lists every setting it runs with, never the secret key.

**The KV documents**, per sink and queue, in the namespace `queen-s3` of the queue's own tenant
(the queue percent-encoded as in the object keys):

| Key | |
|---|---|
| `s3:<sink>:<queue>:lease` | Who runs the queue, with the lease TTL |
| `s3:<sink>:<queue>:intent` | The window being committed: `k`, its bounds, the format |
| `s3:<sink>:<queue>:committed` | The commit pointer: `k`, `tEnd`, the manifest's key, records and bytes |

**Queue options** for the retention hold, on `POST /api/v1/configure`:

| Option | Default | |
|---|---|---|
| `retentionSinkHold` | `""` (off) | The sink's name; retention keeps whatever that sink has not committed |
| `retentionSinkHoldMaxSeconds` | `604800` | 60 to 31,536,000: the most the hold can keep |

## Operating it

### Status

`GET /status` on a node running the sink carries an `s3` block, and since each node answers for
itself, read it on each one. The block has the node's `phase`, its `threads`, the last error reading
the control plane's rows, and one entry per tenant sink under `sinks`: its `tenant`, `source` (`env`
or `cp`) and `phase` (`starting`, `running`, `backoff`, `stopped` or `error`, with the reason in
`error`), then the sink's own report. In the run above:

```bash
curl -s localhost:6632/status | jq -c '.s3.sinks[0] | {phase, ok, health, bucket: .bucket.reachable, placement}'
```

```text
{"phase":"running","ok":true,"health":{"ok":true,"queues":1},"bucket":true,"placement":{"givingBack":null,"held":1,"nodes":1,"queues":1,"share":1}}
```

`ok` is false when the health verdict is, or when the bucket did not answer the probe the sink sends
before it starts any queue; `bucket.error` then says what to check. `placement` shows the nodes
running the sink, the fair share, how many queues this node holds and the one it is giving back.
`queues` has a row for every queue: `state` (`claiming`, `held`, `restoring`, `filling`, `intent`,
`upload`, `commit`, or how the last run ended: `drained`, `fenced`, `failed`, `missing`, `crashed`,
`released`), `heldBy`, the last committed `k` and `tEnd`, `completeThrough`, `safeTime`,
`lagSeconds`, what this node committed since it started, and `lastError`, kept until a newer one
replaces it. `/health` does not look at the sink.

### Lag and the retention hold

`queen_s3_lag_seconds{queue}` is the number to alarm on. It is the owning node's `safeTime` minus
`completeThrough`, the stamp below which every record of the queue is in a committed window. On a
busy queue it climbs to a little past `QUEEN_S3_MAX_WINDOW_MS`, drops when a window commits, and
climbs again. On an idle queue it stays near the guard: between 5 and about 10 seconds in our run,
with the default guard. The health verdict turns red when a queue lags more than three windows (30
seconds at least), so alarm at a multiple of the window. `queen_s3_safe_lag_seconds` adds how far the
node's applied log trails its own clock, which grows on a follower that falls behind.

The lag matters because of what sits behind it. If retention deletes records before the sink has read
them, they are gone from the lake too, and that is the one way it can lose data. The loss is never
silent: the range goes into the window's manifest as `"lost": [{"partition": ..., "from": ..., "to": ...}]`,
counts on `queen_s3_records_lost_total`, is logged when the window commits, and the sink carries on
from the new start of the log. To rule it out, let the queue's retention wait for the sink:

```bash
curl -s -X POST localhost:6632/api/v1/configure -H 'content-type: application/json' \
  -d '{"queue":"orders","options":{"retentionEnabled":true,"retentionSeconds":86400,"retentionSinkHold":"default"}}'
```

`retentionSeconds` and `completedRetentionSeconds` then delete nothing stamped after the commit
pointer's `tEnd` minus 60 seconds, read from the queue's own tenant. The hold never reaches back
further than `retentionSinkHoldMaxSeconds`, seven days by default, so a sink that stopped for good
cannot keep a queue's disk growing for ever, and a queue whose sink has not committed yet keeps
everything younger than that cap. `maxWaitTimeSeconds` ignores the hold.

### Metrics

On `/metrics/prometheus`, one exposition per node, with `tenant="<id>"` on the series of every
tenant but the default one: `queen_s3_lag_seconds`, `queen_s3_safe_lag_seconds`,
`queen_s3_windows_committed_total`, `queen_s3_records_written_total`,
`queen_s3_bytes_written_total`, `queen_s3_records_lost_total`, `queen_s3_window_records` and
`queen_s3_window_bytes` (histograms), `queen_s3_buffer_bytes`, `queen_s3_fetch_calls_total`,
`queen_s3_discovery_partitions`, `queen_s3_s3_requests_total` (by operation and status code),
`queen_s3_commit_precondition_lost_total`, `queen_s3_checkpoint_age_windows` and
`queen_s3_queues_given_back_total`. A node exports a queue's gauges only while it runs the queue.

### Stopping

The drain starts at the signal, beside the broker's own shutdown. Every queue stops reading,
finishes a window that is already in its commit sequence, and commits the window it is filling when
that is worth it (1 MiB buffered, or open for ten seconds on the broker's clock), up to `safeTime`
minus the guard as always. Then it gives its lease back, retrying for up to ten seconds when a call
fails, as it can on a leader that is handing over its leadership. Records too young for the drain
stay in the log for the next owner: in our run, a record pushed one second before the SIGTERM landed
in window 2 as soon as the node came back. The broker waits for what is left of
`QUEEN_S3_SHUTDOWN_GRACE_MS` once its listener has drained. A drain cut short loses nothing, since
the next owner redoes the window from its intent, but its queues wait out the lease TTL.

### On Kubernetes

Give every pod the same `QUEEN_S3_*` environment, since the leases decide which pod writes which
queue, and take the secret key and `QUEEN_ENCRYPTION_KEY` from a Secret. The 90-second
`terminationGracePeriodSeconds` of the [Kubernetes manifest](/operate/kubernetes/#what-sigterm-does)
covers the default 30-second grace; raise both together. Size the memory limit for the broker plus
`QUEEN_S3_MEMORY_MB` plus the object a closing window builds in memory, because the buffers are the
broker's own memory now. Leave `QUEEN_S3_INSTANCE` unset: with `QUEEN_RAFT_NODE_ID=ordinal` the
default names the pod and stays the same across restarts, so a restarted pod takes its own queues
back at once.

### The bucket

The sink needs `s3:PutObject`, `s3:GetObject` and `s3:ListBucket` on its prefix, plus
`s3:AbortMultipartUpload`, and it never deletes an object. On Google Cloud Storage a retried upload
overwrites its object byte for byte, and an overwrite needs `storage.objects.delete`, so grant the
service account `roles/storage.objectUser` rather than `objectCreator`, and connect with
`QUEEN_S3_ENDPOINT=https://storage.googleapis.com`, `QUEEN_S3_REGION=auto`,
`QUEEN_S3_PATH_STYLE=true`, a service-account HMAC key, and `QUEEN_S3_SSE` unset. Add a lifecycle rule that aborts
incomplete multipart uploads: the sink aborts the ones that fail, but a node killed between two parts
leaves one behind. Checkpoints accumulate, one every `QUEEN_S3_CHECKPOINT_EVERY` windows per queue,
and removing old ones costs nothing but a longer re-read after a restart.

## Limits

- Exactly once covers the bucket: each committed record is in one object, once. The lake mirrors
  the log faithfully, so a `transactionId` reused outside the broker's dedup window appears twice in
  the lake as it does in the log.
- Windows are per queue, so one queue's throughput is one node's. More nodes spread more queues.
- The buffers live in the broker's process. An allocation failure or the kernel's OOM killer takes
  the broker with it, so size the container for both.
- A record that retention deletes before the sink reads it is lost to the lake (recorded, never
  silent). Use the retention hold, or keep retention well above the window plus the longest outage
  the sink must survive.
- `QUEEN_S3_START=earliest` on a queue with a long log reads all of it through the owning node, and
  the object store bills every byte. Turn `QUEEN_S3_FETCH_CONCURRENCY` down for a backfill.
- With `QUEEN_RAFT_CLIENT_OFFLOAD=false` and JWT on, followers cannot write the sink's KV documents,
  so the leader runs every queue.
- Nothing writes a catalog (Iceberg, Glue). The manifests hold most of what a catalog entry needs:
  row counts, sizes and time bounds.
- Not part of the [Jepsen campaign](/concepts/guarantees/). The sink is tested by its own suite, by
  in-process tests against the broker's state machine, and end to end (`test/s3sink`) against a local
  S3 server: a crash at each of the five points of the commit sequence, `kill -9` takeovers and
  SIGTERM hand-overs on three nodes, a bucket that is down or answers errors, and control-plane
  tenants sharing a bucket. Parquet, the hold and retention overruns are covered by the unit and
  in-process tests only.

## Next

- [Postgres source and sink](/guides/postgres/): the other connector inside the broker, and how its
  exactly-once argument works against a database.
- [Guarantees](/concepts/guarantees/): what the log underneath the lake promises, and how it was
  tested.

Source: https://queenmq.com/guides/s3/index.mdx
