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
Quellen und Links
- Apache Airflow — Workflow Orchestration
- Prefect Cloud — Modern Data Workflows
- Dagster — Data Engineering Platform
- LangChain Text Splitter — Document Chunking
- Pinecone — Vector Database
- RAG: Retrieval-Augmented Generation — Original Paper
- Vector Search Best Practices — Performance Tuning
