Lightweight data observability for Microsoft Fabric and Azure Databricks — anomaly detection + automatic lineage tracking, zero dependencies.
Project description
pipeline_sentinel
Lightweight data observability for Microsoft Fabric and Azure Databricks — anomaly detection + automatic lineage tracking, zero dependencies, drop it into any notebook.
What it does
Anomaly detection — after every pipeline run, sentinel profiles your output DataFrame and compares it against historical baselines using Z-score statistics. Bad data gets caught before it reaches downstream consumers.
Lineage tracking — patches PySpark's read/write operations automatically so you get a full dependency graph with zero code changes. When an anomaly fires, sentinel tells you exactly which downstream tables are at risk.
Installation
Option A — Fabric / Databricks workspace (recommended)
Upload the folder to your workspace Files, then import from any notebook:
# In a Fabric or Databricks notebook cell:
import sys
sys.path.insert(0, "/lakehouse/default/Files/pipeline_sentinel") # Fabric
# sys.path.insert(0, "/Workspace/Shared/pipeline_sentinel") # Databricks
from pipeline_sentinel.sentinel import PipelineSentinel
from pipeline_sentinel.lineage import LineageTracker
Option B — pip (once published to PyPI)
pip install pipeline-sentinel
Option C — local / CI
git clone https://github.com/tanushreearora01/pipeline-sentinel
cd pipeline-sentinel
pip install -e . # no dependencies beyond stdlib + optional pandas/pyspark
Quick Start
Anomaly detection only
from pipeline_sentinel.sentinel import PipelineSentinel
sentinel = PipelineSentinel(
pipeline_name="sales_etl",
table_name="gold.fact_sales",
alert_min_severity="MEDIUM",
)
# After your transformation:
anomalies = sentinel.record(df=output_df, processing_seconds=42.3)
That's it. No config files. No external services required.
With lineage tracking
from pipeline_sentinel.sentinel import PipelineSentinel
from pipeline_sentinel.lineage import LineageTracker
tracker = LineageTracker(pipeline_name="sales_etl")
sentinel = PipelineSentinel(
pipeline_name="sales_etl",
table_name="silver.orders",
alert_min_severity="MEDIUM",
lineage_tracker=tracker, # attach tracker here
)
with tracker.track():
raw = spark.read.table("bronze.raw_orders") # auto-recorded
clean = transform(raw)
clean.write.saveAsTable("silver.orders") # auto-recorded
anomalies = sentinel.record(df=clean, processing_seconds=35.1)
# If an anomaly fires, the alert automatically includes:
# ⚠ Downstream at risk: gold.fact_sales, gold.kpi_daily
Anomaly Detection
Checks
| Check | Method | Description |
|---|---|---|
| Row count drop | Z-score | Flags unusual row count vs historical mean |
| Row count spike | Z-score | Flags unexpected volume increases |
| Processing time spike | Z-score | Catches slow jobs before SLA breach |
| Null rate spike | Z-score + absolute threshold | Per-column null rate monitoring |
| Duplicate spike | Z-score | Detects upstream deduplication failures — opt-in via check_duplicates=True (avoids a second full-table scan) |
| Schema drift | Column set diff | Added/removed columns since last run |
| Empty dataset | Absolute | Immediately flags 0-row outputs as CRITICAL |
Severity levels
| Severity | Z-score range | Console colour |
|---|---|---|
| LOW | < 2.0 | Blue |
| MEDIUM | 2.0 – 2.9 | Yellow |
| HIGH | 3.0 – 3.9 | Red |
| CRITICAL | 4.0+ / always for empty datasets | Magenta |
Alert sinks
| Sink | How to enable |
|---|---|
| Notebook console | On by default — colour-coded output |
| Microsoft Teams | Pass teams_webhook_url to PipelineSentinel |
| Delta table | Pass delta_anomaly_path to PipelineSentinel |
Full configuration
sentinel = PipelineSentinel(
pipeline_name="sales_etl",
table_name="gold.fact_sales",
# Persist run history to Delta — improves detection over time
history_table_path="abfss://container@storage.dfs.core.windows.net/sentinel/history/sales_etl",
# Save all anomalies to a queryable Delta table
delta_anomaly_path="abfss://container@storage.dfs.core.windows.net/sentinel/anomalies",
# Teams channel alerts
teams_webhook_url="https://your-org.webhook.office.com/...",
# Severity filter: LOW | MEDIUM | HIGH | CRITICAL
alert_min_severity="MEDIUM",
# How many std deviations to flag (default 2.5)
zscore_threshold=2.5,
# Always flag columns with null rate above this (default 10%)
null_rate_threshold=0.10,
# How many historical runs to compare against (default 30)
max_history=30,
# Opt-in duplicate counting (requires a second full-table scan, default False)
check_duplicates=False,
# Optional lineage tracker (see Lineage section below)
lineage_tracker=tracker,
)
Usage patterns
Simple record
anomalies = sentinel.record(
df=output_df,
processing_seconds=elapsed,
run_id="run_20240501_001", # optional — auto-generated if omitted
)
Context manager — auto-times processing
with sentinel.watch(df=output_df, run_id="run_001"):
pass
print(sentinel.last_anomalies)
Pre-loading history (testing / backfill)
from pipeline_sentinel.models import PipelineRun
from datetime import datetime
history = [
PipelineRun(
pipeline_name="sales_etl",
table_name="gold.fact_sales",
run_id="run_000",
run_timestamp=datetime(2024, 4, 30, 10, 0),
row_count=100_000,
processing_time_seconds=35.0,
null_rates={"customer_id": 0.02},
duplicate_count=10,
column_names=["order_id", "customer_id", "amount"],
),
# ... more runs
]
sentinel = PipelineSentinel(
pipeline_name="sales_etl",
table_name="gold.fact_sales",
history_runs=history,
)
Lineage Tracking
LineageTracker intercepts PySpark's DataFrameReader and DataFrameWriter inside a with tracker.track(): block. No changes to your pipeline code are needed — reads and writes are captured automatically.
How it works
with tracker.track():
df = spark.read.table("bronze.raw_orders") ← recorded as READ
result = transform(df)
result.write.saveAsTable("silver.orders") ← recorded as WRITE
# Edges added: bronze.raw_orders → silver.orders
Intercepted operations:
spark.read.table(name)spark.read.format(...).load(path)df.write.saveAsTable(name)df.write.format(...).save(path)df.write.insertInto(name)
spark.sql(...) is not auto-intercepted — use tracker.record_read() / tracker.record_write() manually for SQL cells.
Visualise the graph
tracker.show()
Lineage Graph (5 table(s), 4 edge(s))
───────────────────────────────────────────────────────
bronze.raw_customers
└──► silver.orders_enriched
├──► gold.fact_sales
└──► gold.kpi_daily
bronze.raw_orders
└──► silver.orders_enriched (already shown)
Impact and ancestry queries
# What breaks if silver.orders_enriched has bad data?
tracker.impact("silver.orders_enriched")
# → ['gold.fact_sales', 'gold.kpi_daily']
# What feeds into gold.fact_sales?
tracker.ancestors("gold.fact_sales")
# → ['silver.orders_enriched', 'bronze.raw_orders', 'bronze.raw_customers']
Workspace-wide lineage (shared graph)
Pass a single LineageGraph to all your pipeline trackers to build a cross-pipeline dependency map:
from pipeline_sentinel.lineage import LineageGraph, LineageTracker
shared_graph = LineageGraph()
tracker_ingest = LineageTracker("ingest_pipeline", graph=shared_graph)
tracker_reports = LineageTracker("reports_pipeline", graph=shared_graph)
# ... run both trackers over time ...
print(shared_graph.render()) # full workspace lineage tree
Manual API (Pandas / SQL cells)
with tracker.track():
tracker.record_read("bronze.raw_orders") # for spark.sql reads
tracker.record_write("silver.orders") # for spark.sql writes
Querying the Anomaly Delta Table
-- Most recent critical / high anomalies
SELECT pipeline_name, table_name, anomaly_type, severity, message, detected_at
FROM delta.`abfss://container@storage.dfs.core.windows.net/sentinel/anomalies`
WHERE severity IN ('CRITICAL', 'HIGH')
ORDER BY detected_at DESC
LIMIT 100;
-- Row count trend for a specific pipeline
SELECT run_id, run_timestamp, row_count
FROM delta.`abfss://container@storage.dfs.core.windows.net/sentinel/history/sales_etl`
ORDER BY run_timestamp;
Architecture
pipeline_sentinel/
├── sentinel.py # PipelineSentinel — main orchestrator
├── lineage.py # LineageTracker + LineageGraph — auto-lineage
├── detectors.py # AnomalyDetector — Z-score logic
├── alerts.py # AlertManager — notebook / Teams / Delta sinks
└── models.py # PipelineRun, Anomaly dataclasses
usage_notebook.py # Full usage example with synthetic data demo
Zero Hard Dependencies
pipeline_sentinel uses only the Python standard library by default.
- PySpark support is auto-detected when running in Fabric / Databricks
- Pandas fallback works for local testing and CI
- Teams alerts use
urllib— norequestsneeded - Lineage patching restores all original methods on context exit, even on exception
License
MIT
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file pipeline_sentinel-0.1.1.tar.gz.
File metadata
- Download URL: pipeline_sentinel-0.1.1.tar.gz
- Upload date:
- Size: 23.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
ad876b992589f80a8bdbbf16ddd42183dc1e13b13545bcdb0aac5eb9a5555ccb
|
|
| MD5 |
2be4217dc4fac2f784816f31a2e9181a
|
|
| BLAKE2b-256 |
bdf423e98001fcf975774110ebc810fae8c44d1952b930fa308bd38918a83058
|
File details
Details for the file pipeline_sentinel-0.1.1-py3-none-any.whl.
File metadata
- Download URL: pipeline_sentinel-0.1.1-py3-none-any.whl
- Upload date:
- Size: 20.3 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.0
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
f93fa350e38be28e086c79cafa5d873272e1aa0c38ff25473d8f489097fd5947
|
|
| MD5 |
b7aa7520f8b7d987509f747475452ce0
|
|
| BLAKE2b-256 |
20d0a70e2ce7da344c210c9f50974ecf817f458fd93dcb01039f9be22d7e4be2
|