Skip to content

Latest commit

Β 

History

85 Commits

Folders and files

NameName
Last commit message
Last commit date
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 

AeroStream

License Docker Image Go Rust Protocol UI Docs Status

AeroStream is an ultra-high-performance, distributed event-streaming and messaging engine engineered for extreme throughput, microsecond latencies, and modern multi-cloud workloads.

Built with a Dual-Engine Architectureβ€”pairing a resilient Go-based Raft control plane with a zero-copy Rust-based storage and networking data planeβ€”AeroStream delivers next-generation event streaming with 100% Kafka wire-protocol compatibility, built-in multi-cloud tiered storage, schema governance, stream transforms, and an integrated Web Console UI.


⚑ Benchmark Summary (OpenMessaging Benchmark)

AeroStream was measured with the vendor-neutral Linux Foundation OpenMessaging Benchmark (OMB) framework on an AWS c6id.2xlarge (8 vCPU, 16 GiB RAM, local PCIe Gen4 NVMe SSD). One broker, 1 topic, 32 partitions, 1,024-byte messages, 8 producers, 8 consumers, acks=1, two rounds per workload:

1. Kafka Wire Protocol (:9092) β€” 100% Drop-in Compatibility

Standard Kafka clients (Python, Java, Go, .NET, Node.js) connect directly to port 9092 with zero code changes:

Offered Load Publish Rate $p_{50}$ Latency $p_{95}$ Latency $p_{99}$ Latency $p_{99.9}$ Latency Broker CPU Errors
100,000 msg/s (fixed) 100,000 msg/s (97.7 MB/s) 0.7 ms 1.2 ms 1.4 ms 2.3 ms 14% 0
200,000 msg/s (fixed) 200,000 msg/s (195.5 MB/s) 0.7 ms 1.3 ms 1.7 ms 3.0 ms 22% 0
Maximum Rate (unthrottled) 271,350 msg/s (265.0 MB/s) 105.1 ms 812.5 ms 1,104 ms 1,376 ms 55% 0

2. Native Protocol (:9091) β€” Ultra-Low Latency with Paced Writeback

AeroStream's native binary protocol (port 9091) uses a 7-byte framing header and dedicated SDKs. By implementing paced page-cache writeback (sync_file_range / posix_fadvise every 8 MiB), AeroStream eliminated kernel flusher stalls, slashing $p_{99}$ tail latency by 48×–73Γ— and dropping broker CPU by over 56%:

Workload Metric Original Baseline Final (Writeback Fix) Improvement
100,000 msg/s $p_{50}$ / $p_{95}$ 1.2 ms / 2.7 ms 0.7 ms / 1.2 ms 2.3Γ— lower latency
(1 KB payloads) $p_{99}$ / $p_{99.9}$ 63.2 ms / 94.2 ms 1.3 ms / 1.8 ms 48Γ— lower tail latency ($p_{99}$)
Broker / load-gen CPU 97% / 88% 32% / 47% 67% less broker CPU
200,000 msg/s $p_{50}$ / $p_{95}$ 1.8 ms / 70.7 ms 0.8 ms / 1.3 ms 54Γ— lower $p_{95}$ latency
(1 KB payloads) $p_{99}$ / $p_{99.9}$ 109.9 ms / 145.8 ms 1.5 ms / 3.8 ms 73Γ— lower tail latency ($p_{99}$)
Broker / load-gen CPU 96% / 97% 42% / 58% 56% less broker CPU
Maximum rate Publish throughput 244,385 msg/s (238 MB/s) 287,428 msg/s (280.7 MB/s) +18% higher throughput
(unthrottled) Saturation $p_{99}$ 1,009 ms 149 ms 85% lower queueing tail
Broker / load-gen CPU 95% / 96% 67% / 51% 29% lower broker CPU at saturation

(Reproduced in October 2026: 100k $p_{99}$ 1.2 ms, 200k $p_{99}$ 1.4 ms, max throughput 287,459 msg/s)

3. Head-to-Head: Kafka Wire Port (:9092) vs Native Port (:9091)

Feature / Metric Kafka Wire Protocol (:9092) AeroStream Native Protocol (:9091)
Protocol Overhead Full Kafka Header v2 + RecordBatch Minimal 7-Byte Header (0xAE 0x01)
Client Compatibility Any Kafka client (Java, Python, Go, Node, .NET) Native SDKs (Go, Rust, Java, .NET, Node.js)
200,000 msg/s $p_{99}$ Latency 1.7 ms 1.5 ms (1.4 ms reproduced)
Max Sustained Throughput 271,350 msg/s (265.0 MB/s) 287,428 msg/s (280.7 MB/s) (+5.9%)
Saturation $p_{99}$ Queueing Tail 1,104 ms 149 ms (86.5% reduction)
Broker CPU at 200k msg/s 22% (pinned cores) 42% (pinned cores)

