How to Build a Notification System That Actually Gets Read

· 12 min read · Automation

Build a multi-channel notification system that routes alerts to Slack, email, and SMS based on severity — with deduplication, rate limiting, digests, and escalation so people do not ignore your alerts.

How to Build a Notification System That Actually Gets Read

Your pipeline sends alerts. Nobody reads them. The Slack channel has 200 unread messages, all saying “Pipeline completed successfully.” When something actually breaks, the real alert is buried.

This is alert fatigue. And it is the reason most notification systems fail — not because they cannot send messages, but because they send too many of the wrong ones.

This guide builds a notification system that routes alerts by severity, rate-limits noise, digests low-priority updates, and escalates critical failures. By the end, your alerts will be read because they will be worth reading.

Who This Is For

  • Engineers whose Slack channels are flooded with meaningless pipeline alerts
  • On-call teams who need to know when something genuinely requires human attention
  • Vibe coders who built automations that run unattended and need smart alerting when things go wrong
  • Managers who want to know about problems before their customers or stakeholders do

You need basic Python knowledge. The patterns here work with any messaging system (Slack, email, Teams, PagerDuty) — the guide focuses on the routing and filtering logic.

The Notification Architecture

flowchart TD
  E["Event\n(pipeline output)"] --> DD["Deduplicator\n(same alert, 15 min window)"]
  DD -->|repeat| SUP["Suppress\n(count folded into next)"]
  DD -->|new| R["Rules Engine\n(severity + routing)"]
  R -->|Critical| S["Slack #critical\n+ SMS"]
  R -->|Warning| SL["Slack #alerts"]
  R -->|Info| D["Digest\n(daily summary)"]
  R -->|Success| L["Log Only"]
  S --> RL["Rate Limiter"]
  SL --> RL
  RL --> Send["Deliver"]

Events come in. The deduplicator folds repeats of the same failure into one. The rules engine decides where new ones go. Rate limiting prevents flooding. Digests batch low-priority items into a single daily message.

What You Will Need

Bash
pip install requests schedule twilio
  • requests — HTTP client for Slack and webhook calls
  • schedule — lightweight task scheduling for digests
  • twilio — SMS delivery for critical alerts

Step 1: Define Alert Severity

Every alert should have a severity level. Without this, all alerts look the same and all get ignored.

Python
from enum import Enum
from dataclasses import dataclass, field
from datetime import datetime


class Severity(Enum):
    CRITICAL = "critical"   # Pipeline failed, data loss risk
    WARNING = "warning"     # Something unusual, might need attention
    INFO = "info"           # Status update, progress tracking
    SUCCESS = "success"     # Pipeline completed normally


@dataclass
class Alert:
    """A structured alert with all routing metadata."""

    title: str
    message: str
    severity: Severity
    source: str                          # Which pipeline or system
    timestamp: datetime = field(default_factory=datetime.now)
    context: dict = field(default_factory=dict)  # Extra data for debugging

    @property
    def emoji(self):
        return {
            Severity.CRITICAL: "[CRIT]",
            Severity.WARNING: "[WARN]",
            Severity.INFO: "[INFO]",
            Severity.SUCCESS: "[OK]",
        }[self.severity]

Creating Alerts

Python
# From a pipeline
alert = Alert(
    title="Daily sales pipeline failed",
    message="Database connection refused after 3 retries",
    severity=Severity.CRITICAL,
    source="sales_pipeline",
    context={"stage": "extract", "retries": 3, "error": "ConnectionRefused"},
)

# Informational
alert = Alert(
    title="Pipeline completed",
    message="Processed 4,521 orders in 12.3 seconds",
    severity=Severity.SUCCESS,
    source="sales_pipeline",
    context={"rows": 4521, "duration_seconds": 12.3},
)

Step 2: Notification Channels

Slack Integration

Python
import requests
import os
import logging

logger = logging.getLogger("notifications")

