Pulse is a tiny real-time data streaming framework (mini Flink/Beam) in Rust. Async (Tokio), pluggable operators, state, and I/O. Fast, modular, local-first.
-
Updated
Oct 29, 2025 - Rust
Pulse is a tiny real-time data streaming framework (mini Flink/Beam) in Rust. Async (Tokio), pluggable operators, state, and I/O. Fast, modular, local-first.
Stream processing for Kafka pipelines that must stay correct when processes die. SQL, event time, keyed state and exactly-once delivery in modern C++: one process on a laptop, embedded, or a cluster. No JVM.
Real-time decisions on an electricity distribution network. Deterministic code owns whether a window is closed; models only forecast, score and rank. Seven claims proved offline and against a deployed AWS estate — 0 rows published early of 3,779, 41 gate mutations refused, all six erasure legs confirmed.
Deterministic reconciliation engine with typed breaks, replayable decisions, and audit trail.
Delayed conversion labels: freshness, correctness, corrections, processing lag, and retained state.
Interactive F# lab for watermark lag, late events, idle partitions, incomplete windows, and delayed alerts.
Exactly-once streaming pipeline: Kafka KRaft to Flink to ClickHouse. Zero data loss verified by killing Flink mid-stream, not by assertion.
Apache Flink and PyFlink streaming lab demonstrating event time, managed keyed state, checkpoints, controlled TaskManager failure, and validated state recovery.
Event-time stream processor for out-of-order F1 telemetry. Watermark-based windowing, late-data amendment, and checkpoint recovery verified against offline ground truth at 100% match vs 97.9% for arrival-order processing. Python, asyncio, Postgres, FastAPI.
Cross-engine output validation for event-time streaming pipelines.
Kafka -> Spark Structured Streaming -> Delta, with event-time windowing and watermarking made runnable: a dependency-free reference implementation you can execute without a JVM
Evidence-bound ad click aggregation reference: event-time windows, fenced workers, immutable publications, explicit corrections, and PostgreSQL authority.
Kafka streaming lab demonstrating partitioned topics, duplicate handling, late events, DLQ routing, and measured end-to-end latency.
Real-time anomaly detection on payment streams with PyFlink: event-time windows, keyed state, watermarks
Durable event-time aggregation with SQLite WAL, transactional deduplication, monotonic watermarks, late-event accounting, recovery tests, and reproducible benchmarks.
Kafka-compatible streaming demo for real-time birding intelligence: event-time windows, late-event handling, stateful deduplication, hotspot metrics, and target-species alerts.
Watermark Lens - diagnostic tooling for Apache Flink watermarks, windows, and late data
BTCUSDT top-of-book microstructure signal study for short-horizon event-time midprice direction prediction.
Measures what a replay of a stateful streaming job cannot reproduce. 62% of the output rows survive reprocessing, the session table goes from 205 rows to 12, and the bound you can compute from the log is four times too small because it is about the wrong pair of runs.
To associate your repository with the event-time topic, visit your repo's landing page and select "manage topics."