πŸ“Š Explore Full Benchmark Reports & Reproduction:

  • πŸ“ˆ Website Benchmark Page: Per-run metrics, CPU utilization, test environment, and caveats.
  • πŸ”¬ Benchmark Report: EC2 results, resource-capped container runs, design notes, and partition-density measurements.
  • ☁️ EC2 Benchmark Scripts: One command creates the machine, runs OMB, copies the results back, and destroys the machine.

πŸš€ Key Features

  • Shard-per-Core Zero-Contention Engine:
    • Deterministic partition-to-shard mapping ($S = \text{hash}(\text{topic}, \text{partition}) \pmod N$) with thread-to-CPU affinity pinning via libc::sched_setaffinity.
    • Lock-free actor message passing via flume::unbounded channelsβ€”eliminating cross-core mutex locks, atomics, and thread migrations on hot produce/consume paths.
  • Extreme Memory & Storage Optimizations:
    • In-place base offset patching directly on disk (write_all_at), completely bypassing multi-megabyte heap reallocations in $\mathcal{O}(1)$ time.
    • Paced page-cache writeback via Linux sync_file_range(2) and posix_fadvise(2) (pacing dirty flushes every 8 MiB), eliminating OS writeback stalls in memory-constrained containers.
    • Zero-copy cold tiering via hard links (fs::hard_link), decoupling hot partition log rollover from object store network latency.
  • 100% Kafka Wire Protocol Compatibility:
    • Native listener on port 9092 supporting 34+ Kafka API keys across produce, fetch, metadata, consumer groups, schemas, and ACLs.
    • High-performance Fetch long polling with lazy tokio::sync::Notify registration and lockless out-of-lock disk I/O.
    • Enterprise SASL authentication (PLAIN and SCRAM-SHA-256) and Two-Phase Commit (2PC) Transactions.
    • Drop-in compatibility with standard Kafka client ecosystems (kafka-python, librdkafka, kafka-go, Confluent.Kafka, Java / Spring Kafka).
  • Dual-Engine Decoupled Architecture:
    • Control Plane (Go): Distributed Raft consensus, automated partition leadership elections, dynamic cluster membership, 2-second broker heartbeats with dynamic config piggybacking, and gRPC coordination.
    • Data Plane (Rust): Tokio async runtime, CPU core pinning, kernel zero-copy sendfile(2) socket transfers, and memory-mapped (mmap) offset indexing.
  • Multi-Cloud Tiered Storage:
    • Hot partition segments on fast local NVMe/SSD.
    • Transparent, non-blocking background offload to AWS S3 / MinIO, Google Cloud Storage (GCS), Azure Blob Storage, or network filesystem mounts.
  • Iceberg-Native Topics:
    • Topics can write directly into Apache Iceberg tables (Parquet/Avro) for lakehouse-native analytics without a separate sink connector.
  • Share Groups (KIP-932 Queue Semantics):
    • Cooperative, queue-like consumption where multiple consumers acquire/acknowledge individual records from the same partition without exclusive assignment.
    • Per-record delivery-attempt limits, lock timeouts, and dead-letter-queue (DLQ) forwarding for records that exhaust retries.
  • Built-in Schema Registry:
    • Confluent-compatible REST API on /subjects, /schemas, and /compatibility.
    • First-class support for Avro, Protobuf, and JSON Schema with BACKWARD, FORWARD, and FULL compatibility validation.
  • In-Broker Stream Processing & Transforms:
    • Dynamic Stream Processing Engine (/api/streams), inline real-time filtering, PII data masking (MASK_PII), JSON schema transformation, and WASM runtime.
  • Enterprise Security & Granular RBAC:
    • Role-based access control (SUPER_ADMIN, OPERATOR, PRODUCER, CONSUMER, AUDITOR).
    • Granular topic, consumer group, and cluster ACLs with prefix and wildcard pattern matching.
  • Modern Web Console UI:
    • Sleek Angular management console with dark/light themes, dynamic cluster topology visualizer, live message inspector, real-time consumer lag monitoring (/api/lag), schema registry browser, and policy simulators.

🐳 Quick Start with Docker

Run the full AeroStream stack (Controller, Zero-Copy Broker, and Web Console) using the official container image:

docker run -d --name aerostream \
  -p 9091:9091 -p 9092:9092 -p 9001:9001 -p 8001:8001 -p 7001:7001 \
  -v aerostream_data:/data \
  quay.io/gradientgeeks/aerostream:latest

πŸ“– Deployment Quickstart Guides:

  • 🐳 Docker Quickstart Guide: 30-second local setup with single all-in-one container, port mapping, and client samples.
  • ☸️ Kubernetes Quickstart Guide: Production deployment using standard kubectl manifests, headless services, StatefulSets, and automated zero-downtime draining.
  • ⎈ Helm Quickstart Guide: Official Helm v3 chart installation, values customization, S3 tiered storage, and rack-aware zone placement.

