FutureQ¶
A high-performance, distributed delayed-message queue broker written in Go.
FutureQ lets producers publish messages with a relative delay and guarantees reliable dispatch to consumers when the delay expires. It combines durable embedded storage, Raft-based replication, and bidirectional gRPC streaming into a single, easy-to-operate binary.
:octicons-arrow-right-24: Get started :octicons-mark-github-24: GitHub repository :octicons-file-code-24: Protocol definition
Features¶
- β± Delayed messaging β enqueue a message with a
delay_ms; it becomes visible to consumers only after the delay expires. Optional per-messagettl_msdiscards messages nobody consumed. - πΎ Durable storage β disk-backed by Pebble (CockroachDB's LSM store) with a time-optimized 24-byte key schema, a pure in-memory mode, and an alternative bbolt engine.
- π High availability β multi-node replication via Dragonboat (multi-group Raft): an event shard for messages plus a metadata group for cluster topology.
- β Dynamic membership β nodes join and leave a running cluster over gRPC (
JoinCluster/LeaveCluster) with catch-up before voting. No static bootstrap list needed after the first node. - π₯ Consumer groups & topics β topic-based routing with fan-out across groups and round-robin dispatch within a group. "Universal" consumers (no group) each receive every message.
- β At-least-once delivery β in-flight tracking with automatic re-dispatch of unacknowledged messages; batched deletes (a single Raft proposal per batch) amortize LSM tombstone costs.
- βοΈ Tunable durability β publish acknowledgements at three levels: quorum commit, leader-local persistence (fast path), or fire-and-forget β with a broker-side minimum enforced by policy.
- π Observability β Prometheus metrics endpoint and structured logging (zap) out of the box.
Technology stack¶
| Concern | Choice |
|---|---|
| Language | Go 1.26 |
| Storage | Pebble (LSM tree) Β· bbolt (alternative) Β· in-memory |
| Consensus | Dragonboat v4 (multi-group Raft) |
| Transport | gRPC (bidirectional streaming) |
| Metrics | Prometheus |
| Protocol | futureq-io/protocol (Protobuf, v0.2.0) |
How it works¶
ββββββββββββββββ PublishStream ββββββββββββββββββββββββββββ
β Producer β ββββββββββββββββββΆβ β
ββββββββββββββββ β FutureQ Node β
β β
ββββββββββββββββ Subscribe β ββββββββββββββββββββββ β
β Consumer β βββββββββββββββββ β β gRPC API (8443) β β
ββββββββββββββββ β βββββββββββ¬βββββββββββ β
β βΌ β
ββββββββββββββββ Raft (50005) β ββββββββββββββββββββββ β
β Other Nodes β βββββββββββββββββΆ β β Dispatcher / Hub β β
ββββββββββββββββ β β Deleter Β· Janitor β β
β βββββββββββ¬βββββββββββ β
ββββββββββββββββ Prometheus β βΌ β
β Metrics β βββ (9090) βββββββ β Pebble + Raft log β β
ββββββββββββββββ β ββββββββββββββββββββββ β
ββββββββββββββββββββββββββββ
- Writes go through Raft as one log entry per batch (or straight to Pebble in standalone mode) and are stored under time-bucketed keys for efficient expiry scans.
- Reads are push-based: the dispatcher scans for due messages and routes them to connected consumers through a hub using a round-robin strategy per consumer group.
- Acks are batched by a deleter and committed as a single Raft proposal, keeping write amplification low.
Delivery is at-least-once β consumers should be idempotent.