Eine AI-Pipeline ist eine ETL-Maschine spezialisiert auf Daten für LLMs: Dokumente rein, Embeddings raus, Suchbar machen.

Das Problem: Ad-Hoc Data Processing

Szenario 1 (CHAOS):
┌─────────────┐
│ PDF kommt   │
└────┬────────┘
     │
     ├─→ Jemand extrahiert Text manuell
     ├─→ Kopiert in Word
     ├─→ Schickt es per Email
     ├─→ Jemand fügt ins System ein
     └─→ Fehler! PDF war alt!

Szenario 2 (GEORDNET):
┌─────────────────────────────────────────────────────────┐
│ Automated AI Pipeline                                   │
├─────────────────────────────────────────────────────────┤
│ 1. Data Ingestion: PDF via S3, Webhook, oder Polling    │
│ 2. Preprocessing: Text extract, clean, chunk            │
│ 3. Embedding: Batch embed 1000s documents               │
│ 4. Indexing: Upsert zu Vector DB                        │
│ 5. Monitoring: Log failures, retry, alert               │
└─────────────────────────────────────────────────────────┘

1. Pipeline Architektur (5 Stufen)

Stage 1: Ingestion         Stage 2: Processing        Stage 3: Enrichment
┌──────────────────┐       ┌──────────────────┐       ┌──────────────────┐
│ S3 / HTTP / DB   │ ──→   │ Extract Text     │ ──→   │ Embed            │
│ Files kommen rein│       │ Split Chunks     │       │ Compute Vectors  │
└──────────────────┘       │ Clean Metadata   │       └──────────────────┘
                           └──────────────────┘                │
                                                               ↓
                           Stage 5: Monitoring        Stage 4: Storage
                           ┌──────────────────┐       ┌──────────────────┐
                           │ Errors & Alerts  │←──────│ Vector DB        │
                           │ Performance      │       │ Cache / Metrics  │
                           │ Logs & Metrics   │       └──────────────────┘
                           └──────────────────┘

2. Apache Airflow: DAG-basierte Pipelines

DAG = Directed Acyclic Graph. Definiere Abhängigkeiten, Airflow führt es aus.

# dags/document_ingestion.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.bash import BashOperator
from airflow.utils.dates import days_ago
from datetime import timedelta

def fetch_documents():
    """Stage 1: Hole Dokumente von S3"""
    import boto3
    s3 = boto3.client('s3')

    response = s3.list_objects_v2(
        Bucket='document-bucket',
        Prefix='uploads/',
        MaxKeys=100
    )

    documents = []
    for obj in response.get('Contents', []):
        documents.append({
            'key': obj['Key'],
            'size': obj['Size'],
            'last_modified': obj['LastModified']
        })

    return documents

def extract_text(documents):
    """Stage 2: Extrahiere Text aus Dokumenten"""
    import PyPDF2
    import json

    extracted = []
    for doc in documents:
        if doc['key'].endswith('.pdf'):
            # PDF verarbeiten
            with open(f"/tmp/{doc['key']}", 'rb') as f:
                reader = PyPDF2.PdfReader(f)
                text = "\n".join([
                    page.extract_text()
                    for page in reader.pages
                ])

            extracted.append({
                'document_id': doc['key'],
                'text': text,
                'page_count': len(reader.pages)
            })

    return extracted

def chunk_documents(documents):
    """Stage 2b: Teile große Dokumente in Chunks"""
    from langchain.text_splitter import RecursiveCharacterTextSplitter

    splitter = RecursiveCharacterTextSplitter(
        chunk_size=512,
        chunk_overlap=50,
        separators=["\n\n", "\n", " ", ""]
    )

    chunks = []
    for doc in documents:
        doc_chunks = splitter.split_text(doc['text'])

        for i, chunk in enumerate(doc_chunks):
            chunks.append({
                'document_id': doc['document_id'],
                'chunk_index': i,
                'content': chunk,
                'tokens': len(chunk.split())
            })

    return chunks

def embed_chunks(chunks):
    """Stage 3: Embedde Chunks"""
    from sentence_transformers import SentenceTransformer
    import json

    model = SentenceTransformer('all-MiniLM-L6-v2')

    embeddings = []
    for chunk in chunks:
        embedding = model.encode(chunk['content'])

        embeddings.append({
            'document_id': chunk['document_id'],
            'chunk_index': chunk['chunk_index'],
            'content': chunk['content'],
            'embedding': embedding.tolist(),  # Liste statt Array
            'dimension': len(embedding)
        })

    # Batch speichern
    with open('/tmp/embeddings.jsonl', 'w') as f:
        for emb in embeddings:
            f.write(json.dumps(emb) + '\n')

    return embeddings