class SlackNotifier:
    """Send notifications to Slack channels via webhooks."""

    def __init__(self, webhooks=None):
        self.webhooks = webhooks or {
            "critical": os.environ.get("SLACK_WEBHOOK_CRITICAL"),
            "alerts": os.environ.get("SLACK_WEBHOOK_ALERTS"),
            "info": os.environ.get("SLACK_WEBHOOK_INFO"),
        }

    def send(self, alert, channel="alerts"):
        """Send an alert to a Slack channel."""
        webhook_url = self.webhooks.get(channel)
        if not webhook_url:
            logger.warning(f"No Slack webhook configured for channel: {channel}")
            return False

        payload = {
            "blocks": [
                {
                    "type": "header",
                    "text": {
                        "type": "plain_text",
                        "text": f"{alert.emoji} {alert.title}",
                    },
                },
                {
                    "type": "section",
                    "text": {
                        "type": "mrkdwn",
                        "text": (
                            f"*Source:* {alert.source}\n"
                            f"*Time:* {alert.timestamp.strftime('%Y-%m-%d %H:%M:%S')}\n"
                            f"*Message:* {alert.message}"
                        ),
                    },
                },
            ],
        }

        # Add context fields if present
        if alert.context:
            fields = []
            for key, value in list(alert.context.items())[:8]:
                fields.append({
                    "type": "mrkdwn",
                    "text": f"*{key}:* {value}",
                })
            payload["blocks"].append({
                "type": "section",
                "fields": fields,
            })

        response = requests.post(webhook_url, json=payload, timeout=10)
        if response.status_code == 200:
            logger.info(f"Slack alert sent to #{channel}: {alert.title}")
            return True
        else:
            logger.error(f"Slack delivery failed: {response.status_code}")
            return False

Email Notifications

Python
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart

class EmailNotifier:
    """Send notification emails."""

    def __init__(self, smtp_host=None, smtp_port=587, username=None, password=None):
        self.smtp_host = smtp_host or os.environ.get("SMTP_HOST")
        self.smtp_port = smtp_port
        self.username = username or os.environ.get("SMTP_USERNAME")
        self.password = password or os.environ.get("SMTP_PASSWORD")

    def send(self, alert, recipients):
        """Send an alert via email."""
        msg = MIMEMultipart("alternative")
        msg["Subject"] = f"[{alert.severity.value.upper()}] {alert.title}"
        msg["From"] = self.username
        msg["To"] = ", ".join(recipients)

        # Plain text
        text = (
            f"Alert: {alert.title}\n"
            f"Severity: {alert.severity.value}\n"
            f"Source: {alert.source}\n"
            f"Time: {alert.timestamp.strftime('%Y-%m-%d %H:%M:%S')}\n\n"
            f"{alert.message}\n"
        )

        # HTML
        context_rows = "".join(
            f"<tr><td><strong>{k}</strong></td><td>{v}</td></tr>"
            for k, v in alert.context.items()
        )
        html = f"""
        <h2>{alert.emoji} {alert.title}</h2>
        <table>
            <tr><td><strong>Severity</strong></td><td>{alert.severity.value}</td></tr>
            <tr><td><strong>Source</strong></td><td>{alert.source}</td></tr>
            <tr><td><strong>Time</strong></td><td>{alert.timestamp.strftime('%Y-%m-%d %H:%M:%S')}</td></tr>
            {context_rows}
        </table>
        <p>{alert.message}</p>
        """

        msg.attach(MIMEText(text, "plain"))
        msg.attach(MIMEText(html, "html"))

        with smtplib.SMTP(self.smtp_host, self.smtp_port) as server:
            server.starttls()
            server.login(self.username, self.password)
            server.send_message(msg)

        logger.info(f"Email sent to {len(recipients)} recipients: {alert.title}")

SMS for Critical Alerts

Slack and email are fine for warnings. A critical failure at 2 AM needs something that actually wakes someone up — that means paying for SMS, and reserving it for the one severity level that justifies the cost.

Python
from twilio.rest import Client

