Platform Overview & Quickstart¶
What is AeroStream?¶
AeroStream is an open-source, ultra-high-throughput, cloud-native distributed event streaming platform engineered around a Dual-Engine Architecture. It combines the distributed consensus stability and operational agility of a Go Raft Control Plane with the zero-copy performance and mechanical sympathy of a Rust Shard-per-Core Storage Kernel.

Drop-In Apache Kafka Compatibility & Ultra-Fast Native Streaming
AeroStream natively provides dual ingress protocols:
- Port
9092(Kafka Wire Protocol): Drop-in replacement for standard Kafka SDKs (kafka-clients,confluent-kafka,kafkajs,franz-go). Point existing apps directly to AeroStream without modifying code or schemas. - Port
9091(AeroStream Native Protocol0xAE 0x01): Ultra-lightweight 7-byte binary header with zero Kafka framing overhead, delivering sub-millisecond tail latencies via official SDKs in Go, Rust, Java, .NET, and Node.js.
Dual-Engine Architectural Rationale¶
AeroStream pairs two purpose-built runtimes, each used where it excels: Go 1.26 for the control plane (HashiCorp Raft consensus, REST APIs, Schema Registry, RBAC policies, and stream transforms) and Rust 1.98.1 (Edition 2024) for the storage data plane (zero-copy sendfile(2), hardware CRC32C, memory-mapped indices, and no garbage collection).
| Subsystem | Engine Runtime | Core Responsibilities | Performance Highlights |
|---|---|---|---|
| Control Plane | Go 1.26 (Alpine) | HashiCorp Raft Quorum, Schema Registry, RBAC ACLs, Stream Transforms, Connectors, Web Console REST API | Green Tea GC (<1 ms pause), SIMD Swiss Tables hash maps, native Kubernetes cgroup auto-tuning |
| Data Plane | Rust 1.98.1 (Edition 2024) | TCP Listeners (Ports 9091/9092), Zero-Copy Segmented Commit Log, Hardware CRC32C, Tiered Storage | Zero-copy sendfile(2) I/O, 0.7 ms median publish latency at 100k-200k msg/s, lock-free Shard-per-Core |
| Console UI | Angular 21 | Cluster Topology Visualizer, Live Message Inspector, Schema Browser, ACL Simulators | Single Page Application, dark/light themes, reactive WebSocket & REST updates |