Accessing Endpoints:

  • Web Console UI: http://localhost:9001/aerostream/console
  • Kafka Wire Protocol: localhost:9092 (point any Kafka producer/consumer here)
  • Native TCP Data Plane: localhost:9091
  • REST Management API & Schema Registry: http://localhost:9001
  • gRPC Control Plane: localhost:8001

πŸ— Architecture Overview

AeroStream achieves its performance through strict architectural decoupling and hardware-aligned execution:

  • Go Control Plane (Ports 9001 & 8001): Drives distributed consensus via HashiCorp Raft, serves the dynamic Schema Registry, manages enterprise RBAC / ACL policies, Stream Processing Engine, and connector runtimes. Brokers maintain a 2-second heartbeat loop with dynamic configuration piggybacking and cluster topology synchronization.
  • Rust Shard-per-Core Data Plane (Ports 9091 & 9092): Partitions are mapped deterministically across dedicated shard worker threads pinned to specific CPU cores via libc::sched_setaffinity. Each shard runs an isolated event loop processing requests via lock-free actor message passing (flume::unbounded), bypassing cross-core locks and atomics.
  • Kernel & Memory Pacing: Utilizes Linux kernel zero-copy sendfile(2) DMA transfers, memory-mapped (mmap) sparse index lookup, in-place base offset patching directly on disk (write_all_at), and paced dirty page writeback (sync_file_range(2) + posix_fadvise(2)).
  • Multi-Cloud Tiered Storage: Automatically rolls sealed 128 MB log segments into an asynchronous offloader queue via zero-copy hard links (fs::hard_link), persisting them to AWS S3, MinIO, Google Cloud Storage, or Azure Blob without blocking producer ingestion.

πŸ“– Deep-Dive Architecture Specifications & Diagrams:


πŸ›  Local Development & Building from Source

Prerequisites

  • Go: 1.24 or higher
  • Rust: 2021 or 2024 edition (Cargo & Rustc)
  • Node.js: 20+ and npm (for Web Console)
  • Make

1. Build all components

make build

This builds:

  • go-controller/bin/controller: Go cluster controller and Raft consensus daemon.
  • rust-broker/target/release/rust-broker: Rust zero-copy storage broker.
  • client/bin/client: Go CLI client and benchmark utility.

2. Start a Local 5-Node Cluster

Start 3 Go Raft controllers and 2 Rust storage brokers locally:

make start

3. Check Cluster Health

make status

4. Interactive CLI Usage

# Create a topic with 3 partitions and replication factor 2
./client/bin/client create-topic orders 3 2

# Produce messages using native protocol
./client/bin/client produce orders 0 '{"order_id": "ORD-1001", "amount": 99.50}'

# Consume messages
./client/bin/client consume orders 0 0 --follow

5. Standard Kafka Client Usage

Produce and consume via standard Kafka tools or Python:

from kafka import KafkaProducer, KafkaConsumer

# Produce via standard Kafka protocol to port 9092
producer = KafkaProducer(bootstrap_servers=['localhost:9092'])
producer.send('orders', b'{"order_id": "ORD-1002", "status": "CONFIRMED"}')
producer.flush()

# Consume via standard Kafka protocol
consumer = KafkaConsumer('orders', bootstrap_servers=['localhost:9092'], auto_offset_reset='earliest')
for message in consumer:
    print(f"Received: offset={message.offset} value={message.value.decode('utf-8')}")
    break

6. Official .NET / C# Client Usage (Confluent.Kafka)

AeroStream is fully compatible with .NET 8 / 9 via official Confluent.Kafka:

using Confluent.Kafka;

var config = new ProducerConfig {
    BootstrapServers = "localhost:9092",
    EnableIdempotence = true // KIP-98 exactly-once semantics
};
using var producer = new ProducerBuilder<string, string>(config).Build();
var deliveryReport = await producer.ProduceAsync("orders", new Message<string, string> {
    Key = "order-1002",
    Value = "{\"status\": \"CONFIRMED\"}"
});
Console.WriteLine($"Delivered to {deliveryReport.TopicPartitionOffset}");

Full .NET test suite located in examples/dotnet-app/.


πŸ§ͺ Testing

Run the test suite across all sub-systems:

# Test Rust Storage Engine & Kafka Protocol
cd rust-broker && cargo test

# Test Go Controller, Consensus & Schema Registry
cd go-controller && go test -v ./...

# Test CLI Client
cd client && go test -v ./...

πŸ“œ Documentation & Guides

AeroStream provides comprehensive, interactive web documentation built with MkDocs Material (hosted at https://aerostream.gradientgeeks.com/docs/) and maintained in the aerostream-docs repository:


πŸ“„ License

This project is licensed under the Apache License 2.0 - see the LICENSE file for details.

Releases

Sponsor this project

Packages

Used by

Contributors

Languages