class SMSNotifier:
    """Send SMS for alerts urgent enough to justify paying for delivery."""

    def __init__(self, account_sid=None, auth_token=None, from_number=None):
        self.client = Client(
            account_sid or os.environ.get("TWILIO_ACCOUNT_SID"),
            auth_token or os.environ.get("TWILIO_AUTH_TOKEN"),
        )
        self.from_number = from_number or os.environ.get("TWILIO_FROM_NUMBER")

    def send(self, alert, to):
        """Send an SMS. `to` is a list of E.164 phone numbers."""
        body = f"{alert.emoji} {alert.title} ({alert.source})"[:160]
        sent = 0
        for number in to:
            try:
                self.client.messages.create(body=body, from_=self.from_number, to=number)
                sent += 1
            except Exception as e:
                logger.error(f"SMS delivery failed to {number}: {e}")
        logger.info(f"SMS sent to {sent}/{len(to)} recipients: {alert.title}")
        return sent > 0

Step 3: Routing Rules

Define where each alert goes based on severity and source:

flowchart LR
  A["Alert"] --> RE["Rules\nEngine"]
  RE -->|"CRITICAL\n(any source)"| C["Slack #critical\n+ Email\n+ SMS"]
  RE -->|"WARNING\n(production)"| W["Slack #alerts\n+ Email"]
  RE -->|"WARNING\n(dev)"| WD["Slack #dev"]
  RE -->|"INFO"| I["Daily Digest"]
  RE -->|"SUCCESS"| S["Log Only"]
Python
class NotificationRouter:
    """Route alerts to the right channels based on rules."""

    def __init__(self):
        self.slack = SlackNotifier()
        self.email = EmailNotifier()
        self.sms = SMSNotifier()
        self.digest = DigestCollector()

        # Define routing rules
        self.rules = [
            {
                "match": lambda a: a.severity == Severity.CRITICAL,
                "actions": [
                    lambda a: self.slack.send(a, channel="critical"),
                    lambda a: self.email.send(a, recipients=["[email protected]"]),
                    lambda a: self.sms.send(a, to=self._oncall_numbers()),
                ],
            },
            {
                "match": lambda a: a.severity == Severity.WARNING,
                "actions": [
                    lambda a: self.slack.send(a, channel="alerts"),
                ],
            },
            {
                "match": lambda a: a.severity == Severity.INFO,
                "actions": [
                    lambda a: self.digest.add(a),
                ],
            },
            {
                "match": lambda a: a.severity == Severity.SUCCESS,
                "actions": [
                    lambda a: logger.info(f"Success: {a.source}{a.message}"),
                ],
            },
        ]

    def _oncall_numbers(self):
        """E.164 numbers from ONCALL_PHONE, comma-separated. Empty if unset."""
        raw = os.environ.get("ONCALL_PHONE", "")
        return [n.strip() for n in raw.split(",") if n.strip()]

    def route(self, alert):
        """Route an alert through matching rules."""
        matched = False
        for rule in self.rules:
            if rule["match"](alert):
                for action in rule["actions"]:
                    try:
                        action(alert)
                    except Exception as e:
                        logger.error(f"Notification action failed: {e}")
                matched = True
                break

        if not matched:
            logger.warning(f"No routing rule matched alert: {alert.title}")

Step 4: Rate Limiting

Prevent flooding channels when multiple pipelines fail at once.

Python
from collections import defaultdict
from datetime import datetime, timedelta

class RateLimiter:
    """Prevent notification flooding."""

    def __init__(self, max_per_minute=5, max_per_hour=30):
        self.max_per_minute = max_per_minute
        self.max_per_hour = max_per_hour
        self.history = defaultdict(list)

    def should_send(self, channel):
        """Check if sending to this channel is allowed."""
        now = datetime.now()

        # Clean old entries
        self.history[channel] = [
            t for t in self.history[channel]
            if t > now - timedelta(hours=1)
        ]

        # Check rate limits
        recent_minute = sum(
            1 for t in self.history[channel] if t > now - timedelta(minutes=1)
        )
        recent_hour = len(self.history[channel])

        if recent_minute >= self.max_per_minute:
            logger.warning(f"Rate limited ({channel}): {recent_minute}/min limit reached")
            return False

        if recent_hour >= self.max_per_hour:
            logger.warning(f"Rate limited ({channel}): {recent_hour}/hour limit reached")
            return False

        return True

    def record_send(self, channel):
        """Record that a notification was sent."""
        self.history[channel].append(datetime.now())