Verified Benchmark Results (OpenMessaging Benchmark)¶
AeroStream was evaluated using the official Linux Foundation OpenMessaging Benchmark (OMB) suite through its Kafka wire protocol port (9092) on an AWS c6id.2xlarge instance (8 vCPUs, 16 GiB RAM, local NVMe SSD):
| Workload Target | Actual Publish Rate | Publish \(p_{50}\) | \(p_{95}\) | Publish \(p_{99}\) | Publish \(p_{99.9}\) | Broker Cores Busy | Errors |
|---|---|---|---|---|---|---|---|
| 100,000 msg/s (fixed) | 100,000 msg/s (97.7 MB/s) | 0.7 ms | 1.2 ms | 1.3 ms | 1.8 ms | 32% | 0 |
| 200,000 msg/s (fixed) | 200,000 msg/s (195.5 MB/s) | 0.8 ms | 1.3 ms | 1.5 ms | 3.8 ms | 42% | 0 |
| Maximum Rate (unthrottled) | 287,428 msg/s (280.7 MB/s) | — | — | 149 ms | — | 67% | 0 |
- Zero Tail Latency Spikes: Paced writeback brings \(p_{99}\) latency down to 1.5 ms at sustained 200,000 msg/s (down from 109.9 ms baseline), with consumers keeping pace in real time.
- Low CPU Utilization: Broker consumes only 32% CPU at 100k msg/s and 42% at 200k msg/s (down from 96–97% baseline).
- Hardware-Accelerated Throughput: Max rate reaches 287,428 msg/s (280.7 MB/s) (+18% over baseline) with saturation \(p_{99}\) dropping by 85% to 149 ms.
Quickstart: Launch AeroStream in 30 Seconds¶
Launch the complete AeroStream stack (Go Controller, Rust Storage Broker, and Angular 21 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
Port Mapping Reference¶
| Port | Protocol | Subsystem | Purpose |
|---|---|---|---|
9092 |
TCP (Kafka Wire) |
Rust Broker | Standard Kafka Client Ingress (spring-kafka, confluent-kafka, librdkafka, kcat) |
9091 |
TCP (Native) |
Rust Broker | AeroStream Native High-Speed Protocol (aerostream-sdk 0xAE 0x01 framing) |
9001 |
HTTP / REST |
Go Controller | Angular 21 Web Console, Schema Registry, REST Produce/Fetch, Management API |
8001 |
gRPC |
Go Controller | Internal cluster metadata synchronization and broker registration |
7001 |
TCP (Raft) |
Go Controller | HashiCorp Raft consensus quorum transport |
Docker Compose Configuration¶
Create a docker-compose.yml for local orchestration:
version: '3.8'
services:
aerostream:
image: quay.io/gradientgeeks/aerostream:latest
container_name: aerostream
ports:
- "9091:9091" # Ultra High-Speed Native TCP Protocol (0xAE 0x01)
- "9092:9092" # Apache Kafka Wire Protocol
- "9001:9001" # HTTP REST, Angular 21 Console & Schema Registry
- "8001:8001" # Internal gRPC Broker Control
- "7001:7001" # Raft Consensus Transport
volumes:
- aerostream_data:/data
restart: unless-stopped
volumes:
aerostream_data:
Launch with:
Verifying Cluster Health¶
Verify cluster health via the REST API:
Expected JSON output:
{
"node_id": "node1",
"raft_state": "Leader",
"raft_leader": "127.0.0.1:7001",
"brokers_count": 1,
"brokers": [
{
"id": 1,
"host": "0.0.0.0",
"port": 9091,
"kafka_port": 9092,
"active": true
}
],
"topics_count": 0,
"groups_count": 0
}
Accessing the Web Console¶
Open your browser to:
http://localhost:9001/aerostream/console
The integrated Angular 21 Web Console allows real-time inspection of cluster topology, partition states, live consumer lag, schema registry contracts, and broker performance metrics.
Client Quickstarts¶
1. Official AeroStream Native SDKs (Port 9091, 0xAE 0x01)¶
For the lowest possible latency and overhead, use the official aerostream-sdk packages:
package main
import (
"context"
"fmt"
"log"
"github.com/gradientgeeks/aerostream-sdk/go/client"
)
func main() {
c, err := client.NewClient("127.0.0.1:9091", client.WithAuthToken("secret-token"))
if err != nil {
log.Fatal(err)
}
defer c.Close()
producer := c.NewProducer()
offset, err := producer.Produce(context.Background(), "telemetry", 0, []byte("sensor-payload"))
if err != nil {
log.Fatal(err)
}
fmt.Printf("Produced record at offset %d\n", offset)
}
use aerostream_client::{AeroClient, ClientConfig};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client = AeroClient::connect(
ClientConfig::new("127.0.0.1:9091")
.with_auth_token("secret-token")
).await?;
let producer = client.producer();
let offset = producer.send("telemetry", 0, b"sensor-payload").await?;
println!("Produced record at offset {offset}");
Ok(())
}
<dependency>
<groupId>org.gradientgeeks.aerostream</groupId>
<artifactId>aerostream-client</artifactId>
<version>0.1.0-preview</version>
</dependency>
import org.gradientgeeks.aerostream.client.AeroClient;
import org.gradientgeeks.aerostream.client.AeroProducer;
import java.nio.charset.StandardCharsets;
public class Main {
public static void main(String[] args) {
try (AeroClient client = AeroClient.connect("127.0.0.1:9091", "secret-token");
AeroProducer producer = client.producer()) {
long offset = producer.send("telemetry", 0, "sensor-payload".getBytes(StandardCharsets.UTF_8));
System.out.printf("Produced record at offset %d%n", offset);
}
}
}
using System.Text;
using GradientGeeks.AeroStream.Client;
await using var client = await AeroClient.ConnectAsync(new AeroClientOptions {
BootstrapServers = ["127.0.0.1:9091"],
AuthToken = "secret-token"
});
var producer = client.CreateProducer();
long offset = await producer.SendAsync("telemetry", 0, Encoding.UTF_8.GetBytes("sensor-payload"));
Console.WriteLine($"Produced record at offset {offset}");
import { AeroClient } from '@gradientgeeks/aerostream-client';
const client = await AeroClient.connect('127.0.0.1:9091', 'secret-token');
const producer = client.producer();
const offset = await producer.send('telemetry', 0, 'sensor-payload');
console.log(`Produced record at offset ${offset}`);
await client.close();
2. Standard Apache Kafka Clients (Port 9092)¶
Point standard Kafka client libraries directly at localhost:9092:
from confluent_kafka import Producer, Consumer
# Produce records to AeroStream
producer = Producer({'bootstrap.servers': 'localhost:9092'})
producer.produce('orders', key='ORD-101', value=b'{"amount": 149.50}')
producer.flush()
# Consume records
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'orders-analytics',
'auto.offset.reset': 'earliest'
})
consumer.subscribe(['orders'])
msg = consumer.poll(1.0)
if msg:
print(f"Received: {msg.value().decode('utf-8')}")
consumer.close()
Add to application.yml:
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
acks: 1
consumer:
group-id: inventory-service
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 1. Create a 3-partition topic via Controller REST API
curl -X POST http://localhost:9001/api/topics \
-H "Content-Type: application/json" \
-d '{"name": "orders", "partitions": 3, "replication_factor": 1}'
# 2. Produce records via standard Kafka CLI tools
echo "order-101: {\"amount\": 89.50}" | kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic orders
# 3. Consume records
kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic orders \
--from-beginning
package main
import (
"context"
"fmt"
"github.com/segmentio/kafka-go"
)
func main() {
w := &kafka.Writer{
Addr: kafka.TCP("localhost:9092"),
Topic: "orders",
Balancer: &kafka.LeastBytes{},
}
defer w.Close()
err := w.WriteMessages(context.Background(),
kafka.Message{
Key: []byte("ORD-101"),
Value: []byte(`{"amount": 89.50}`),
},
)
if err != nil {
panic(err)
}
fmt.Println("Produced message successfully to AeroStream!")
}