Skip to main content

Dead simple pipeline monitoring. Know when your pipelines die.

Project description

Deadpipe Python SDK

Dead simple pipeline monitoring. Know when your pipelines die.

Installation

pip install deadpipe

Quick Start

Option 1: Decorator (Recommended)

from deadpipe import Deadpipe

dp = Deadpipe("your-api-key")

@dp.heartbeat("daily-sales-etl")
def run_pipeline():
    # Your pipeline code here
    process_data()
    return {"records_processed": 1500}  # Optional: track records

# That's it! Deadpipe will ping on success or failure.
run_pipeline()

Option 2: Context Manager

from deadpipe import Deadpipe

dp = Deadpipe("your-api-key")

with dp.pipeline("hourly-sync"):
    # Your code here
    sync_data()

Option 3: Manual Ping

from deadpipe import Deadpipe

dp = Deadpipe("your-api-key")

try:
    run_my_job()
    dp.ping("my-job", status="success", records_processed=1000)
except Exception as e:
    dp.ping("my-job", status="failed")
    raise

Option 4: Environment Variable

Set DEADPIPE_API_KEY and use module-level functions:

import deadpipe

@deadpipe.heartbeat("my-pipeline")
def my_job():
    pass

Async Support (FastAPI, aiohttp, asyncio)

For async applications, install with async extras:

pip install deadpipe[async]

Async Decorator

from deadpipe import AsyncDeadpipe

dp = AsyncDeadpipe("your-api-key")

@dp.heartbeat("async-pipeline")
async def my_async_job():
    await fetch_data()
    return {"records_processed": 500}

Async Context Manager

async with dp.pipeline("async-etl"):
    await process_async_data()

Async Manual Ping

await dp.ping("my-job", status="success")

FastAPI Example

from fastapi import FastAPI
from deadpipe import AsyncDeadpipe

app = FastAPI()
dp = AsyncDeadpipe()  # Uses DEADPIPE_API_KEY env var

@app.on_event("startup")
async def startup():
    # Optionally track app startup
    await dp.ping("fastapi-app", status="success")

@app.on_event("shutdown") 
async def shutdown():
    await dp.close()  # Clean up HTTP session

@app.post("/process")
@dp.heartbeat("data-processor")
async def process_data():
    result = await heavy_processing()
    return {"records_processed": len(result)}

Session Management

AsyncDeadpipe uses connection pooling via aiohttp.ClientSession. For long-running apps, use as a context manager or call close():

# Option 1: Context manager (auto-closes)
async with AsyncDeadpipe("your-key") as dp:
    await dp.ping("my-pipeline")

# Option 2: Manual close
dp = AsyncDeadpipe("your-key")
try:
    await dp.ping("my-pipeline")
finally:
    await dp.close()

Airflow Integration

from deadpipe import Deadpipe

dp = Deadpipe(api_key=Variable.get("DEADPIPE_API_KEY"))

@dp.heartbeat("{{ dag.dag_id }}")
def my_task():
    ...

Or add to the end of any task:

from deadpipe import ping

def my_task():
    # ... your code ...
    ping("daily-etl", status="success")

dbt Integration

Add to your dbt_project.yml on-run-end hook:

on-run-end:
  - "{{ deadpipe_heartbeat('dbt-run') }}"

Or call from Python:

# In your dbt runner script
from deadpipe import Deadpipe

dp = Deadpipe("your-api-key")

with dp.pipeline("dbt-daily"):
    subprocess.run(["dbt", "run"], check=True)

API Reference

Deadpipe(api_key, base_url, timeout)

Create a client instance.

  • api_key: Your API key (or set DEADPIPE_API_KEY env var)
  • base_url: Override for self-hosted (default: https://www.deadpipe.com/api/v1)
  • timeout: Request timeout in seconds (default: 10)

dp.ping(pipeline_id, status, duration_ms, records_processed, app_name)

Send a heartbeat.

  • pipeline_id: Unique identifier for this pipeline
  • status: "success" or "failed"
  • duration_ms: How long the run took (optional)
  • records_processed: Number of records (optional)
  • app_name: Group pipelines under an app (optional)

@dp.heartbeat(pipeline_id, app_name, on_error)

Decorator that auto-sends heartbeats.

  • on_error: What to do on exception:
    • "ping" (default): Send failed heartbeat, then re-raise
    • "raise": Re-raise without heartbeat
    • "ignore": Send success heartbeat anyway

with dp.pipeline(pipeline_id, app_name)

Context manager for heartbeats.

AsyncDeadpipe(api_key, base_url, timeout)

Async client with the same API as Deadpipe, but all methods are async:

  • await dp.ping(...) - Async heartbeat
  • @dp.heartbeat(...) - Decorator for async functions
  • async with dp.pipeline(...) - Async context manager
  • await dp.run(pipeline_id, fn, *args) - Run async function with heartbeat
  • await dp.close() - Close HTTP session
  • async with AsyncDeadpipe(...) as dp: - Auto-close on exit

Dependencies

  • Sync: Zero dependencies (uses Python standard library only)
  • Async: Requires aiohttp (install with pip install deadpipe[async])

License

MIT

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

deadpipe-0.1.1.tar.gz (5.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

deadpipe-0.1.1-py3-none-any.whl (5.9 kB view details)

Uploaded Python 3

File details

Details for the file deadpipe-0.1.1.tar.gz.

File metadata

  • Download URL: deadpipe-0.1.1.tar.gz
  • Upload date:
  • Size: 5.4 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.7

File hashes

Hashes for deadpipe-0.1.1.tar.gz
Algorithm Hash digest
SHA256 cf712d6edbf4decc31f6a3cf21229fa79fe0d4d0e2b1b5d351be1d2859b70805
MD5 fe140c0ad44bf670479436a9f023f3c8
BLAKE2b-256 421b3405d6dad1eb306807e7ba9310bb4c258a2f0344d39de3cb4d9b6b97b96f

See more details on using hashes here.

File details

Details for the file deadpipe-0.1.1-py3-none-any.whl.

File metadata

  • Download URL: deadpipe-0.1.1-py3-none-any.whl
  • Upload date:
  • Size: 5.9 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.13.7

File hashes

Hashes for deadpipe-0.1.1-py3-none-any.whl
Algorithm Hash digest
SHA256 14f1aa90959790826bd9599fd6bb9d66a71c80b92623ab45a8a3f673ce8ad208
MD5 a6fe2b150a151d657bdeee4663367257
BLAKE2b-256 41c6fefef29d8ef010306bc8e0d8a5e6a650fd61d56488f919984ce6952ac424

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page