class RateLimitedRouter(NotificationRouter):
    """Router with rate limiting applied."""

    def __init__(self):
        super().__init__()
        self.limiter = RateLimiter(max_per_minute=5, max_per_hour=30)

    def route(self, alert):
        """Route with rate limiting — critical alerts bypass limits."""
        if alert.severity == Severity.CRITICAL:
            # Critical alerts always go through
            super().route(alert)
            return

        channel = alert.severity.value
        if self.limiter.should_send(channel):
            super().route(alert)
            self.limiter.record_send(channel)
        else:
            # Silently add to digest instead
            self.digest.add(alert)
            logger.info(f"Rate limited — alert added to digest: {alert.title}")

Step 5: Deduplication

Rate limiting caps volume. It does not notice that the last twelve messages were the same failure. A flaky check that fails, retries, and fails again every 90 seconds will happily stay under a 5-per-minute limit while still paging someone twelve times for one incident.

Deduplication catches what rate limiting cannot: it groups alerts by what they are, not just how many arrived, and folds repeats into a single message with a count.

Python
import hashlib

class Deduplicator:
    """Suppress repeats of the same alert within a time window, and report how
    many were suppressed the next time a genuinely new one fires."""

    def __init__(self, window_minutes=15):
        self.window = timedelta(minutes=window_minutes)
        self.seen = {}  # key -> {"first_seen": datetime, "count": int}

    def _key(self, alert):
        return hashlib.sha1(f"{alert.source}:{alert.title}".encode()).hexdigest()

    def check(self, alert):
        """Returns (should_send, suppressed_count). should_send is False while
        inside the window; suppressed_count is how many repeats were folded
        into the alert that does eventually go out."""
        key = self._key(alert)
        now = datetime.now()
        entry = self.seen.get(key)

        if entry is None or now - entry["first_seen"] > self.window:
            self.seen[key] = {"first_seen": now, "count": 0}
            return True, 0

        entry["count"] += 1
        return False, entry["count"]


class DeduplicatingRouter(RateLimitedRouter):
    """Router that folds repeats of the same alert into one, instead of
    resending an unchanged failure every time a flaky check retries."""

    def __init__(self):
        super().__init__()
        self.dedup = Deduplicator(window_minutes=15)

    def route(self, alert):
        should_send, suppressed = self.dedup.check(alert)

        if not should_send:
            logger.info(f"Deduplicated ({suppressed} suppressed so far): {alert.title}")
            return

        if suppressed:
            alert.message = f"{alert.message}\n\n(repeated {suppressed}x in the last {self.dedup.window.seconds // 60} minutes)"

        super().route(alert)

The key is source + title, not the full message — a pipeline that fails with a slightly different row count each retry should still be treated as the same incident, not twelve different ones. Widen the window past 15 minutes for noisy upstream dependencies; keep it short for anything customer-facing.

Step 6: Daily Digests

Batch low-priority updates into a single daily summary.

Python
import json

class DigestCollector:
    """Collect alerts for daily digest delivery."""

    def __init__(self, digest_file="digest_buffer.json"):
        self.digest_file = digest_file
        self._load()

    def _load(self):
        """Load existing digest buffer."""
        try:
            with open(self.digest_file, "r") as f:
                self.buffer = json.load(f)
        except (FileNotFoundError, json.JSONDecodeError):
            self.buffer = []

    def _save(self):
        """Persist digest buffer to disk."""
        with open(self.digest_file, "w") as f:
            json.dump(self.buffer, f, default=str)

    def add(self, alert):
        """Add an alert to the digest buffer."""
        self.buffer.append({
            "title": alert.title,
            "message": alert.message,
            "severity": alert.severity.value,
            "source": alert.source,
            "timestamp": alert.timestamp.isoformat(),
        })
        self._save()

    def build_digest(self):
        """Build a formatted digest from buffered alerts."""
        if not self.buffer:
            return None

        # Group by source
        by_source = defaultdict(list)
        for item in self.buffer:
            by_source[item["source"]].append(item)

        lines = [f"Daily Pipeline Digest -- {datetime.now().strftime('%Y-%m-%d')}\n"]

        for source, items in by_source.items():
            lines.append(f"\n*{source}* ({len(items)} events)")
            for item in items[-5:]:  # Last 5 per source
                severity_label = {"critical": "[CRIT]", "warning": "[WARN]", "info": "[INFO]", "success": "[OK]"}
                label = severity_label.get(item["severity"], "[--]")
                lines.append(f"  {label} {item['title']}")

            if len(items) > 5:
                lines.append(f"  _...and {len(items) - 5} more_")

        return "\n".join(lines)

    def send_and_clear(self, slack_notifier, channel="info"):
        """Send the digest and clear the buffer."""
        digest_text = self.build_digest()
        if not digest_text:
            logger.info("No items in digest — skipping")
            return

        webhook_url = slack_notifier.webhooks.get(channel)
        if webhook_url:
            requests.post(webhook_url, json={"text": digest_text}, timeout=10)

        item_count = len(self.buffer)
        self.buffer = []
        self._save()
        logger.info(f"Digest sent: {item_count} items")