def index_embeddings(**context):
    """Stage 4: Speichere Embeddings in Vector DB"""
    from pinecone import Pinecone

    pc = Pinecone(api_key="YOUR_API_KEY")
    index = pc.Index("documents")

    # Lese embeddings aus Task vorher
    task_instance = context['task_instance']
    embeddings = task_instance.xcom_pull(
        task_ids='embed_chunks'
    )

    # Upsert zu Pinecone
    vectors_to_upsert = [
        (
            f"{emb['document_id']}#{emb['chunk_index']}",
            emb['embedding'],
            {
                'document_id': emb['document_id'],
                'chunk_index': emb['chunk_index'],
                'content': emb['content']
            }
        )
        for emb in embeddings
    ]

    index.upsert(vectors=vectors_to_upsert, batch_size=100)
    return len(vectors_to_upsert)

# DAG Definition
default_args = {
    'owner': 'ai-team',
    'retries': 2,
    'retry_delay': timedelta(minutes=5)
}

dag = DAG(
    'document_ingestion_pipeline',
    default_args=default_args,
    description='Ingeste und indexiere Dokumente',
    schedule_interval='0 2 * * *',  # Täglich 2 AM
    start_date=days_ago(1),
    catchup=False
)

# Tasks
fetch_task = PythonOperator(
    task_id='fetch_documents',
    python_callable=fetch_documents,
    dag=dag
)

extract_task = PythonOperator(
    task_id='extract_text',
    python_callable=extract_text,
    op_args=[fetch_task.output],
    dag=dag
)

chunk_task = PythonOperator(
    task_id='chunk_documents',
    python_callable=chunk_documents,
    op_args=[extract_task.output],
    dag=dag
)

embed_task = PythonOperator(
    task_id='embed_chunks',
    python_callable=embed_chunks,
    op_args=[chunk_task.output],
    pool='embedding_pool',  # GPU pool
    dag=dag
)

index_task = PythonOperator(
    task_id='index_embeddings',
    python_callable=index_embeddings,
    provide_context=True,
    dag=dag
)

# Dependencies
fetch_task >> extract_task >> chunk_task >> embed_task >> index_task

3. Prefect: Moderne Alternative zu Airflow

Prefect ist einfacher und schneller für KMU:

# flows/document_pipeline.py
from prefect import flow, task
from prefect.task_runs import in_process_task_run
import httpx

@task(retries=2, retry_delay_seconds=60)
async def fetch_from_s3(bucket: str) -> list:
    """Hole Dokumente"""
    import boto3
    s3 = boto3.client('s3')

    response = s3.list_objects_v2(Bucket=bucket, Prefix='uploads/')
    return [obj['Key'] for obj in response.get('Contents', [])]

@task
async def extract_text(document_path: str) -> str:
    """Extrahiere Text"""
    import PyPDF2

    with open(document_path, 'rb') as f:
        reader = PyPDF2.PdfReader(f)
        text = "\n".join([p.extract_text() for p in reader.pages])

    return text

@task
async def generate_embedding(text: str) -> list:
    """Embedde Text"""
    from sentence_transformers import SentenceTransformer

    model = SentenceTransformer('all-MiniLM-L6-v2')
    embedding = model.encode(text)
    return embedding.tolist()

@task
async def store_in_vector_db(
    document_id: str,
    embedding: list,
    text: str
):
    """Speichere in Vector DB"""
    from pinecone import Pinecone

    pc = Pinecone(api_key="API_KEY")
    index = pc.Index("documents")

    index.upsert([(
        document_id,
        embedding,
        {"content": text}
    )])

@flow
async def document_pipeline(bucket: str):
    """Main Pipeline Flow"""
    documents = await fetch_from_s3(bucket)

    for doc in documents:
        # Tasks laufen parallel wenn möglich
        text = await extract_text(doc)
        embedding = await generate_embedding(text)
        await store_in_vector_db(doc, embedding, text)

# Starte Pipeline
if __name__ == "__main__":
    document_pipeline("my-bucket")

4. RAG Pipeline (Retrieval-Augmented Generation)

RAG ist ein spezieller Pipeline-Typ: Retrieve relevant docs, dann generate.

# rag_pipeline.py
from datetime import datetime
from typing import List
import json

class RAGPipeline:
    def __init__(self, vector_db, llm_client):
        self.vector_db = vector_db
        self.llm = llm_client
        self.metrics = {
            'queries': 0,
            'avg_retrieval_time': 0,
            'avg_generation_time': 0
        }

    def retrieve(self, query: str, top_k: int = 5) -> List[str]:
        """Stage 1: Retrieve Relevant Documents"""
        import time

        start = time.time()

        # Embedde Query
        from sentence_transformers import SentenceTransformer
        model = SentenceTransformer('all-MiniLM-L6-v2')
        query_embedding = model.encode(query).tolist()

        # Vector Similarity Search
        results = self.vector_db.search(
            vector=query_embedding,
            top_k=top_k,
            include_metadata=True
        )

        elapsed = time.time() - start
        self.metrics['avg_retrieval_time'] = elapsed

        return [
            result['metadata']['content']
            for result in results
        ]

    def generate(self, query: str, context: List[str]) -> str:
        """Stage 2: Generate Response mit Context"""
        import time

        start = time.time()

        # Build Prompt mit Context
        context_str = "\n\n".join([
            f"[Document {i}]\n{doc}"
            for i, doc in enumerate(context)
        ])

        prompt = f"""Context:
{context_str}

Query: {query}

Answer based on the provided context:"""

        # Generate
        response = self.llm.generate(prompt)

        elapsed = time.time() - start
        self.metrics['avg_generation_time'] = elapsed

        return response

    def execute(self, query: str) -> dict:
        """Full RAG Pipeline"""
        self.metrics['queries'] += 1

        # Stage 1: Retrieve
        context = self.retrieve(query)

        # Stage 2: Generate
        response = self.generate(query, context)

        return {
            'query': query,
            'context': context,
            'response': response,
            'metrics': {
                'retrieval_ms': self.metrics['avg_retrieval_time'] * 1000,
                'generation_ms': self.metrics['avg_generation_time'] * 1000,
                'total_queries': self.metrics['queries']
            }
        }

