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:
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/lakeThen the broker, with the sink pointed at that bucket:
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:latestQUEEN_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:
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/nullUntil 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:
curl -s localhost:6632/status | jq -c '.s3.sinks[0].queues[0] | {state, k, windowsCommitted, records, lagSeconds}'{"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:
docker run --rm --network queen-s3 -e MC_HOST_m=http://queen:queen-secret-key@minio:9000 \
minio/mc ls --recursive m/lakeIt holds two objects, the window’s data and its manifest. Their keys, with the timestamps of our run:
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.jsonRead the data object, with the key from your own listing:
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{"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,
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 before writing them, so give the
bucket encryption of its own (QUEEN_S3_SSE, or the bucket’s default encryption).
The objects
<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.zstThe 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):
{
"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). 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:
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:
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, 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:
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, 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:
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 lagThe 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:
curl -s localhost:6632/status | jq -c '.s3.sinks[0] | {phase, ok, health, bucket: .bucket.reachable, placement}'{"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:
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
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
transactionIdreused 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=earlieston a queue with a long log reads all of it through the owning node, and the object store bills every byte. TurnQUEEN_S3_FETCH_CONCURRENCYdown for a backfill.- With
QUEEN_RAFT_CLIENT_OFFLOAD=falseand 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. 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 -9takeovers 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: the other connector inside the broker, and how its exactly-once argument works against a database.
- Guarantees: what the log underneath the lake promises, and how it was tested.