Scheduling the Digest

Python
import schedule
import time

def run_digest():
    """Send the daily digest."""
    digest = DigestCollector()
    slack = SlackNotifier()
    digest.send_and_clear(slack, channel="info")

# Send digest every day at 8 AM
schedule.every().day.at("08:00").do(run_digest)

while True:
    schedule.run_pending()
    time.sleep(60)

Step 7: Escalation Chains

If a critical alert goes unacknowledged, escalate to the next person.

Python
class EscalationChain:
    """Escalate alerts through a chain of contacts."""

    def __init__(self):
        self.chains = {
            "default": [
                {"method": "slack", "channel": "critical", "wait_minutes": 0},
                {"method": "email", "contact": "[email protected]", "wait_minutes": 10},
                {"method": "email", "contact": "[email protected]", "wait_minutes": 30},
            ],
        }
        self.pending_escalations = {}

    def start_escalation(self, alert, chain_name="default"):
        """Begin escalation for a critical alert."""
        chain = self.chains.get(chain_name, self.chains["default"])
        escalation_id = f"{alert.source}_{alert.timestamp.isoformat()}"

        self.pending_escalations[escalation_id] = {
            "alert": alert,
            "chain": chain,
            "current_step": 0,
            "started_at": datetime.now(),
            "acknowledged": False,
        }

        # Send first notification immediately
        self._send_step(escalation_id)
        return escalation_id

    def acknowledge(self, escalation_id):
        """Mark an escalation as acknowledged — stop further steps."""
        if escalation_id in self.pending_escalations:
            self.pending_escalations[escalation_id]["acknowledged"] = True
            logger.info(f"Escalation acknowledged: {escalation_id}")

    def check_escalations(self):
        """Check if any pending escalations need the next step."""
        now = datetime.now()

        for esc_id, esc in self.pending_escalations.items():
            if esc["acknowledged"]:
                continue

            chain = esc["chain"]
            current = esc["current_step"]

            if current >= len(chain) - 1:
                continue

            next_step = chain[current + 1]
            elapsed = (now - esc["started_at"]).total_seconds() / 60

            if elapsed >= next_step["wait_minutes"]:
                esc["current_step"] += 1
                self._send_step(esc_id)

    def _send_step(self, escalation_id):
        """Send notification for the current escalation step."""
        esc = self.pending_escalations[escalation_id]
        step = esc["chain"][esc["current_step"]]
        alert = esc["alert"]

        logger.warning(
            f"Escalation step {esc['current_step'] + 1}/{len(esc['chain'])}: "
            f"{step['method']} for {alert.title}"
        )

Wiring It Together

Python
class NotificationSystem:
    """Complete notification system with routing, rate limiting, and digests."""

    def __init__(self):
        self.router = DeduplicatingRouter()
        self.escalation = EscalationChain()

    def notify(self, title, message, severity, source, **context):
        """Send a notification through the system."""
        alert = Alert(
            title=title,
            message=message,
            severity=severity,
            source=source,
            context=context,
        )

        self.router.route(alert)

        if severity == Severity.CRITICAL:
            self.escalation.start_escalation(alert)

        return alert


# Usage in a pipeline
notifications = NotificationSystem()

try:
    result = run_pipeline()
    notifications.notify(
        title="Sales pipeline completed",
        message=f"Processed {result['rows']} orders",
        severity=Severity.SUCCESS,
        source="sales_pipeline",
        rows=result["rows"],
        duration=result["duration"],
    )
