Skip to content

S3 sink

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.

Updated View as Markdown

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/lake

Then 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: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:

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:

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/lake

It 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.json

Read 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.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):

{
  "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 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:

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 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. 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: 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.
Navigation

Type to search…

↑↓ navigate↵ selectEsc close