Skip to content

Latest commit

 

History

15 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 

Repository files navigation

StreamForge ⚡

StreamForge is a distributed, durable, high-throughput message broker and event streaming platform written from scratch in Python. It provides low-latency partition append logs, zero-copy memory mapping (mmap), binary TCP framing, consumer group rebalancing, leader election, and idempotent producer semantics.

Inspired by Apache Kafka, StreamForge is designed without external dependencies to deliver maximum transparency, performance, and operational reliability.


🌟 Key Architecture & Features

  • Custom Binary TCP Protocol: Length-prefixed framing with magic byte verification, API routing, correlation IDs, and CRC32 payload checksum validation.
  • Append-Only Segment Logs: Physical .log files paired with sparse binary-search .index files leveraging mmap memory mapping.
  • Idempotent Producer Semantics: Deduplicates record re-transmissions using Producer IDs (PID) and Monotonic Sequence Numbers.
  • Consumer Group Rebalancing: Dynamic consumer membership, heartbeats, auto-commit, and range partition assignment algorithms.
  • Cluster High Availability: Automatic node peer discovery, partition leadership election, and follower async replication.
  • Embedded Web Management Dashboard: Real-time browser UI (http://localhost:8080) for cluster stats, partition inspection, and benchmark execution.

📊 Benchmark Metrics (1M Messages / 10 Partitions / 1 Broker)

StreamForge includes a built-in benchmark harness (benchmark.py) capable of evaluating broker throughput, latency profiles, and crash recovery times.

================================================================================
                STREAMFORGE DISTRIBUTED MESSAGE BROKER BENCHMARK                
================================================================================
 Target Workload : 1 Broker | 10 Partitions | 1,000,000 Messages
 Batch Size      : 2,000 records/request
 Storage Path    : /tmp/streamforge_bench_data
--------------------------------------------------------------------------------

 Throughput              : 568,937.80 msg/s  (56.97 MB/s)
 Consumer Throughput     : 26,360.59 msg/s
 P50 Latency             : 0.0015 ms
 P95 Latency             : 0.0017 ms
 P99 Latency             : 0.0028 ms
 Crash Recovery Time     : 0.006 sec (Validated 1,000,000 msgs on disk)
================================================================================

🚀 Quick Start

1. Run Unit Tests

python3 -m unittest test_streamforge.py

2. Launch Benchmark

python3 benchmark.py --messages 1000000 --partitions 10 --batch-size 2000

3. Launch Web Dashboard

python3 streamforge/dashboard.py

Open http://localhost:8080 in your browser.


📁 Repository Layout

├── streamforge/
│   ├── __init__.py       # Package metadata
│   ├── protocol.py       # Binary TCP framing & API codes
│   ├── storage.py        # Mmap append log & index engine
│   ├── broker.py         # Async TCP broker server
│   ├── producer.py       # Producer API & batching
│   ├── consumer.py       # Consumer Group coordinator
│   ├── cluster.py        # Replication & Leader Election
│   └── dashboard.py      # Management Console & Web UI
├── benchmark.py          # 1M Message Benchmark Suite
└── test_streamforge.py   # System unit test suite

📜 License

MIT License.

About

High-performance, distributed, durable event streaming and message broker built from scratch

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages