Streaming
spark/jobs/streaming_ingest.py
Reads Kafka-compatible events and writes checkpointed Parquet outputs.
StreamFlow Phase 1
A Kafka-compatible broker, Spark Structured Streaming ingest job, bounded Spark summary job, Airflow DAG, and data quality contract for event-driven data engineering practice.
Why it matters
The project demonstrates how event data becomes trustworthy enough for analytics, RAG corpora, and model evaluation workflows.
spark/jobs/streaming_ingest.py
Reads Kafka-compatible events and writes checkpointed Parquet outputs.
src/streamflow/quality.py
Separates valid events from rejected records with explicit reasons.
airflow/dags/streamflow_daily_summary.py
Coordinates bounded batch summary work instead of owning the streaming loop.