# Nutzung
rag = RAGPipeline(vector_db, llm_client)

result = rag.execute("Was sind die neuen Features?")
print(result['response'])

5. Batch vs. Real-Time Pipelines

Batch Pipeline (Night-Runs):

23:00 - Start
  ├─ Lade 1000s neue Dokumente
  ├─ Embedde alle parallel
  ├─ Indexiere in Vector DB
01:00 - Fertig

Vorteil: Schnell, kosteneffizient (GPU-Auslastung)
Nachteil: Neue Daten erst morgen suchbar

Real-Time Pipeline (On-Demand):

User hochlä Datei
  ↓ (100ms)
Webhook triggered
  ↓ (50ms)
Extract Text
  ↓ (200ms)
Embed (Async)
  ↓ (500ms)
Indexiert
  ↓ (50ms)
"Datei ist jetzt suchbar"

Vorteil: Instant
Nachteil: Mehr Servercapacity needed

Hybrid (Best of Both):

Real-Time Pipeline für:
  - User-uploaded Files
  - Live-Streaming-Daten
  - Hot Topics

Batch Pipeline für:
  - News/RSS Feeds (täglich)
  - Archive Backups (weekly)
  - Model Re-training (monthly)

6. Error Handling und Retries

from prefect import task
from prefect.task_runs import task_run_context

@task(
    retries=3,
    retry_delay_seconds=60,
    retry_jitter_seconds=15
)
def embed_with_retry(chunk: str) -> list:
    """Embedde mit automatischen Retries"""
    try:
        return embed_model.encode(chunk)
    except Exception as e:
        context = task_run_context.get()

        # Log für Monitoring
        print(f"Embed failed (attempt {context.task_run.total_run_count}): {e}")

        # Nach 3 Fehlversuchen: Fallback
        if context.task_run.total_run_count >= 3:
            return [0] * 384  # Dummy embedding

        raise  # Prefect retried

@task
def embed_with_circuit_breaker(chunk: str, circuit_breaker) -> list:
    """Mit Circuit Breaker Pattern"""
    if circuit_breaker.is_open():
        # Service down → Nutze Cache statt Retry
        return cache.get(chunk)

    try:
        return embed_model.encode(chunk)
    except Exception as e:
        circuit_breaker.record_failure()
        if circuit_breaker.should_open():
            circuit_breaker.open()
        raise

7. Pipeline Monitoring

from prometheus_client import Counter, Histogram
import logging

# Metrics
docs_processed = Counter('pipeline_docs_processed_total', 'Total docs')
chunks_created = Counter('pipeline_chunks_created_total', 'Total chunks')
embedding_time = Histogram('pipeline_embedding_seconds', 'Embedding latency')
pipeline_errors = Counter('pipeline_errors_total', 'Errors', ['stage'])

logger = logging.getLogger(__name__)

class MonitoredPipeline:
    def process_document(self, doc_path: str):
        """Process mit Monitoring"""
        try:
            # Stage 1
            text = extract(doc_path)
            docs_processed.inc()

            # Stage 2
            with embedding_time.time():
                chunks = split_and_embed(text)
            chunks_created.inc(len(chunks))

            # Stage 3
            index_chunks(chunks)

            logger.info(f"Processed {doc_path}: {len(chunks)} chunks")

        except Exception as e:
            pipeline_errors.labels(stage=type(e).__name__).inc()
            logger.error(f"Failed to process {doc_path}: {e}")
            raise

# Prometheus Scrape Endpoint
from prometheus_client import generate_latest

@app.get("/metrics")
def metrics():
    return generate_latest()

Zusammenfassung: Pipeline-Architektur

Anforderung Tool Komplexität
Einfache Batch Cron + Python Niedrig
Komplexe DAGs Airflow Hoch
Modern + Fast Prefect Mittel
Datenengineering Dagster Hoch
Streaming Kafka + Spark Sehr Hoch

Empfehlung:

  • Start: Cron + Python
  • Scale: Prefect
  • Enterprise: Airflow + Kubernetes