except Exception as e:
    notifications.notify(
        title="Sales pipeline failed",
        message=str(e),
        severity=Severity.CRITICAL,
        source="sales_pipeline",
        error_type=type(e).__name__,
    )

Testing Your Alert Routing

Routing rules and deduplication windows are exactly the kind of code that breaks silently — a rule that stops matching does not throw an exception, it just stops sending. Test the logic directly instead of relying on watching a Slack channel for a week.

Python
from unittest.mock import MagicMock


def test_critical_alert_reaches_slack_email_and_sms():
    router = NotificationRouter()
    router.slack = MagicMock()
    router.email = MagicMock()
    router.sms = MagicMock()

    alert = Alert(title="DB down", message="", severity=Severity.CRITICAL, source="sales")
    router.route(alert)

    router.slack.send.assert_called_once()
    router.email.send.assert_called_once()
    router.sms.send.assert_called_once()


def test_deduplicator_suppresses_a_repeat_within_the_window():
    dedup = Deduplicator(window_minutes=15)
    alert = Alert(title="DB down", message="", severity=Severity.CRITICAL, source="sales")

    first_send, _ = dedup.check(alert)
    second_send, suppressed = dedup.check(alert)

    assert first_send is True
    assert second_send is False
    assert suppressed == 1


def test_deduplicator_allows_a_new_alert_once_the_window_expires():
    dedup = Deduplicator(window_minutes=15)
    alert = Alert(title="DB down", message="", severity=Severity.CRITICAL, source="sales")
    dedup.check(alert)

    key = dedup._key(alert)
    dedup.seen[key]["first_seen"] = datetime.now() - timedelta(minutes=20)

    should_send, _ = dedup.check(alert)
    assert should_send is True

Mock the notifiers, not the routing logic — the point is to prove the rules and the dedup window behave correctly, not to actually hit Slack or Twilio in CI. For the same pattern applied to the pipelines that feed these alerts, see Testing Data Pipelines: A Practical Guide with pytest.

What This Replaces

Old approachNew approach
print("Pipeline done")Structured alerts with severity
All alerts to one Slack channelRouting by severity and source
Same failure paging you 12 timesDeduplicated with a rollup count
200 unread messagesRate limiting + digests
Nobody notices critical failuresEscalation chains, SMS for critical
”I think it ran yesterday?”Success/failure notifications with context
Checking logs manuallyAlerts come to you

Common Alert Design Mistakes

MistakeWhy it failsFix
Alerting on every successAlert fatigue — nothing seems importantLog successes, alert on failures
No severity levelsCannot prioritiseUse CRITICAL / WARNING / INFO / SUCCESS
Same channel for everythingImportant alerts get buriedRoute by severity
No rate limitingOne broken pipeline floods channelCap per-minute and per-hour
No deduplicationA flaky retry pages you 12 times for one incidentFold repeats into a rollup count
Alert without context”Pipeline failed” — which one? How?Include source, stage, error, timestamp
No escalationNobody sees the 2 AM alertEscalation chains with increasing urgency

Next Steps

Start with Slack notifications on pipeline failures — that alone is a significant improvement over checking logs. Add severity levels next, then deduplication as soon as anything retries on failure — it is the fix that removes the most noise for the least code, and it is what stops rate limiting from being the only thing standing between a flaky check and a pager full of duplicates.

The digest pattern is particularly valuable: collect all the “it worked fine” messages into one daily summary so your alert channels stay clean for things that actually need attention.

For building the pipelines that feed into this notification system, see How to Design Data Pipelines for Reliable Reporting. For making pipelines self-healing before they even need to alert, see How to Build Self-Healing Data Pipelines.

Automation services include building monitoring and alerting infrastructure for production pipeline systems.

Get in touch to discuss setting up notifications for your data systems.

python notification system slack webhook python multi-channel alerts python email alerts python automation alert routing system python alert escalation notification rate limiting slack alerts data pipeline digest notifications python alert fatigue automation alert deduplication python

Enjoyed this article?

Get notified when I publish new articles on automation, ecommerce, and data engineering.

Get in touch

Related Articles