High-volume Kafka-to-storage archiver designed for PB/s scale data pipelines. Built on the scalo data-plane runtime (config cascade, logging, metrics, Kafka transport, tiered sink, deployment contract).
- Multiple destinations: File, MinIO, S3, GCS, Azure Blob
- Compression: Zstd (default), LZ4, Snappy, Gzip, or none, each in its standard framed format (LZ4 frames
.lz4, the snappy framing format.sz, gzip members.gz, zstd frames.zst), so the codec's own tools read an archive file whole - Smart routing: By JSON field expressions (e.g.,
org_id) or topic - Rolling archives: By final compressed file size (1GB default) or time (1 hour default)
- At-least-once delivery on Kafka: an offset is committed once the archive file holding its record is complete
- Memory-capped: every destination's buffer together holds at most a quarter of the memory limit, the largest flushing first
- Disk-bounded: intake pauses as staged uploads near 8 GiB, and a full disk holds intake rather than dropping records
flowchart LR
K[("Kafka<br/>batch recv, 10K msgs")] --> BM["Buffer Manager<br/>per-destination buffering"]
BM --> AW["Archive Writer<br/>compressed rolling files"]
AW --> ST["Storage backend<br/>File / S3 / GCS / Azure / MinIO"]
ST -. file complete .-> C["Kafka offset commit<br/>at-least-once"]
The workspace is three crates -- core, io and archiver -- described under
Where things live.
docs/architecture.md carries the codemap, the one-way
rules between them, and the build graph.
Handles high destination cardinality (e.g., 10,000+ orgs) without exhausting memory:
flowchart TB
R["incoming records"] --> HB["Hot buffers, in memory<br/>64 destinations at once<br/>a quarter of the memory limit together"]
HB -->|"flush_bytes, flush_records, age,<br/>LRU eviction or the cap"| AW["Archive writers<br/>one per destination, up to max_writers (1024)"]
AW -->|file closes| ST["Local file synced, or staged file uploaded<br/>writer_parallelism uploads at once"]
ST -. held records or staged bytes near their caps .-> BP["Intake pauses<br/>Kafka partitions paused"]
Nothing spools to disk ahead of the writers: the only local files are the archive files, and the staged copies of object-store files until the store takes them.
# Build
cargo build --release
# Run with environment variables
KAFKA_BROKERS=localhost:9092 \
KAFKA_TOPICS=events \
ARCHIVER_DESTINATION=file://./data/archive \
./target/release/dfe-archiver
# Or with config file
./target/release/dfe-archiver --config config.yamlConfiguration follows a cascade (highest to lowest priority):
- CLI arguments
- Environment variables
.envfile- Config file (
config.yaml) - Built-in defaults
| Variable | Description | Default |
|---|---|---|
ARCHIVER_TRANSPORT |
kafka (a broker) or grpc (the Push listener) |
kafka |
ARCHIVER_GRPC_LISTEN |
Push listener bind address, on grpc |
(none) |
KAFKA_BROKERS |
Kafka broker addresses | localhost:9092 |
KAFKA_GROUP_ID |
Consumer group ID | dfe-archiver |
KAFKA_TOPICS |
Topics to consume (comma-separated); empty discovers | (discover) |
KAFKA_TOPIC_INCLUDE |
Regex patterns a discovered topic must match | (none) |
KAFKA_TOPIC_EXCLUDE |
Regex patterns that drop a discovered topic | DLQ + internal |
KAFKA_TOPIC_REFRESH_SECS |
How often discovery re-reads the broker | 60 |
KAFKA_SASL_MECHANISM |
SASL mechanism | (none) |
KAFKA_SASL_USER |
SASL username | (none) |
KAFKA_SASL_PASSWORD |
SASL password | (none) |
ARCHIVER_DESTINATION |
Output URL | (none -- idles) |
ARCHIVER_COMPRESSION_CODEC |
Compression codec | zstd |
ARCHIVER_MEMORY_LIMIT_BYTES |
Memory guard cap; 0 auto-detects from the cgroup |
0 |
ARCHIVER_MEMORY_PRESSURE_THRESHOLD |
Backpressure trigger, 0.0-1.0 | 0.8 |
ARCHIVER_VERSION_CHECK__ENABLED |
false disables the startup version check |
true |
METRICS_ADDR |
Metrics server address | 0.0.0.0:9090 |
LOG_LEVEL |
Log level | info |
The memory guard belongs to the scalo runtime and reads only the ARCHIVER_MEMORY_* variables above, never a memory: block in the config file. The other scalo sections (version_check, metrics, logger, self_regulation, scaling, worker_pool) resolve from the cascade, where the config file is the settings layer and an ARCHIVER_SECTION__KEY variable outranks it. The archiver's own sections take a single underscore; the double-underscore names belong to those scalo sections.
transport picks how records arrive. On kafka the archiver joins a consumer
group and reads the landing topics; on grpc it binds the scalo Push listener
and the previous stage sends to it point to point, which is how a deployment
with no broker still archives. A deployment sets it from the same dial that
decides the rest of the stack's transport, so it is not normally hand-authored.
The archiver starts, passes readiness and serves health and metrics with no
work to do -- no destination, or no topics and no discovery pattern. It opens
no broker connection and binds no listener while idle, reports the
work_config health component Degraded with the reason, holds the
pipeline_idle gauge at 1, and starts the moment a config change gives it
work. Config::idle_reason is the whole predicate.
transport: kafka # or "grpc" for the Push listener
kafka:
brokers:
- kafka:9092
group_id: dfe-archiver
topics: # omit to discover, filtered by topic_include
- events
- logs
sasl_mechanism: SCRAM-SHA-256
sasl_username: archiver
sasl_password: ${KAFKA_PASSWORD} # env var substitution
acknowledgements:
enabled: true # false commits at receipt
archive:
destination: s3://my-bucket/archives
# Under the routed destination; {year} {month} {day} {hour} {minute}
# {timestamp} {seq} are the only placeholders, anything else is refused.
# Each file name ends -<seq>-<writer id>.
path_template: "{year}/{month}/{day}/{hour}"
roll_size_bytes: 1073741824 # 1GB (final compressed size)
roll_interval_secs: 300 # unset: 300 while offsets are held or on the direct (grpc) transport, else 3600
buffer:
flush_bytes: 1048576 # 1 MiB: a destination's buffer flushes into its file at this size
flush_records: 100000 # ... or at this many records, whichever comes first
flush_age_secs: 60
writer_parallelism: 2 # object-store uploads at once, each holding 4 x multipart_chunk_size
compression:
codec: zstd
level: 3
routing:
mode: expression # or "topic"
expression_fields:
- org_id
- event_typeThe archiver supports hot-reloading configuration without restart via SIGHUP or file polling (5-second interval).
Hot-reloaded (takes effect on next batch):
kafka.batch_size- re-read once per receive
Requires pod restart - everything else. The pipeline snapshots the config at
startup, so transport, kafka.*, grpc.*, archive.*, buffer.*,
routing.*, compression.* and dlq.* keep their startup values until the
process restarts. A reload of one
of those logs a warning naming the sections that changed, and the security event
says a restart is needed rather than reporting a reload that reached nothing.
archive:
destination: file:///var/data/archivearchive:
destination: s3://bucket-name/prefix
s3:
endpoint: http://minio:9000
region: us-east-1
access_key_id: ${AWS_ACCESS_KEY_ID}
secret_access_key: ${AWS_SECRET_ACCESS_KEY}archive:
destination: gs://bucket-name/prefix
gcs:
project_id: my-project
service_account_key: /path/to/key.jsonarchive:
destination: az://container-name/prefix
azure:
account_name: myaccount
account_key: ${AZURE_STORAGE_KEY}On kafka, a record's offset is committed only once the archive file holding it is durable: the store confirmed the upload, or the local file and its directories synced to disk. An object-store file is written to <buffer.spool_dir>/uploads and uploaded whole when it rolls -- roll_size_bytes or roll_interval_secs, whichever comes first -- so the store is never on the write path. scalo's Kafka transport tracks every offset it hands out and commits each partition only up to its lowest offset not yet released, so one destination's roll never commits past a record another destination still holds.
A failed upload keeps the staged file and its held offsets, and retries in the background with exponential backoff and jitter, from 0.5 s up to 60 s between attempts, until the store takes it. Each attempt is a fresh upload, and a failed one is aborted so it leaves no parts in the store. An outage of any length costs disk and consumer lag, never records. A staged copy that cannot be read back is moved to <buffer.spool_dir>/uploads/quarantine and its records are read again after a restart. Intake pauses (the Kafka partitions, through the self-regulation gate) as held records near a quarter of the memory limit at 32 bytes each, or staged files near 8 GiB. A store outage becomes consumer lag, not memory or disk growth. A refusal the store gives for the object itself -- a key past the 1024-byte limit, an entity too large -- is permanent. The file's records then go to the DLQ, each as its own entry, and only with the DLQ off, or for a record too large for any DLQ backend, are they dropped, counted in messages_dropped_total with the reason logged. messages_archived_total counts records when the store confirms their file, messages_written_total when they go into an open file.
Each held record costs about 32 bytes of memory until its file is durable. So while offsets are held an unset roll_interval_secs is 300 rather than 3600: at 10k records/s that is 96 MB instead of 1.15 GB. A configured value always wins, and one above 900 logs a startup warning.
If the archiver is killed, every record not yet in a durable file is read again after the restart: up to one roll interval of intake plus whatever was still uploading, as duplicates, never as loss. The restart removes the staged files the killed process left, since their records are read again.
A local write that fails in a way that can clear -- a full or failing disk -- is held and retried with backoff and jitter, its batch and offsets kept, and the archiver receives nothing more until it lands. A disk that stays broken shows as consumer lag, not lost records. Once shutdown begins, a write still failing gets 10 s, then its records are left for the restart to read again.
A batch refused for good -- by the store, or by the codec -- goes to the DLQ one record per entry, and only a write the DLQ confirms releases a record's offset. A record is dropped with its reason, counted in messages_dropped_total, if the DLQ is off or the record alone is too large for any DLQ backend once base64 grows it by a third. A DLQ write that fails, or a file that cannot complete, holds the commit below it, and nothing reads it again while the process runs, so the archiver drains and exits non-zero and the restart reads it again.
Expression routing never parses a record nested past 64 levels, because sonic-rs recurses once per level with no limit and about 20,000 levels overflow a 2 MiB worker stack. The record goes to the DLQ as it arrived, under its topic, with the reason payload nesting exceeds the maximum parse depth of 64, and with the DLQ off it is dropped and counted in messages_dropped_total{reason="too_deep"}.
kafka.acknowledgements.enabled: false commits at receipt instead, so a kill loses what the open files and buffers held. A record there that is written nowhere -- its file could not complete, or the DLQ could not take it -- is counted in messages_dropped_total{reason="unreplayable"}, and so is one on grpc.
On grpc the listener answers each push once its records are queued: a record is released only when its file is durable, long after any sender's deadline. So the archive copy on the direct path is at-most-once -- a kill loses what the queue and the open files held. A graceful stop still writes every record it answered: the listener closes first, its queue is drained into the files, then the files complete and upload. The drain gives uploads 20 s. A file still uploading then keeps its staged copy: the restart uploads it on grpc or with acknowledgements off, and removes it on kafka, where its records are read again. A staged file no process completed is moved to quarantine rather than uploaded.
Every file name carries a component unique to the writer, so two replicas writing one destination in one window never complete an upload onto the same key.
pipeline_delivery_guarantee{guarantee, reason} reads 1 for the guarantee in force: at_least_once/confirmed on kafka into an object store, at_least_once_local/sink_confirms_locally on kafka into a local path, and best_effort with acks_disabled or, on grpc, sink_cannot_confirm.
The only files on local disk are archive files: a local destination's open files, and under <buffer.spool_dir>/uploads each object-store file until the store takes it.
- The archiver refuses to start with less than 1 GiB free on the
buffer.spool_dirvolume. - Intake pauses as staged bytes near 8 GiB, under the spool volume's 10 GiB
emptyDirlimit, so a store outage becomes consumer lag. - The quarantine directory holds at most 1 GiB. A file past it is removed and counted in
staged_files_removed_total{reason="quarantine_full"}. - A local write that fails on a full disk is held and retried, never dropped.
Neither limit is configured. staged_bytes and uploads_pending show the disk in use.
Prometheus metrics at http://0.0.0.0:9090/metrics (configurable). Three layers:
Platform (dfe_*): records_received_total, records_delivered_total, transport_sent_total, scaling_pressure - auto-emitted by scalo.
Metric groups (dfe_archiver_*): AppMetrics (received/processed/error counts, memory, config reloads), BufferMetrics (bytes, records, flush duration), ConsumerMetrics (poll duration, offsets committed), SinkMetrics (write duration/errors by backend). The archiver sets no lag, partition or rebalance value on ConsumerMetrics: per-partition lag and the rebalance count are the rdkafka_* gauges below.
Archiver-specific (dfe_archiver_*): files_created_total, files_closed_total, archive_roll_total{trigger}, writer_evictions_total, compression_ratio, compression_duration_seconds, events_per_second, hot_buffers_active, hot_buffer_evictions_total, unique_destinations, routing_fallback_total{field}, kafka_commit_errors_total.
Archiver delivery (dfe_archiver_*): messages_written_total, messages_archived_total, messages_dlq_total, messages_dropped_total{reason} (refused, dlq_too_large, too_deep, unreplayable), uploads_pending, staged_bytes, staged_files_recovered_total, staged_files_quarantined_total{reason}, staged_files_removed_total{reason}, and pipeline_delivery_guarantee{guarantee, reason}.
rdkafka stats (rdkafka_*): broker_rtt_avg_seconds{broker}, topic_partition_consumer_lag{topic,partition}, consumer_rebalance_count - published by scalo from the consuming transport's own librdkafka statistics, every 5 s.
Health endpoints: /healthz (liveness), /readyz (readiness via HealthRegistry)
- Rust 1.94+ (edition 2024)
- Docker (for local Kafka via
dfe-docker) hyperi-cifor CI validation
The object-store e2e tests are deliberately not #[ignore]d, so the default
run starts Azurite, fake-gcs-server and LocalStack through testcontainers and
needs a Docker daemon.
# Unit, integration and object-store e2e tests
cargo nextest run
# Adds the Kafka, minio:// and GCS-credential tests, which need a live stack
TEST_MODE=docker cargo nextest run -- --ignored
# The same ignored tests against the remote dev stack
TEST_MODE=remote cargo nextest run -- --ignored
# Pre-push validation
hyperi-ci checkcargo build --release
# Features: jemalloc (in `full`), transport-memory, pgo-driver
cargo build --release --features jemallocdfe-archiver # Run the archiver service
dfe-archiver --config config.yaml # With explicit config
dfe-archiver emit-contract # Print deployment contract JSON
dfe-archiver emit-dockerfile # Generate Dockerfile
dfe-archiver emit-helm chart/ # Generate Helm chart
dfe-archiver version # Version infoThis software is licensed under the Business Source License 1.1 (BUSL-1.1). See LICENSE for details.
Copyright (c) 2026 HyperI Pty Ltd
dfe-archiver is the sink at the end of the DFE pipeline. Records arrive from Kafka or the scalo Push gRPC listener, are buffered per destination, compressed, and written as rolling files to file, s3, gs, az or MinIO. It parses nothing, enriches nothing, routes between no topics and answers no queries. Not a library either -- all three crates set publish = false and the artefact is a container image on ghcr. With nothing configured it starts and idles rather than crash-looping -- see idling until configured.
| Path | What it holds |
|---|---|
crates/core/ |
Config types, the rolling writer, the tiered buffer, routing, codecs, the storage trait. No I/O, no metrics. |
crates/io/ |
Kafka and gRPC transports, and the object-store backends behind create_backend. The cold-build bottleneck, none of it ours. |
crates/archiver/ |
The binary. archiver.rs the pipeline loop, main.rs the scalo ServiceApp wiring, metrics.rs the Prometheus surface, contract.rs the deployment contract. |
crates/archiver/src/config/loader.rs |
init_cascade -- the call that makes ARCHIVER_*__* reach anything. |
crates/archiver/src/bin/pgo_driver.rs |
The second [[bin]], gated on required-features = ["pgo-driver"]. The gate keeps it out of the release payload: hyperi-ci skips feature-gated binaries when packaging, and verifies BOLT on the first packaged binary only. |
Dockerfile |
Generated from contract.rs. Its header names the generator and schema version. |
docs/architecture.md |
Crate map and the one-way rules. docs/DESIGN.md is the pipeline design. |
.hyperi-ci.yaml |
PGO and BOLT settings, and which gates block versus warn. .config/nextest.toml and scripts/pgo-workload.sh are what it drives. |
| Command | What it covers |
|---|---|
hyperi-ci check |
The pre-push gate. make check is the same. |
cargo nextest run |
Unit, integration and object-store e2e. Those are not #[ignore]d, so it starts Azurite, fake-gcs-server and LocalStack via testcontainers and needs a Docker daemon. |
TEST_MODE=docker cargo nextest run -- --ignored |
Adds the Kafka, minio:// and GCS-credential tests, which need a live stack (docker compose -f docker-compose.dev.yaml up -d). TEST_MODE=remote runs them against the remote dev stack. |
cargo deny check advisories |
Advisories and yanked crates. quality.rust.audit and quality.rust.deny are blocking, so a red advisory fails hyperi-ci check and CI too. |
Green says less than it looks, two ways:
- the
defaultnextest profile setsretries = 0deliberately, because hyperi-ci never selects--profile ci. A retry there hides a real intermittent. A test pins it. - the
pushtrigger ignoresdocs/**and**.md, so a docs-only push runs nothing.pull_requesthas no such filter.
| Don't | Do | Why |
|---|---|---|
| Commit the Kafka offsets of a batch once it is written into a file. | Hold them on the file and release them when it completes. The armed consumer commits each partition up to its lowest offset not yet released. | A file is an upload in progress until it rolls, and a kill abandons it. Default routing is expression on org_id, so one partition fans out to several files that complete at different times (#82). |
Build a second MetricsManager in a test and assert on render(). |
Share one manager, assert on the delta. | set_global_recorder succeeds once per process. Later managers keep the existing recorder and render an empty string. nextest forks per test and hides it, cargo-llvm-cov runs one process and does not (#84). |
Trust cargo update -p rustls to clear the advisory. |
cargo update -p rustls --precise 0.23.45, then build both arches. |
Plain -p stops at 0.23.43, still vulnerable, because 0.23.45 needs aws-lc-rs to move too. --precise drags aws-lc-sys 0.41 to 0.45, which compiles C (#85). |
Put memory:, metrics: or scaling: in the config file. |
Set them as ARCHIVER_<SECTION>__<KEY> env vars. |
scalo builds the memory guard, metrics listener and scaling engine before the config file is read. Those blocks once parsed and validated while reaching nothing. |
Point scalo at a local path to try an unreleased change. |
Keep the crates.io range. Read the local clone instead. | A path override builds against uncommitted work, so the release does not reproduce. |
Hand-edit Dockerfile. |
dfe-archiver emit-dockerfile > Dockerfile. |
It is generated from the deployment contract, and a scalo release can move the generator or its schema version. |
| Give PGO a port check or a startup probe as its workload. | Drive consume, route, compress and write for 60s or more. | Shallow workloads bias the profile toward startup paths and give NEGATIVE gains. Measured on dfe-loader and dfe-receiver. |
Inbound:
- hyperi-io/scalo-rs, cargo dependency. The workspace declares one
scalorange in[workspace.dependencies]and all three crates inherit it, so a scalo release is a range check, a bump and a rebuild. - hyperi-io/scalo-rs again, generator, lockstep.
Dockerfilecomes fromscalo::deployment::generate_dockerfile()over this repo's contract, so a generator or schema-version change means regenerate and commit the diff.
Outbound -- the repo a change here breaks:
- hyperi-io/dfe-infra, image pin, lockstep.
helm/charts/dfe-archiver/Chart.yamlcarries theappVersion, drift-checked against dfe-infra'sversions.yaml. A release here is not deployed until that pin moves. The chart lives there --dfe-archiver emit-helmcan write one, this repo commits none.
Regenerate both from dfe-infra: python3 scripts/dfe-stack suite --consumer dfe-archiver, and again with --producer.