Skip to content

About

Multi-tenant PySpark streaming engine with automatic schema evolution, pluggable sources (Kafka, files, HTTP, Spark connectors), live-editable ETL, and built-in observability.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

Stream Engine

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.


Features

  • 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

Architecture

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

Requirements

  • Python 3.11+
  • Java 21 (OpenJDK 21 recommended — PySpark 4.x requires jdk.incubator.vector)
pip install -r requirements.txt

requirements.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
The JAVA_HOME is set automatically in main.py to /opt/homebrew/opt/openjdk@21. Adjust if your JDK is elsewhere.


Quick Start

python3 main.py

On 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.


Configuration

Synthetic streams — config/streams.py

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
)

Real sources — config/streams.yaml

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"

Third-party connectors — config/connectors.yaml

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

ETL transforms — config/transforms.py

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.py is imported after SparkSession.getOrCreate() in main.py. PySpark 4.x requires an active SparkContext even to construct F.col() expressions at module level.


Schema Registry

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. amount was bigint, now arrives as string) → the event is rejected to the dead-letter list; the schema is not modified

REST API

The API runs on http://localhost:8000 in a background thread. Interactive docs at /docs.

Stream endpoints

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"
      }
    }
  }'

Metrics endpoints

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}
  }
}

Observability

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)})

Extending the Platform

Add a new downstream sink

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}")

Add a new source type

  1. Subclass Source in engine/sources/
  2. Implement name, tenant_id, trigger_interval, read(), extract_rows()
  3. Add a factory branch in config/loader.py (_build_source) and api/routes/streams.py (_build_source)

Add a custom extractor for a connector

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={...},
)

How Multi-Tenancy Works

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.


Checkpointing

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.

About

Multi-tenant PySpark streaming engine with automatic schema evolution, pluggable sources (Kafka, files, HTTP, Spark connectors), live-editable ETL, and built-in observability.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages