What an AI Data Pipeline Actually Is

Most people hear "data pipeline" and think ETL — extract, transform, load. AI data pipelines are that plus three additional stages that make the data useful to language models: embed, index, and serve. The full lifecycle looks like this:

  1. Ingest — pull raw data from APIs, scrapers, databases, file drops, or webhooks
  2. Transform — clean, normalize, deduplicate, and structure the data
  3. Enrich — use an LLM to classify, tag, summarize, or extract structured fields
  4. Embed — convert text into vector representations for semantic search
  5. Index — write vectors + metadata to a vector store or search index
  6. Serve — expose the processed data through an API that your AI application queries at runtime

I have built dozens of these pipelines for clients. The pattern is remarkably consistent regardless of whether you are processing product catalogs, support tickets, legal documents, or real estate listings. The data shapes differ, but the architecture is the same.

The goal is not just automation — it is a pipeline that detects its own failures, retries intelligently, and alerts you only when it genuinely needs human intervention. If you are waking up to fix your pipeline every morning, it is not automated. It is a cron job with a human dependency.

Real Example: Automated Product Data Pipeline

I built a pipeline for an e-commerce client that needed to keep their AI-powered product search current with inventory from three different suppliers. Each supplier had a different data format. The pipeline runs every four hours and handles roughly 12,000 products per cycle.

Here is the core pipeline class, simplified from the production version:

import requests
import json
import hashlib
import time
from datetime import datetime, timezone

class ProductPipeline:
    def __init__(self, config):
        self.config = config
        self.stats = {"ingested": 0, "transformed": 0,
                      "enriched": 0, "embedded": 0, "errors": 0}
        self.dead_letter = []

    def run(self):
        """Full pipeline cycle: ingest -> transform -> enrich -> embed -> index."""
        raw_products = self.ingest()
        cleaned = self.transform(raw_products)
        enriched = self.enrich_batch(cleaned)
        self.embed_and_index(enriched)
        self.report()

    def ingest(self):
        """Pull from all configured supplier APIs."""
        products = []
        for supplier in self.config["suppliers"]:
            try:
                resp = requests.get(
                    supplier["url"],
                    headers=supplier.get("headers", {}),
                    timeout=30
                )
                resp.raise_for_status()
                items = resp.json().get("products", [])
                for item in items:
                    item["_source"] = supplier["name"]
                products.extend(items)
                self.stats["ingested"] += len(items)
            except Exception as e:
                self.dead_letter.append({
                    "stage": "ingest",
                    "supplier": supplier["name"],
                    "error": str(e),
                    "timestamp": datetime.now(timezone.utc).isoformat()
                })
                self.stats["errors"] += 1
        return products

    def transform(self, raw):
        """Normalize fields, deduplicate, validate."""
        seen_skus = set()
        cleaned = []
        for item in raw:
            sku = item.get("sku") or item.get("product_id") or item.get("id")
            if not sku or sku in seen_skus:
                continue
            seen_skus.add(sku)

            product = {
                "sku": str(sku),
                "title": (item.get("title") or item.get("name") or "").strip(),
                "description": (item.get("description") or "").strip(),
                "price": self._parse_price(item.get("price")),
                "in_stock": bool(item.get("in_stock", item.get("available", False))),
                "category": (item.get("category") or "uncategorized").lower(),
                "source": item["_source"],
                "updated_at": datetime.now(timezone.utc).isoformat()
            }

            if product["title"] and product["price"] is not None:
                cleaned.append(product)
                self.stats["transformed"] += 1
        return cleaned

    def _parse_price(self, val):
        if val is None:
            return None
        if isinstance(val, (int, float)):
            return round(float(val), 2)
        try:
            return round(float(str(val).replace("$", "").replace(",", "")), 2)
        except (ValueError, TypeError):
            return None

Notice how every stage is isolated. If supplier B fails, supplier A and C still process. Failed items go to the dead letter queue, not into the pipeline output. This is critical — a single bad record should never crash an entire pipeline run.

LLM Enrichment Stage

The most interesting stage is enrichment, where an LLM extracts structured data that the suppliers did not provide. For this client, the LLM categorizes products into a standardized taxonomy, extracts key specs from freeform descriptions, and generates search-friendly summaries.

def enrich_batch(self, products, batch_size=20):
    """Enrich products with LLM-extracted metadata in batches."""
    enriched = []
    for i in range(0, len(products), batch_size):
        batch = products[i:i + batch_size]
        prompt = self._build_enrichment_prompt(batch)

        for attempt in range(3):
            try:
                resp = requests.post(
                    "https://api.anthropic.com/v1/messages",
                    headers={
                        "x-api-key": self.config["anthropic_key"],
                        "content-type": "application/json",
                        "anthropic-version": "2023-06-01"
                    },
                    json={
                        "model": "claude-sonnet-4-20250514",
                        "max_tokens": 4096,
                        "messages": [{"role": "user", "content": prompt}]
                    },
                    timeout=60
                )
                resp.raise_for_status()
                result = resp.json()["content"][0]["text"]
                parsed = json.loads(result)

                for product, meta in zip(batch, parsed):
                    product["tags"] = meta.get("tags", [])
                    product["specs"] = meta.get("specs", {})
                    product["summary"] = meta.get("summary", "")
                    enriched.append(product)
                    self.stats["enriched"] += 1
                break  # success, exit retry loop

            except requests.exceptions.HTTPError as e:
                if e.response.status_code == 429:
                    wait = int(e.response.headers.get("retry-after", 30))
                    time.sleep(wait)
                elif e.response.status_code >= 500:
                    time.sleep(2 ** attempt)
                else:
                    self.dead_letter.append({
                        "stage": "enrich", "batch_start": i,
                        "error": str(e),
                        "timestamp": datetime.now(timezone.utc).isoformat()
                    })
                    self.stats["errors"] += 1
                    break
            except (json.JSONDecodeError, KeyError) as e:
                self.dead_letter.append({
                    "stage": "enrich", "batch_start": i,
                    "error": f"Parse error: {e}",
                    "timestamp": datetime.now(timezone.utc).isoformat()
                })
                self.stats["errors"] += 1
                break
    return enriched

Key details: I batch 20 products per LLM call to amortize the per-call overhead. The retry logic distinguishes between rate limits (wait and retry), server errors (exponential backoff), and client errors (dead-letter immediately, do not retry). A parse failure on the LLM output also goes to the dead letter queue rather than crashing the pipeline.

Scheduling and Orchestration

I see teams reach for Apache Airflow before they need it. Airflow is powerful but heavy — it requires a metadata database, a scheduler daemon, and its own deployment infrastructure. For 90% of data pipelines I build, the right tool is simpler.

The progression I recommend:

  1. Python + cron / Task Scheduler — your pipeline is a single script. Schedule it with crontab on Linux or Task Scheduler on Windows. This handles everything up to about 5 pipeline scripts running at different intervals. Zero infrastructure overhead.
  2. Python + systemd timers — more reliable than cron for production. Supports dependencies, logging, and automatic restart. Still zero extra infrastructure.
  3. Lightweight orchestrator — when you have 10+ pipelines with dependencies between them, use something like Prefect or Dagster. They give you a DAG-based execution model without requiring you to run Airflow's infrastructure.
  4. Airflow — only when you need multi-team collaboration, complex branching logic, or compliance-grade audit logging. If you are a one-person operation, you almost certainly do not need it.
I run pipelines processing 50,000+ records daily on a systemd timer calling a Python script. The entire "orchestration layer" is a 6-line .service file and a .timer file. It has been running for months without intervention. Do not over-engineer scheduling.

Here is what the systemd setup looks like:

# /etc/systemd/system/product-pipeline.timer
[Unit]
Description=Product data pipeline (every 4 hours)

[Timer]
OnCalendar=*-*-* 00/4:00:00
Persistent=true
RandomizedDelaySec=120

[Install]
WantedBy=timers.target

# /etc/systemd/system/product-pipeline.service
[Unit]
Description=Product pipeline runner

[Service]
Type=oneshot
ExecStart=/opt/pipeline/venv/bin/python /opt/pipeline/run.py
Environment=PIPELINE_ENV=production
TimeoutStartSec=1800

The Persistent=true flag means if the machine was off during a scheduled run, systemd fires the pipeline immediately on boot. RandomizedDelaySec prevents thundering herd if you have multiple pipelines scheduled at the same time.

Error Handling and Self-Healing

A pipeline that runs unattended needs three layers of error handling:

1. Per-item error isolation

Every item processed inside a try/except. A malformed product record from supplier C does not kill the entire run. Failed items go to a dead letter queue (a JSON file or database table) that you can inspect and replay later.

2. Stage-level retries with backoff

External calls (APIs, databases, LLMs) get automatic retries. I use a simple pattern:

def retry(fn, max_attempts=3, backoff=2):
    """Retry with exponential backoff. Returns (result, error)."""
    for attempt in range(max_attempts):
        try:
            return fn(), None
        except Exception as e:
            if attempt == max_attempts - 1:
                return None, e
            time.sleep(backoff ** attempt)
    return None, Exception("exhausted retries")

3. Pipeline-level alerting

After every run, the pipeline sends a structured report. If errors exceed a threshold, it escalates to a notification channel. I use ntfy for push notifications because it is self-hostable and has zero dependencies:

def report(self):
    """Send pipeline run summary. Alert if error rate is high."""
    total = self.stats["ingested"] or 1
    error_rate = self.stats["errors"] / total

    summary = (
        f"Pipeline run complete: {self.stats['ingested']} ingested, "
        f"{self.stats['transformed']} transformed, "
        f"{self.stats['enriched']} enriched, "
        f"{self.stats['embedded']} embedded, "
        f"{self.stats['errors']} errors ({error_rate:.1%})"
    )

    if error_rate > 0.1:  # more than 10% failures
        requests.post(
            "https://ntfy.example.com/pipeline-alerts",
            data=summary,
            headers={"Priority": "high", "Title": "Pipeline alert"}
        )
    elif self.dead_letter:
        # write dead letter items for later inspection
        with open("/var/log/pipeline/dead_letter.jsonl", "a") as f:
            for item in self.dead_letter:
                f.write(json.dumps(item) + "\n")
The self-healing part is not magic. It is exhaustive error handling at every layer: item, batch, stage, and run. Each layer catches what the layer below missed. The result is a pipeline that handles 95% of failures automatically and only pages you for the remaining 5%.

Monitoring and Observability

Running a pipeline is easy. Knowing when it has silently started producing garbage is hard. These are the signals I monitor:

  • Record count drift — if the pipeline usually processes 12,000 products and today it processed 3,000, something changed at the source. Alert when count deviates more than 30% from the 7-day moving average.
  • Embedding staleness — track the oldest updated_at timestamp in the vector store. If any product has not been refreshed in 24 hours (for a 4-hour pipeline), the pipeline is silently failing on that subset.
  • Dead letter growth — the dead letter queue should drain, not grow. If it grows across three consecutive runs, there is a systematic issue, not transient noise.
  • LLM output quality — spot-check a random sample of LLM-enriched records each run. If the JSON parse failure rate exceeds 5%, the prompt needs adjustment or the model changed behavior.
  • Latency per stage — log how long each stage takes. A sudden spike in the transform stage might mean you are receiving 10x more data than expected. A spike in the embed stage might mean the embedding API is throttling you.

All of this fits in a simple JSON log line per run. No Grafana, no Prometheus, no Datadog — just a structured log file that you can query with jq when something looks wrong.

Cost Control

The LLM enrichment stage is where costs can spiral. Here are the techniques I use to keep AI pipeline costs predictable:

Batch over real-time

Processing 12,000 products in a single batch run every 4 hours costs a fraction of processing each product change as it happens via webhook. Batching lets you deduplicate (product updated 6 times in 4 hours? Process once), use larger batches per LLM call, and amortize API overhead.

Cache embeddings

If a product title and description have not changed since the last run, do not re-embed it. Hash the text content and compare to the stored hash. For a 12,000-product catalog where 200-300 products change per cycle, this reduces embedding API calls by 97%.

def needs_reembed(self, product, existing_hash):
    """Only re-embed if content actually changed."""
    content = f"{product['title']} {product['description']} {product['summary']}"
    current_hash = hashlib.sha256(content.encode()).hexdigest()
    return current_hash != existing_hash, current_hash

Model routing for enrichment

Not every enrichment task needs a frontier model. I use a tiered approach: simple category classification goes to Haiku (pennies per thousand products), spec extraction from structured descriptions goes to Haiku, and only freeform description summarization hits Sonnet. This cuts enrichment costs by 60-70% compared to sending everything to the same model.

Set hard budget caps

Every pipeline run has a maximum API spend. If the pipeline hits the cap mid-run, it finishes processing what it has and alerts you rather than silently burning through your budget. I have seen runaway pipelines consume $200 in API costs in a single night because a data source started returning 100x more records than expected.

When to Build Custom vs. Use Off-the-Shelf

Tools like Zapier and Make are excellent for connecting SaaS applications: "when a new row appears in Google Sheets, send a Slack message." They fall apart for AI data pipelines because:

  • Volume — Zapier charges per task. Processing 12,000 products with 5 steps each is 60,000 tasks. At $0.01/task on their premium plan, that is $600/month for a pipeline that costs $15/month to run on a $5 VPS.
  • LLM integration — Zapier's AI steps are limited to simple text generation. You cannot batch 20 products into a single LLM call, implement custom retry logic, or parse structured JSON output with error recovery.
  • Error handling — Zapier retries the entire zap, not individual items. One bad record retries all 12,000 products through every step.
  • Latency — each step in a Zapier workflow has network overhead. A 6-stage pipeline with 12,000 items becomes painfully slow compared to a Python script running everything in-memory with batched API calls.

My rule: if your pipeline touches an LLM, processes more than 100 items per run, or needs per-item error isolation, build it in Python. The upfront investment is 2-3 days versus 2-3 hours for Zapier, but the operational cost and reliability difference is enormous over the lifetime of the system.

The best pipeline is the one nobody thinks about. If it runs, handles its own errors, alerts when something genuinely needs attention, and costs less than a coffee per day — it is doing its job. Everything else is over-engineering.

Getting Started

If you are building your first AI data pipeline, start here:

  1. Define your five stages explicitly. Write a one-sentence description of what ingest, transform, enrich, embed, and serve mean for your specific data. If you cannot describe each stage clearly, you do not understand your pipeline yet.
  2. Build the pipeline without the LLM first. Get ingest, transform, and index working with dummy enrichment data. This forces you to get the plumbing right before adding the expensive part.
  3. Add the LLM enrichment last. With the pipeline structure solid, adding the LLM stage is a matter of writing one function and slotting it in.
  4. Run it manually ten times before automating. Each run will reveal edge cases in your data that you did not anticipate. Fix those before handing control to a scheduler.
  5. Monitor the dead letter queue. It tells you exactly where your pipeline is fragile. The patterns in the dead letter queue drive your next round of improvements.

Related Articles

Cost OptimizationLLM

Your AI Bill Is 10x What It Should Be

Model routing, prompt caching, and context management techniques that cut AI costs by 70-90%.

RAGArchitecture

Building a RAG System That Actually Works

Retrieval-augmented generation done right: chunking, embedding, retrieval, and evaluation.

AutonomousAgents

The Autonomous AI Operator

Building AI agents that take action, monitor results, and course-correct without human intervention.