A market data pipeline with a C++ core: reconstructs an L2 order book from Nasdaq ITCH messages, aggregates OHLCV bars, publishes them to Kafka, and persists them to Postgres with idempotent upserts so replaying a partition (including a late-arriving trade that revises an already-closed bar) converges instead of duplicating or losing data.
The core data path is implemented through Kafka publishing and replay-safe PostgreSQL persistence.
| Component | State |
|---|---|
| L2 order book: add / cancel / execute / modify, top of book, OHLCV bars | working, 8 test cases |
| ITCH parser and full-session census | working, framed/decoded against generated fixtures |
| Kafka publisher | working, versioned keyed bars with idempotent librdkafka delivery |
| Postgres persistence, idempotent upserts, late-trade revisions | working, versioned migrations and transactional offsets |
cmake -B build -G Ninja
cmake --build build
ctest --test-dir build
./build/alphastreamNeeds CMake 3.20+, Ninja, and a C++23 compiler. Catch2 arrives through FetchContent, so there is
nothing to install first.
Start the local services and create the explicitly configured topic:
docker compose up -d kafka postgres
./scripts/bootstrap_topics.shBuild the optional C++ librdkafka publisher, install the Python consumer, and migrate the database:
cmake -S . -B build-kafka -G Ninja \
-DALPHASTREAM_BUILD_KAFKA=ON -DALPHASTREAM_BUILD_BENCH=OFF
cmake --build build-kafka --target publish_sample_bars
python3 -m venv build/pipeline-venv
build/pipeline-venv/bin/pip install -e pipeline
build/pipeline-venv/bin/python -m alphastream_pipeline.migrations \
--dsn postgresql://alphastream:alphastream@localhost:55432/alphastreamPublish an initial bar, a second symbol, and a late revision, then persist all three Kafka records:
./build-kafka/publish_sample_bars
build/pipeline-venv/bin/python -m alphastream_pipeline.consumer \
--dsn postgresql://alphastream:alphastream@localhost:55432/alphastream \
--max-messages 3The Kafka key is the stable numeric symbol ID, so every update for one symbol stays ordered within a
partition. The PostgreSQL consumer stores each partition's next offset in the same transaction as
the bar upsert. Replaying an already-consumed record is therefore harmless, while a revision with a
newer (revision, source_seq) replaces the prior value for that bar's natural key.
For AddressSanitizer and UndefinedBehaviorSanitizer:
cmake -B build -G Ninja -DALPHASTREAM_SANITIZE=ONsrc/ engine, versioned bar serializer, and Kafka publisher
pipeline/ PostgreSQL migrations and replay-safe Kafka consumer
tests/ Catch2 regression tests
scripts/ fetch_itch.sh pulls an ITCH session into data/
docker-compose.yml Kafka (KRaft mode) and Postgres 16 for local work
scripts/fetch_itch.sh downloads a Nasdaq BX ITCH 5.0 session: 391 MB compressed, 28.7 million
messages across 8,849 symbols. BX speaks the same protocol as the main Nasdaq feed at a fraction of
the volume, which makes it the corpus to build a parser against before pointing it at 5.6 GB
sessions. Files land in data/, which git ignores.