A production-grade, multi-tenant streaming data platform built on PySpark Structured Streaming. Ingest data from any source, validate it against a versioned schema registry, run ETL transforms, and monitor everything — all without restarting. New streams can be registered live via REST API.
- Multi-source ingestion — Kafka, file directories, HTTP polling, synthetic rate sources, and any third-party Spark connector (Kinesis, Event Hubs, Delta Lake, Pub/Sub, Iceberg)
- Schema registry with auto-evolution — schemas are inferred on first event, versioned on disk, and evolved automatically when new fields appear
- Type validation + dead-letter routing — type-mismatched events are rejected and routed to a dead-letter list with reasons; valid events proceed through the pipeline
- Chainable ETL transforms — per-dataset pipelines with filters, type casts, computed columns, and literal fields
- Multi-tenancy — schema namespaces are fully isolated per
tenant_id; multiple tenants share one Spark session - Observability — per-dataset counters, rejection rates, anomaly detection (rejection spikes, schema thrash, data gaps), and pluggable alert handlers
- Live REST API — register, start, stop, and monitor streams at runtime without restarting the engine
Every micro-batch flows through the same 5-stage pipeline regardless of source type:
Source.extract_rows()
│
▼
Producer.produce_batch() ← schema inference, registration, evolution, validation
│
▼
Pipeline.run() ← ETL transforms (filter, cast, add columns)
│
▼
Consumer.consume() ← dispatch to downstream handlers
│
▼
ObservabilityPipeline ← metrics counters + anomaly detection + alerts
stream_engine/
│
├── main.py # Entry point
│
├── engine/
│ ├── sources/
│ │ ├── base.py # Abstract Source interface
│ │ ├── rate.py # Synthetic in-process generator
│ │ ├── kafka.py # Apache Kafka
│ │ ├── file_source.py # Directory file watcher
│ │ ├── http.py # HTTP polling
│ │ └── connector.py # Generic third-party Spark connector
│ │
│ ├── stream_manager.py # Lifecycle coordinator (one query per source)
│ ├── producer.py # Schema registration + batch validation
│ ├── registry.py # Versioned schema store (on-disk JSON)
│ ├── validation.py # Type inference + schema comparison
│ ├── consumer.py # Downstream batch dispatcher
│ ├── observability.py # Metrics + anomaly detection + alerts
│ └── transformer.py # Chainable ETL pipeline DSL
│
├── config/
│ ├── streams.py # Synthetic RateSource definitions
│ ├── streams.yaml # Real Kafka / File / HTTP sources
│ ├── connectors.yaml # Third-party connector JARs + streams
│ ├── transforms.py # Per-dataset ETL pipelines
│ └── loader.py # YAML → Source object factories
│
├── api/
│ ├── app.py # FastAPI application
│ ├── state.py # Shared StreamManager reference
│ ├── models.py # Pydantic request/response models
│ └── routes/
│ ├── streams.py # Stream CRUD endpoints
│ └── metrics.py # Observability endpoints
│
└── schema_registry/
└── datasets/{tenant_id}/{dataset}/ # v1.json, v2.json, metadata.json
- Python 3.11+
- Java 21 (OpenJDK 21 recommended — PySpark 4.x requires
jdk.incubator.vector)
pip install -r requirements.txtrequirements.txt
pyspark==4.0.0
fastapi>=0.111.0
uvicorn[standard]>=0.29.0
pydantic>=2.7.0
pyyaml>=6.0.1
requests>=2.31.0
On macOS with Homebrew:
brew install openjdk@21
TheJAVA_HOMEis set automatically inmain.pyto/opt/homebrew/opt/openjdk@21. Adjust if your JDK is elsewhere.
python3 main.pyOn first run, Spark downloads any enabled connector JARs via Maven (cached locally for subsequent runs). The engine then starts all configured sources and the REST API.
Stream Engine started.
[StreamManager] started 'orders' (tenant=default)
[StreamManager] started 'clickstream' (tenant=default)
[StreamManager] started 'iot_sensors' (tenant=default)
[StreamManager] started 'payments' (tenant=default)
[API] listening on http://0.0.0.0:8000 (docs: /docs)
Open http://localhost:8000/docs for the interactive API explorer.
Define in-process data generators using Python lambdas. These run without any external dependencies — useful for development, testing, and demos.
RateSource(
name="orders",
tenant_id="default",
fields={
"order_id": lambda: str(uuid.uuid4()),
"amount": lambda: random.randint(100, 5000),
"status": lambda: random.choice(["created", "paid", "shipped"]),
"timestamp": lambda: int(time.time()),
},
optional_fields={
"currency": lambda: random.choice(["USD", "EUR", "INR"]),
},
type_breaks={"amount": "not_a_number"}, # injected after schema is registered
)Set enabled: true and fill in connection details to activate a source without touching Python.
streams:
- name: kafka_orders
tenant_id: acme
enabled: true
source:
type: kafka
brokers: "localhost:9092"
topic: orders
starting_offsets: latest
value_format: json
- name: s3_events
tenant_id: acme
enabled: false
source:
type: file
path: /data/events/
format: json
trigger_interval: "10 seconds"
- name: payments_api
tenant_id: acme
enabled: false
source:
type: http
url: "https://api.example.com/payments"
method: GET
poll_interval_seconds: 30
response_path: "data.events"Enable a connector package (downloaded once via Maven) and define streams using it.
packages:
kinesis:
package: "com.qubole.spark:spark-sql-kinesis_2.12:1.2.0_spark-3.0"
enabled: false
eventhubs:
package: "com.microsoft.azure:azure-eventhubs-spark_2.12:2.3.22"
enabled: false
streams:
- name: kinesis_clickstream
tenant_id: analytics
enabled: false
source:
type: spark_connector
format: kinesis
extractor: json_data
options:
streamName: "clickstream"
region: "us-east-1"
startingPosition: "TRIM_HORIZON"Built-in extractors map connector-specific column names to event dicts:
| Extractor | Column | Works with |
|---|---|---|
json_value |
value |
Kafka, MSK, Kinesis-Qubole |
json_data |
data |
AWS Kinesis native |
json_body |
body |
Azure Event Hubs |
json_payload |
payload |
Google Pub/Sub |
raw_rows |
all columns | Delta Lake, Iceberg, JDBC |
Define per-dataset pipelines using PySpark column expressions. Datasets with no entry receive a pass-through.
from pyspark.sql import functions as F
from engine.transformer import Pipeline
TRANSFORMS = {
"orders": (
Pipeline("orders")
.add_field("platform", "stream_engine")
.filter(F.col("status") != "cancelled")
),
"iot_sensors": (
Pipeline("iot_sensors")
.cast("temperature", DoubleType())
.add_computed("temp_fahrenheit", F.round(F.col("temperature") * 9 / 5 + 32, 2))
.filter(F.col("temperature") < 90.0)
),
}Note:
config/transforms.pyis imported afterSparkSession.getOrCreate()inmain.py. PySpark 4.x requires an activeSparkContexteven to constructF.col()expressions at module level.
Schemas are automatically inferred from the first event of each dataset and persisted on disk:
schema_registry/datasets/
└── {tenant_id}/
└── {dataset}/
├── v1.json ← initial StructType schema
├── v2.json ← evolved schema (new fields added)
└── metadata.json ← current_version, version history, timestamps
Schema evolution is additive and automatic:
- A new field in an incoming event →
bump_version()is called, a new versioned JSON is written - A type mismatch (e.g.
amountwasbigint, now arrives asstring) → the event is rejected to the dead-letter list; the schema is not modified
The API runs on http://localhost:8000 in a background thread. Interactive docs at /docs.
| Method | Path | Description |
|---|---|---|
POST |
/streams/register |
Register + immediately start a new source |
GET |
/streams |
List all registered streams and their status |
GET |
/streams/{name} |
Get one stream's status |
POST |
/streams/{name}/start |
Restart a stopped stream |
POST |
/streams/{name}/stop |
Stop a stream without unregistering |
DELETE |
/streams/{name} |
Stop + unregister a stream |
Register a Kafka stream:
curl -X POST http://localhost:8000/streams/register \
-H "Content-Type: application/json" \
-d '{
"name": "live_orders",
"tenant_id": "acme",
"source": {
"type": "kafka",
"brokers": "localhost:9092",
"topic": "orders"
}
}'Register a third-party connector stream:
curl -X POST http://localhost:8000/streams/register \
-H "Content-Type: application/json" \
-d '{
"name": "kinesis_events",
"tenant_id": "acme",
"source": {
"type": "spark_connector",
"format": "kinesis",
"extractor": "json_data",
"options": {
"streamName": "events",
"region": "us-east-1"
}
}
}'| Method | Path | Description |
|---|---|---|
GET |
/metrics/snapshot |
Global counters, rejection rate, per-dataset breakdown |
GET |
/metrics/streams |
Running/stopped status of all streams |
GET |
/metrics/{dataset} |
Per-dataset counters: processed, rejected, evolved, new |
Example snapshot response:
{
"uptime_seconds": 120.4,
"global_stats": {"processed": 940, "rejected": 12, "evolved": 2, "new": 4},
"rejection_rate": 0.0126,
"per_dataset": {
"orders": {"processed": 310, "rejected": 4, "evolved": 1, "new": 1},
"clickstream": {"processed": 280, "rejected": 3, "evolved": 0, "new": 1},
"iot_sensors": {"processed": 195, "rejected": 2, "evolved": 1, "new": 1},
"payments": {"processed": 155, "rejected": 3, "evolved": 0, "new": 1}
}
}The ObservabilityPipeline wraps three components:
PipelineMetrics — in-memory counters per dataset (processed, rejected, evolved, new). Printed every 5 batches and exposed via /metrics/snapshot.
AnomalyDetector — stateful rolling-window rules:
| Alert | Level | Trigger |
|---|---|---|
REJECTION_SPIKE |
CRITICAL | ≥30% of events rejected in the last 60s |
SCHEMA_THRASH |
WARNING | Schema evolved ≥3 times in 120s |
DATA_GAP |
WARNING | No events seen for a dataset in 30s |
AlertManager — fanout dispatcher for Alert objects. Default handler prints to stdout. Add custom handlers for webhooks, files, or Slack:
@obs.alerts.register
def send_to_slack(alert):
requests.post(SLACK_WEBHOOK, json={"text": str(alert)})Register a handler on the Consumer in main.py:
@consumer.register
def write_to_delta(df, rejected, dataset, version, status):
df.write.format("delta").mode("append").save(f"/delta/{dataset}")- Subclass
Sourceinengine/sources/ - Implement
name,tenant_id,trigger_interval,read(),extract_rows() - Add a factory branch in
config/loader.py(_build_source) andapi/routes/streams.py(_build_source)
Pass a callable instead of a string name:
SparkConnectorSource(
name="my_source",
format="myformat",
extractor=lambda df: [{"id": r["id"], "val": r["payload"]} for r in df.collect()],
options={...},
)Each source carries a tenant_id. StreamManager creates one SchemaRegistry per tenant, pointing to schema_registry/datasets/{tenant_id}/. Tenant A's orders schema is stored and versioned completely independently from tenant B's orders schema. All tenants share the same Spark session and engine pipeline.
Spark streaming checkpoints are written to:
/tmp/stream_engine_checkpoints/{tenant_id}/{stream_name}/
These enable Spark to resume a query from where it left off after a restart. Change checkpoint_base in StreamManager.__init__() to use a persistent location (HDFS, S3, etc.) for production deployments.