RAG
RAG - Retrieval-Augmented Generation
Diese Technologie ist derzeit sehr populär. Doch ehrlich gesagt bleiben Fragen offen, da die Knowledge Base vektoriell und letztlich statistisch ist und sich daher nicht für den Flugzeugbau oder die Durchführung medizinischer Operationen eignet. Ich bevorzuge neuro-symbolische KI, bei der die Wissensbasis beispielsweise aus den exakten S-Ausdrücken von Lisp besteht, die wir hier noch genauer betrachten werden.
1. ALL-IN-ONE Lösung
Wir stellen ein Schweizer Taschenmesser für alle Ihre Bedürfnisse her, sodass Sie nur für LLM-API bezahlen müssen (falls Sie kein eigenes besitzen).
Welche Klingen hat unser Schweizer Taschenmesser?
- Knowledge Bases - PostgreSQL(pgvector)
- Knowledge Base Pumping - ".txt", ".md", ".json", ".log", ".py", ".csv", ".pdf", ".docx", ".html", ".htm"
- Textvektorisierung mit Sentence-Transformers
- Generieren von Suchoptionen durch vLLM.
- NLP-SpacyChunker
- SEARCH ENGINE (BM25 + DENSE + RRF FUSION)
- MCP-Generator (FastMCP)
Wir benötigen: 3 Server (vorzugsweise über ein lokales Netzwerk verbunden), 3 Subdomains mit A-Records.
- Server 1 / Subdomain 1: llm.domain.de (10.0.0.5)
- Server 2 / Subdomain 2: vpg.domain.de (10.0.0.7)
- Server 3 / Subdomain 3: rag.domain.de (10.0.0.9)
2. Anweisungen zum Bereitstellen der Umgebung Server 1
Qwen2.5-0.5B-Instruct - Das Modell wird von Alibaba unter der Apache-2.0-Open-Source-Lizenz vertrieben.
- Laden Sie es herunter und führen Sie es auf Ihrem Server oder Computer ohne Lizenzgebühren aus.
- Nutzen Sie es für private, kommerzielle oder Forschungsprojekte.
- Ändern und in Ihre Anwendungen integrieren
Referenz:
- Das Qwen2.5-0.5B-Instruct-Modell ist ein extrem schlankes, universelles Sprachmodell mit nur 0,5 Milliarden Parametern (0,5 B). Der Zusatz „Instruct“ im Namen des Modells deutet darauf hin, dass es speziell für das Verständnis von Textanweisungen und Dialogen optimiert ist und nicht einfach nur Text fortsetzt.
- Mikroaufgaben und Hilfslogik in RAG: Aufgrund seiner geringen Größe eignet es sich ideal als schnelles lokales "Gehirn" in Skripten - zum Beispiel zum Umformulieren von Suchanfragen (QueryRewriter), zum Extrahieren von Entitäten oder zum Filtern von Datenmüll, wenn große Modelle unpraktisch sind.
- Funktioniert auch auf leistungsschwachen Servern und CPUs: Es benötigt nur minimalen Arbeitsspeicher (buchstäblich einige hundert Megabyte) und ist in der Lage, selbst auf einem Standardprozessor (CPU) ohne den Einsatz teurer Grafikkarten hohe Geschwindigkeiten zu erreichen.
- Lokale Automatisierung: Geeignet für einfache Aufgaben wie Textklassifizierung, Zusammenfassung kurzer Passagen, Generierung von JSON-Strukturen aus einer Vorlage und einfache Chatbot-Konversationen, die keine tiefgreifende Expertenlogik erfordern.
Wenn es über das Internet funktioniert
python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip
pip install vllm
python -m vllm.entrypoints.openai.api_server --help
sudo -s
tee /etc/systemd/system/vllm.service > /dev/null << 'EOF'
[Unit]
Description=vLLM OpenAI-Compatible API Server (CPU Mode)
After=network.target
[Service]
Type=simple
User=root
WorkingDirectory=/path/to/your/project
Environment="VLLM_TARGET_DEVICE=cpu"
ExecStart=/path/to/your/.venv/bin/python -m vllm.entrypoints.openai.api_server --model Qwen/Qwen2.5-0.5B-Instruct --host 127.0.0.1 --port 8000 --device cpu
Restart=always
RestartSec=10
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload
systemctl start vllm
systemctl enable vllm
systemctl status vllm
apt install nginx
tee /etc/nginx/sites-available/vllm > /dev/null << 'EOF'
server {
listen 80;
server_name llm.domain.de;
location / {
proxy_pass http://127.0.0.1:8000;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_read_timeout 600s;
proxy_send_timeout 600s;
}
}
EOF
ln -s /etc/nginx/sites-available/vllm /etc/nginx/sites-enabled/
nginx -t
/etc/init.d/nginx restart
apt install certbot python3-certbot-nginx -y
certbot --nginx
crontab -e
0 3 * * * /usr/bin/certbot renew --quiet --deploy-hook "systemctl reload nginx"
certbot --nginx --redirect
nginx -t
/etc/init.d/nginx restart
exit
Wenn ein lokales Netzwerk vorhanden ist
python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip
pip install vllm
python -m vllm.entrypoints.openai.api_server --help
sudo -s
tee /etc/systemd/system/vllm.service > /dev/null << 'EOF'
[Unit]
Description=vLLM OpenAI-Compatible API Server (CPU Mode)
After=network.target
[Service]
Type=simple
User=root
WorkingDirectory=/path/to/your/project
Environment="VLLM_TARGET_DEVICE=cpu"
ExecStart=/path/to/your/.venv/bin/python -m vllm.entrypoints.openai.api_server --model Qwen/Qwen2.5-0.5B-Instruct --host 127.0.0.1 --port 8000 --device cpu
Restart=always
RestartSec=10
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload
systemctl start vllm
systemctl enable vllm
systemctl status vllm
apt install nginx
tee /etc/nginx/sites-available/vllm > /dev/null << 'EOF'
server {
listen 10.0.0.5:80;
location / {
proxy_pass http://127.0.0.1:8000;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_read_timeout 600s;
proxy_send_timeout 600s;
}
}
EOF
ln -s /etc/nginx/sites-available/vllm /etc/nginx/sites-enabled/
nginx -t
/etc/init.d/nginx restart
exit
3. Anweisungen zum Bereitstellen der Umgebung Server 2
Eine Architekturlösung zur Aufgabe von Multi-Datenbank-Konfigurationen (Qdrant/Chroma) zugunsten eines einzigen PostgreSQL-Speichers:
- 1. BM25 Indexproblem (OOM-Engpass): Bei der Verwendung von In-Memory-Lösungen (z. B. rank_bm25 über lokales SQLite/ChromaDB/Qdrant) muss der gesamte Lemma-Korpus in den Arbeitsspeicher geladen werden, um BM25Okapi zu generieren bei jedem Aufruf von reload_index(). Bei Millionen von Datenblöcken führt dies zu Speichermangel (OOM) und anhaltenden CPU-Sperren. Die Verwendung von PostgreSQL ermöglicht den Wechsel zur nativen Volltextsuche (FTS) oder zu indiziertem JSONB, wodurch das Laden des Korpus in den Arbeitsspeicher entfällt.
- 2. Optimierung der Vektorsuche in PostgreSQL mit HNSW: Standardmäßig führt die Kosinusdistanzsuche (`vector <=> query_vector`) einen sequenziellen Scan durch (eine vollständige Suche aller Tabellenzeilen). Die Verwendung eines HNSW-Index (Hierarchical Navigable Small World) erstellt einen Nächste-Nachbarn-Graphen direkt im DBMS und reduziert so die Suchkomplexität von O(N) auf O(log N).
sudo -s
apt install -y wget gnupg
sh -c 'echo "deb http://apt.postgresql.org/pub/repos/apt $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list'
wget --quiet -O - https://www.postgresql.org/media/keys/ACCC4CF8.asc | sudo apt-key add -
apt update
apt install -y postgresql-16 postgresql-client-16 postgresql-16-pgvector
su - postgres -c psql
CREATE DATABASE rag_db;
\c rag_db;
CREATE EXTENSION vector;
\dx
\q
vi /etc/postgresql/16/main/postgresql.conf
listen_addresses = '*'
vi /etc/postgresql/16/main/pg_hba.conf
# dostup tolko po localnoj seti
host rag_db all 10.0.0.0/24 scram-sha-256
# dospup cherez iternet, no s konkretnogo IP
hostssl rag_db all xxx.x.xxx.xx/32 scram-sha-256
# polnih dostup
hostssl rag_db all 0.0.0.0/0 scram-sha-256
su - postgres -c psql
CREATE USER rag_user WITH ENCRYPTED PASSWORD 'XXXXXXXXXXXXXXXXX';
GRANT ALL PRIVILEGES ON DATABASE rag_db TO rag_user;
\c rag_db;
GRANT ALL ON SCHEMA public TO rag_user;
\q
systemctl restart postgresql
netstat -na | grep 5432
# Wenn UFW vorhanden ist und der Port geschlossen ist
# für lokale Netzwerke
ufw allow from 10.0.0.0/24 to any port 5432 proto tcp
# Für externen Zugriff:
ufw allow 5432/tcp
# wenn hostssl
apt install certbot -y
certbot certonly --standalone -d vpg.domain.de
mkdir -p /etc/postgresql/16/main/certs
cp /etc/letsencrypt/live/vpg.domain.de/fullchain.pem /etc/postgresql/16/main/certs/server.crt
cp /etc/letsencrypt/live/vpg.domain.de/privkey.pem /etc/postgresql/16/main/certs/server.key
chown -R postgres:postgres /etc/postgresql/16/main/certs
chmod 600 /etc/postgresql/16/main/certs/server.key
chmod 644 /etc/postgresql/16/main/certs/server.crt
vi /etc/postgresql/16/main/postgresql.conf
ssl = on
ssl_cert_file = 'certs/server.crt'
ssl_key_file = 'certs/server.key'
systemctl restart postgresql
tee /etc/letsencrypt/renewal-hooks/deploy/postgresql-ssl.sh > /dev/null << 'EOF'
#!/bin/bash
DOMAIN="vpg.domain.de"
CERTS_DIR="/etc/postgresql/16/main/certs"
if [ "$RENEWED_LINEAGE" = "/etc/letsencrypt/live/$DOMAIN" ]; then
cp /etc/letsencrypt/live/$DOMAIN/fullchain.pem $CERTS_DIR/server.crt
cp /etc/letsencrypt/live/$DOMAIN/privkey.pem $CERTS_DIR/server.key
chown -R postgres:postgres $CERTS_DIR
chmod 600 $CERTS_DIR/server.key
chmod 644 $CERTS_DIR/server.crt
systemctl reload postgresql
fi
EOF
chmod +x /etc/letsencrypt/renewal-hooks/deploy/postgresql-ssl.sh
exit
4. Anweisungen zum Bereitstellen der Umgebung Server 3
Unsere API befindet sich hier.
sudo -s
apt install nginx
cd /var/www
mkdir rag-api
cd rag-api
python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip setuptools wheel
Erstellen Sie eine Datei requirements.txt
spacy>=3.7.0
sentence-transformers>=2.3.0
torch>=2.0.0
aiohttp>=3.8.0
psycopg2-binary>=2.9.0
rank-bm25>=0.2.2
spacy>=3.0.0
sentence-transformers>=2.2.0
pypdf>=3.0.0
python-docx>=0.8.11
beautifulsoup4>=4.9.0
Flask>=2.0.0
openai>=1.0.0
gunicorn>=20.0.0
Wir installieren alles Notwendige
pip install -r requirements.txt
Unsere internen Neuronen netze installieren
python3 -m spacy download de_core_news_lg
# "de" - kann ersetzt werden durch "ru", "en", "fr", "es", ...
# "lg" - kann ersetzt werden durch "sm", "md"
python -c "from sentence_transformers import SentenceTransformer; SentenceTransformer('sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2')"
# Ein wichtiger Punkt: Wir müssen die Frage klären: GPU oder CPU? Falls Sie sich für die GPU entscheiden, können wir fortfahren. Falls Sie die CPU wählen, überprüfen wir, ob alles ordnungsgemäß funktioniert.
python -c "from sentence_transformers import SentenceTransformer; m = SentenceTransformer('sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2', device='cpu'); print(m.device)"
Daher ist in den meisten Fällen keine Konfiguration erforderlich – das Modell wird automatisch auf der CPU ausgeführt. Wenn Sie die Ausführung jedoch fest auf der CPU festlegen möchten (beispielsweise auf einem Server ohne GPU, um unnötige Geräteprüfungen zu vermeiden), müssen Sie das Gerät bei der Initialisierung des Modells explizit angeben (device="cpu").
OpenMP/MKL Thread-Begrenzung
import torch
# Geben Sie die Anzahl der Threads an. Wenn Sie 8 Threads haben, geben Sie 4 Threads an. Sie müssen die Hälfte der Threads freilassen, da wir "ThreadPoolExecutor" verwenden.
torch.set_num_threads(4)
Wenn Geschwindigkeit bei schwachen CPUs entscheidend ist, kann das Modell in das int8-Format konvertiert werden, aber für ein leichtes Modell wie das MiniLM-L12-v2 ist dies auf Standardprozessoren normalerweise nicht notwendig – es läuft bereits sehr schnell (zehn bis hunderte von Texten pro Sekunde).
Erstellen Sie eine Datei rag.py
import argparse
import asyncio
import json
import math
import os
from concurrent.futures import ThreadPoolExecutor
from dataclasses import asdict, dataclass, field
from pathlib import Path
from typing import Any, Dict, List, Optional, Set, Union, cast
import aiohttp
from rank_bm25 import BM25Okapi
try:
import psycopg2
import psycopg2.extras
POSTGRES_AVAILABLE = True
except ImportError:
POSTGRES_AVAILABLE = False
# --- INTERFACE DATA ---
@dataclass
class TextChunk:
id: str
doc_id: str
source_file: str
text: str
lemmas: List[str]
entities: Dict[str, List[str]]
vector: Optional[List[float]] = None
metadata: Dict[str, Any] = field(default_factory=dict)
def to_dict(self) -> Dict[str, Any]:
d = dict(self.__dict__)
if d.get("vector"):
d["vector_dim"] = len(d["vector"])
del d["vector"]
return d
# --- INTERFACE SEARCH ---
@dataclass
class SearchResult:
chunk_id: str
doc_id: str
source_file: str
text: str
rrf_score: float
bm25_score: float
dense_score: float
metadata: Dict[str, Any]
# --- DOCUMENT READER ENGINE ---
class DocumentReader:
"""Extrahiert Text und Metadaten aus lokalen Dateien verschiedener Formate."""
@staticmethod
def read_file(file_path: Union[str, Path]) -> Dict[str, Any]:
path = Path(file_path)
if not path.exists():
raise FileNotFoundError(f"Файл не найден: {path}")
ext = path.suffix.lower().replace(".", "")
content = ""
meta = {
"file_name": path.name,
"file_path": str(path.absolute()),
"file_size": path.stat().st_size,
"extension": ext,
}
# Textdateien, Markdown und Quellcode
if ext in ["txt", "md", "json", "log", "py", "csv"]:
with open(path, "r", encoding="utf-8", errors="ignore") as f:
content = f.read()
# PDF-Dokumente
elif ext == "pdf":
try:
import pypdf
reader = pypdf.PdfReader(path)
pages = [
page.extract_text()
for page in reader.pages
if page.extract_text()
]
content = "\n\n".join(pages)
meta["num_pages"] = len(reader.pages)
except ImportError:
raise ImportError("Zum Lesen von PDFs installieren: pip install pypdf")
# Microsoft Word
elif ext == "docx":
try:
import docx
doc = docx.Document(path)
content = "\n".join(
[p.text for p in doc.paragraphs if p.text.strip()]
)
except ImportError:
raise ImportError(
"Zum Lesen von DOCX-Dateien installieren: pip install python-docx"
)
# HTML
elif ext in ["html", "htm"]:
try:
from bs4 import BeautifulSoup
with open(path, "r", encoding="utf-8", errors="ignore") as f:
soup = BeautifulSoup(f.read(), "html.parser")
for script in soup(["script", "style"]):
script.decompose()
content = soup.get_text(separator="\n")
except ImportError:
raise ImportError(
"Um HTML lesen zu können, installieren: pip install beautifulsoup4"
)
else:
with open(path, "r", encoding="utf-8", errors="ignore") as f:
content = f.read()
return {"content": content, "metadata": meta}
# --- SPACY NLP CORE ENGINE (OPTIMIZED) ---
class SpacyNLPProcessor:
"""Ein eigenständiger NLP-Prozessor auf Basis von spaCy mit optimierter Single-Pass-Verarbeitung."""
def __init__(self, model_name: str = "de_core_news_lg"):
try:
import spacy
self.nlp = spacy.load(model_name)
except OSError:
import spacy.cli
spacy.cli.download(model_name)
self.nlp = spacy.load(model_name)
def analyze_doc(self, text: str):
"""Führt einen einzelnen spaCy-Durchlauf über den Text durch."""
return self.nlp(text)
def extract_features_from_doc(self, doc) -> Dict[str, Any]:
"""Extrahiert Lemmata und Entitäten aus einem einzelnen spaCy-Dokument ohne erneutes Parsen."""
lemmas = [
token.lemma_
for token in doc
if not token.is_stop
and not token.is_punct
and not token.is_space
and len(token.lemma_) > 1
]
entities: Dict[str, List[str]] = {}
for ent in doc.ents:
label = ent.label_
if label not in entities:
entities[label] = []
if ent.text not in entities[label]:
entities[label].append(ent.text)
return {"lemmas": lemmas, "entities": entities}
def lemmatize_query(self, text: str) -> List[str]:
doc = self.nlp(text.lower())
return [
token.lemma_
for token in doc
if not token.is_stop
and not token.is_punct
and not token.is_space
and len(token.lemma_) > 1
]
# --- SEMANTIC SPACY CHUNKER ---
class SemanticSpacyChunker:
"""Ein Chunker, der Satzgrenzen respektiert und optimiertes Single-Pass-spaCy verwendet."""
def __init__(
self,
nlp_processor: SpacyNLPProcessor,
max_chunk_size: int = 512,
overlap_sents: int = 1,
):
self.nlp = nlp_processor
self.max_chunk_size = max_chunk_size
self.overlap_sents = overlap_sents
def chunk_document(
self,
text: str,
doc_id: str,
source_file: str,
extra_meta: Optional[Dict[str, Any]] = None,
) -> List[TextChunk]:
full_doc = self.nlp.analyze_doc(text)
sentences = [
sent for sent in full_doc.sents if len(sent.text.strip()) > 0
]
chunks: List[TextChunk] = []
current_sent_objs = []
current_len = 0
chunk_idx = 0
for sent in sentences:
sent_text = sent.text.strip()
sent_len = len(sent_text)
if current_len + sent_len > self.max_chunk_size and current_sent_objs:
chunk_text = " ".join([s.text.strip() for s in current_sent_objs])
chunk_doc = self.nlp.analyze_doc(chunk_text)
features = self.nlp.extract_features_from_doc(chunk_doc)
meta = {"chunk_index": chunk_idx}
if extra_meta:
meta.update(extra_meta)
chunks.append(
TextChunk(
id=f"{doc_id}_c{chunk_idx}",
doc_id=doc_id,
source_file=source_file,
text=chunk_text,
lemmas=features["lemmas"],
entities=features["entities"],
metadata=meta,
)
)
chunk_idx += 1
current_sent_objs = (
current_sent_objs[-self.overlap_sents :]
if self.overlap_sents > 0
else []
)
current_len = sum(len(s.text.strip()) for s in current_sent_objs)
current_sent_objs.append(sent)
current_len += sent_len
if current_sent_objs:
chunk_text = " ".join([s.text.strip() for s in current_sent_objs])
chunk_doc = self.nlp.analyze_doc(chunk_text)
features = self.nlp.extract_features_from_doc(chunk_doc)
meta = {"chunk_index": chunk_idx}
if extra_meta:
meta.update(extra_meta)
chunks.append(
TextChunk(
id=f"{doc_id}_c{chunk_idx}",
doc_id=doc_id,
source_file=source_file,
text=chunk_text,
lemmas=features["lemmas"],
entities=features["entities"],
metadata=meta,
)
)
return chunks
# --- MICRO-LLM QUERY REWRITER (STRICT VLLM ENGINE) ---
class QueryRewriter:
"""Generierung von Suchoptionen über vLLM (OpenAI API-kompatibel)."""
def __init__(
self,
endpoint_url: Optional[str] = None,
model_name: str = "Qwen/Qwen2.5-0.5B-Instruct",
):
self.endpoint_url = endpoint_url
self.model_name = model_name
async def rewrite(self, original_query: str) -> List[str]:
if not self.endpoint_url:
return [original_query]
payload = {
"model": self.model_name,
"messages": [
{
"role": "system",
"content": (
"Du bist das RAG SEO-Modul. "
"Wandeln Sie Ihre Suchanfrage in zwei alternative Schlüsselbegriffe um. "
'Bitte antworten Sie NUR mit einem JSON-Objekt im folgenden Format: {"queries": ["formulation 1", "formulation 2"]}'
),
},
{"role": "user", "content": original_query},
],
"response_format": {"type": "json_object"},
"temperature": 0.3,
"max_tokens": 150,
}
try:
async with aiohttp.ClientSession() as session:
async with session.post(
self.endpoint_url, json=payload, timeout=3.0
) as resp:
if resp.status == 200:
data = await resp.json()
raw_content = data["choices"][0]["message"]["content"]
parsed = json.loads(raw_content)
queries = (
parsed.get("queries", [])
if isinstance(parsed, dict)
else parsed
)
if isinstance(queries, list) and len(queries) > 0:
return list(
dict.fromkeys(
[original_query] + [str(q) for q in queries]
)
)
except Exception:
pass
return [original_query]
# --- DENSE EMBEDDER ENGINE ---
class LocalDenseEmbedder:
"""Textvektorisierung mittels Satztransformatoren (Fallback: Hash-Vektorisierer)."""
def __init__(
self,
model_name: str = "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2",
):
try:
from sentence_transformers import SentenceTransformer
self.model = SentenceTransformer(model_name)
self.dim = self.model.get_embedding_dimension()
self.use_st = True
except (ImportError, OSError, ValueError):
self.use_st = False
self.dim = 384
def embed(self, text: str) -> List[float]:
if self.use_st:
return self.model.encode(text, normalize_embeddings=True).tolist()
import hashlib
words = text.lower().split()
vec = [0.0] * self.dim
for w in words:
idx = int(hashlib.md5(w.encode()).hexdigest(), 16) % self.dim
vec[idx] += 1.0
norm = math.sqrt(sum(v * v for v in vec)) or 1.0
return [v / norm for v in vec]
def embed_batch(self, texts: List[str]) -> List[List[float]]:
if not texts:
return []
if self.use_st:
embeddings = self.model.encode(
texts, normalize_embeddings=True, show_progress_bar=False
)
return [e.tolist() for e in embeddings]
return [self.embed(t) for t in texts]
# --- POSTGRESQL STORAGE BACKEND ---
class StorageBackend:
def __init__(self, connection_str: str, vector_dim: int = 384):
if not POSTGRES_AVAILABLE:
raise ImportError(
"psycopg2-binary is missing. Install via: pip install psycopg2-binary"
)
self.connection_str = connection_str
self.vector_dim = vector_dim
self._init_backend()
def _init_backend(self):
with psycopg2.connect(self.connection_str) as conn:
with conn.cursor() as cur:
cur.execute("CREATE EXTENSION IF NOT EXISTS vector;")
cur.execute(
f"""
CREATE TABLE IF NOT EXISTS rag_chunks (
id VARCHAR(255) PRIMARY KEY,
doc_id VARCHAR(255),
source_file VARCHAR(512),
text TEXT,
lemmas JSONB,
entities JSONB,
metadata JSONB,
vector vector({self.vector_dim})
);
"""
)
# Optimierung der Vektorsuche mittels Kosinusdistanz unter Verwendung des HNSW-Index
cur.execute(
"""
CREATE INDEX IF NOT EXISTS idx_rag_chunks_vector
ON rag_chunks USING hnsw (vector vector_cosine_ops);
"""
)
conn.commit()
def save_chunk(self, chunk: TextChunk):
self.save_chunks_batch([chunk])
def save_chunks_batch(self, chunks: List[TextChunk]):
if not chunks:
return
with psycopg2.connect(self.connection_str) as conn:
with conn.cursor() as cur:
psycopg2.extras.execute_values(
cur,
"""
INSERT INTO rag_chunks (id, doc_id, source_file, text, lemmas, entities, metadata, vector)
VALUES %s
ON CONFLICT (id) DO UPDATE SET
text = EXCLUDED.text, lemmas = EXCLUDED.lemmas,
entities = EXCLUDED.entities, metadata = EXCLUDED.metadata, vector = EXCLUDED.vector;
""",
[
(
c.id,
c.doc_id,
c.source_file,
c.text,
json.dumps(c.lemmas, ensure_ascii=False),
json.dumps(c.entities, ensure_ascii=False),
json.dumps(c.metadata, ensure_ascii=False),
c.vector,
)
for c in chunks
],
)
conn.commit()
def fetch_all(self) -> List[TextChunk]:
chunks = []
with psycopg2.connect(self.connection_str) as conn:
with conn.cursor(cursor_factory=psycopg2.extras.DictCursor) as cur:
cur.execute(
"SELECT id, doc_id, source_file, text, lemmas, entities, metadata, vector FROM rag_chunks;"
)
for r in cur.fetchall():
vec = (
[float(x) for x in r["vector"].strip("[]").split(",")]
if r["vector"]
else None
)
chunks.append(
TextChunk(
id=r["id"],
doc_id=r["doc_id"],
source_file=r["source_file"],
text=r["text"],
lemmas=(
r["lemmas"]
if isinstance(r["lemmas"], list)
else json.loads(r["lemmas"])
),
entities=(
r["entities"]
if isinstance(r["entities"], dict)
else json.loads(r["entities"])
),
metadata=(
r["metadata"]
if isinstance(r["metadata"], dict)
else json.loads(r["metadata"])
),
vector=vec,
)
)
return chunks
def search_dense(
self, query_vec: List[float], top_k: int = 5
) -> List[tuple[str, float]]:
vec_str = f"[{','.join(map(str, query_vec))}]"
with psycopg2.connect(self.connection_str) as conn:
with conn.cursor() as cur:
cur.execute(
"""
SELECT id, 1 - (vector <=> %s::vector) AS similarity
FROM rag_chunks ORDER BY vector <=> %s::vector ASC LIMIT %s;
""",
(vec_str, vec_str, top_k),
)
return [(r[0], float(r[1])) for r in cur.fetchall()]
# --- SEARCH ENGINE (BM25 + DENSE + RRF FUSION) ---
class HybridSearchEngine:
def __init__(
self,
nlp_processor: SpacyNLPProcessor,
embedder: LocalDenseEmbedder,
storage: StorageBackend,
):
self.nlp = nlp_processor
self.embedder = embedder
self.storage = storage
self.chunks_cache: Dict[str, TextChunk] = {}
self.bm25: Optional[BM25Okapi] = None
self.chunk_ids_order: List[str] = []
self.reload_index()
def reload_index(self):
chunks = self.storage.fetch_all()
self.chunks_cache = {c.id: c for c in chunks}
self.chunk_ids_order = list(self.chunks_cache.keys())
corpus_lemmas = [c.lemmas for c in chunks]
if corpus_lemmas:
self.bm25 = BM25Okapi(corpus_lemmas)
def search_sparse_only(
self, query: str, top_k: int = 5
) -> List[Dict[str, Any]]:
if not self.bm25:
return []
q_lemmas = self.nlp.lemmatize_query(query)
scores = self.bm25.get_scores(q_lemmas)
ranked = sorted(
zip(self.chunk_ids_order, scores), key=lambda x: x[1], reverse=True
)[:top_k]
return [
{
"chunk_id": cid,
"bm25_score": round(float(s), 4),
"text": self.chunks_cache[cid].text,
}
for cid, s in ranked
]
def search_dense_only(
self, query: str, top_k: int = 5
) -> List[Dict[str, Any]]:
q_vec = self.embedder.embed(query)
results = self.storage.search_dense(q_vec, top_k=top_k)
return [
{
"chunk_id": cid,
"dense_score": round(float(s), 4),
"text": self.chunks_cache[cid].text,
}
for cid, s in results
]
def search_hybrid_rrf(
self, queries: List[str], top_k: int = 5, rrf_k: int = 60
) -> List[SearchResult]:
if not self.chunks_cache or not self.bm25:
return []
rrf_scores: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
bm25_best: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
dense_best: Dict[str, float] = {cid: 0.0 for cid in self.chunks_cache}
candidate_depth = max(top_k * 5, 50)
for q in queries:
q_lemmas = self.nlp.lemmatize_query(q)
q_vec = self.embedder.embed(q)
# 1. BM25 Sparse Search
bm25_raw = (
self.bm25.get_scores(q_lemmas)
if q_lemmas
else [0.0] * len(self.chunk_ids_order)
)
bm25_ranked = sorted(
zip(self.chunk_ids_order, bm25_raw),
key=lambda x: x[1],
reverse=True,
)[:candidate_depth]
for rank, (cid, score) in enumerate(bm25_ranked):
rrf_scores[cid] += 1.0 / (rrf_k + rank + 1)
bm25_best[cid] = max(bm25_best[cid], float(score))
# 2. PostgreSQL Vector Search
dense_ranked = self.storage.search_dense(
q_vec, top_k=candidate_depth
)
for rank, (cid, score) in enumerate(dense_ranked):
rrf_scores[cid] += 1.0 / (rrf_k + rank + 1)
dense_best[cid] = max(dense_best[cid], float(score))
sorted_rrf = sorted(
rrf_scores.items(), key=lambda x: x[1], reverse=True
)[:top_k]
output = []
for cid, score in sorted_rrf:
if score == 0.0:
continue
c = self.chunks_cache[cid]
output.append(
SearchResult(
chunk_id=c.id,
doc_id=c.doc_id,
source_file=c.source_file,
text=c.text,
rrf_score=round(score, 5),
bm25_score=round(bm25_best[cid], 4),
dense_score=round(dense_best[cid], 4),
metadata={**c.metadata, "entities": c.entities},
)
)
return output
# --- FACADE SYSTEM MANAGER & DIRECTORY SCANNER ---
class RAGEngine:
def __init__(
self,
db_conn: str,
spacy_model: str = "de_core_news_lg",
vllm_url: Optional[str] = None,
max_workers: int = 4,
):
self.nlp = SpacyNLPProcessor(model_name=spacy_model)
self.chunker = SemanticSpacyChunker(self.nlp)
self.embedder = LocalDenseEmbedder()
self.rewriter = QueryRewriter(endpoint_url=vllm_url)
self.storage = StorageBackend(
connection_str=db_conn, vector_dim=self.embedder.dim
)
self.search_engine = HybridSearchEngine(
self.nlp, self.embedder, self.storage
)
self.max_workers = max_workers
self.pool = ThreadPoolExecutor(max_workers=max_workers)
def ingest_file(
self, file_path: str, max_chunk_size: int = 512
) -> Dict[str, Any]:
"""Dokumentenanalyse mittels DocumentReader und Vektorisierung in PostgreSQL."""
doc_data = DocumentReader.read_file(file_path)
content = doc_data["content"]
file_meta = doc_data["metadata"]
if not content.strip():
return {
"file": file_path,
"chunks_created": 0,
"status": "skipped_empty",
}
self.chunker.max_chunk_size = max_chunk_size
doc_id = file_meta["file_name"]
chunks = self.chunker.chunk_document(
content, doc_id=doc_id, source_file=file_path, extra_meta=file_meta
)
texts = [c.text for c in chunks]
vectors = self.embedder.embed_batch(texts)
for chunk, vec in zip(chunks, vectors):
chunk.vector = vec
self.storage.save_chunks_batch(chunks)
return {
"file": file_path,
"chunks_created": len(chunks),
"storage_backend": "POSTGRESQL",
"metadata": file_meta,
}
async def ingest_file_async(
self, file_path: str, max_chunk_size: int = 512
) -> Dict[str, Any]:
loop = asyncio.get_running_loop()
return await loop.run_in_executor(
self.pool, self.ingest_file, file_path, max_chunk_size
)
def ingest_directory(
self,
dir_path: str,
extensions: Optional[Set[str]] = None,
max_chunk_size: int = 512,
) -> Dict[str, Any]:
if extensions is None:
extensions = {
".txt",
".md",
".json",
".log",
".py",
".csv",
".pdf",
".docx",
".html",
".htm",
}
root = Path(dir_path)
if not root.exists() or not root.is_dir():
raise NotADirectoryError(
f"Das Verzeichnis {dir_path} existiert nicht oder ist kein Ordner."
)
files_to_process = [
str(path)
for path in root.rglob("*")
if path.is_file() and path.suffix.lower() in extensions
]
file_summary = []
processed_files = 0
total_chunks = 0
futures = {
self.pool.submit(self.ingest_file, fp, max_chunk_size): fp
for fp in files_to_process
}
for future in futures:
try:
res = future.result()
processed_files += 1
total_chunks += res.get("chunks_created", 0)
file_summary.append(res)
except Exception as e:
fp = futures[future]
file_summary.append(
{"file": fp, "error": str(e), "chunks_created": 0}
)
self.search_engine.reload_index()
return {
"directory": str(root.absolute()),
"processed_files_count": processed_files,
"total_chunks_created": total_chunks,
"total_chunks_in_db": len(self.search_engine.chunks_cache),
"storage_backend": "POSTGRESQL",
"files": file_summary,
}
async def ingest_directory_async(
self,
dir_path: str,
extensions: Optional[Set[str]] = None,
max_chunk_size: int = 512,
) -> Dict[str, Any]:
if extensions is None:
extensions = {
".txt",
".md",
".json",
".log",
".py",
".csv",
".pdf",
".docx",
".html",
".htm",
}
root = Path(dir_path)
if not root.exists() or not root.is_dir():
raise NotADirectoryError(
f"Das Verzeichnis {dir_path} existiert nicht oder ist kein Ordner."
)
files_to_process = [
str(path)
for path in root.rglob("*")
if path.is_file() and path.suffix.lower() in extensions
]
loop = asyncio.get_running_loop()
tasks = [
loop.run_in_executor(
self.pool, self.ingest_file, fp, max_chunk_size
)
for fp in files_to_process
]
results = await asyncio.gather(*tasks, return_exceptions=True)
file_summary = []
processed_files = 0
total_chunks = 0
for fp, res in zip(files_to_process, results):
if isinstance(res, BaseException):
file_summary.append(
{"file": fp, "error": str(res), "chunks_created": 0}
)
else:
res_dict = cast(Dict[str, Any], res)
processed_files += 1
total_chunks += res_dict.get("chunks_created", 0)
file_summary.append(res_dict)
self.search_engine.reload_index()
return {
"directory": str(root.absolute()),
"processed_files_count": processed_files,
"total_chunks_created": total_chunks,
"total_chunks_in_db": len(self.search_engine.chunks_cache),
"storage_backend": "POSTGRESQL",
"files": file_summary,
}
# --- CLI INTERFACE ---
def main():
parser = argparse.ArgumentParser(
description="RAGEngine PostgreSQL Architecture"
)
parser.add_argument(
"--db-conn",
type=str,
default=os.getenv(
"POSTGRES_CONN_STR",
"postgresql://postgres:postgres@localhost:5432/rag_db",
),
help="PostgreSQL Connection String",
)
parser.add_argument(
"--spacy-model", type=str, default="de_core_news_lg"
)
parser.add_argument(
"--vllm-url",
type=str,
default=None,
help="vLLM Chat Completions URL (e.g. http://localhost:8000/v1/chat/completions)",
)
parser.add_argument(
"--max-workers",
type=int,
default=4,
help="ThreadPoolExecutor worker threads count",
)
subparsers = parser.add_subparsers(dest="command", required=True)
# NLP Commands
chunk_p = subparsers.add_parser(
"chunk", help="[Chunker] Test semantic chunking"
)
chunk_p.add_argument("--text", type=str, required=True)
chunk_p.add_argument("--max-size", type=int, default=200)
rw_p = subparsers.add_parser(
"rewrite", help="[vLLM] Test query expansion"
)
rw_p.add_argument("--query", type=str, required=True)
# Search Commands
ss_p = subparsers.add_parser(
"search-sparse", help="[Search] BM25-only search"
)
ss_p.add_argument("--query", type=str, required=True)
ss_p.add_argument("--top-k", type=int, default=3)
sd_p = subparsers.add_parser(
"search-dense", help="[Search] Dense vector search in PostgreSQL"
)
sd_p.add_argument("--query", type=str, required=True)
sd_p.add_argument("--top-k", type=int, default=3)
sh_p = subparsers.add_parser(
"search-hybrid",
help="[Facade Search] Full RRF Hybrid Search (BM25 + PG Vector)",
)
sh_p.add_argument("--query", type=str, required=True)
sh_p.add_argument("--top-k", type=int, default=3)
# Ingestion Commands
ing_f = subparsers.add_parser(
"ingest-file", help="[Storage] Index single file into PostgreSQL"
)
ing_f.add_argument("--file", type=str, required=True)
ing_d = subparsers.add_parser(
"ingest-dir",
help="[Storage] Recursively scan and index directory into PostgreSQL",
)
ing_d.add_argument("--dir", type=str, required=True)
ing_d.add_argument(
"--exts",
type=str,
default=".txt,.md,.json,.log,.py,.csv,.pdf,.docx,.html,.htm",
help="Comma-separated file extensions",
)
subparsers.add_parser("mcp", help="[MCP] Launch FastMCP server")
args = parser.parse_args()
engine = RAGEngine(
db_conn=args.db_conn,
spacy_model=args.spacy_model,
vllm_url=args.vllm_url,
max_workers=args.max_workers,
)
if args.command == "chunk":
engine.chunker.max_chunk_size = args.max_size
chunks = engine.chunker.chunk_document(
args.text, doc_id="cli_test", source_file="stdin"
)
print(
json.dumps(
[c.to_dict() for c in chunks], ensure_ascii=False, indent=2
)
)
elif args.command == "rewrite":
rewritten = asyncio.run(engine.rewriter.rewrite(args.query))
print(
json.dumps(
{"expanded_queries": rewritten}, ensure_ascii=False, indent=2
)
)
elif args.command == "search-sparse":
res = engine.search_engine.search_sparse_only(
args.query, top_k=args.top_k
)
print(
json.dumps(
{"type": "BM25-Sparse", "results": res},
ensure_ascii=False,
indent=2,
)
)
elif args.command == "search-dense":
res = engine.search_engine.search_dense_only(
args.query, top_k=args.top_k
)
print(
json.dumps(
{"type": "Dense-Vector (POSTGRESQL)", "results": res},
ensure_ascii=False,
indent=2,
)
)
elif args.command == "search-hybrid":
queries = asyncio.run(engine.rewriter.rewrite(args.query))
results = engine.search_engine.search_hybrid_rrf(
queries, top_k=args.top_k
)
print(
json.dumps(
{
"expanded_queries": queries,
"backend": "POSTGRESQL",
"results": [asdict(cast(Any, r)) for r in results],
},
ensure_ascii=False,
indent=2,
)
)
elif args.command == "ingest-file":
info = asyncio.run(engine.ingest_file_async(args.file))
engine.search_engine.reload_index()
print(
json.dumps(
{"status": "success", "info": info},
ensure_ascii=False,
indent=2,
)
)
elif args.command == "ingest-dir":
x_set = {f".{ext.strip().lstrip('.')}" for ext in args.exts.split(",")}
info = asyncio.run(
engine.ingest_directory_async(args.dir, extensions=x_set)
)
print(
json.dumps(
{"status": "success", "info": info},
ensure_ascii=False,
indent=2,
)
)
elif args.command == "mcp":
try:
from mcp.server.fastmcp import FastMCP
except ImportError:
try:
from fastmcp import FastMCP
except ImportError:
print("FastMCP is missing. Install via: pip install fastmcp")
return
mcp = FastMCP("RAGEngine PostgreSQL Service")
@mcp.tool()
async def query_rag(query: str) -> str:
q_list = await engine.rewriter.rewrite(query)
res_ = engine.search_engine.search_hybrid_rrf(q_list, top_k=3)
if not res_:
return "Совпадений не найдено."
return "\n\n".join(
[
f"[{r.source_file} | RRF: {r.rrf_score}]\n{r.text}"
for r in res_
]
)
mcp.run()
if __name__ == "__main__":
main()
Erstellen Sie eine Datei wrap.py
import logging
import os
from typing import Any, Dict, List, Optional
# Importieren Sie die Original-Engine aus rag.py
import rag as core_rag
logger = logging.getLogger(__name__)
# Zwischenspeicherung von RAGEngine-Instanzen anhand der Verbindungszeichenfolge, um unnötige Initialisierung zu vermeiden
_engines_cache: Dict[str, core_rag.RAGEngine] = {}
def _get_engine(db_conn: str) -> core_rag.RAGEngine:
if db_conn not in _engines_cache:
_engines_cache[db_conn] = core_rag.RAGEngine(db_conn=db_conn)
return _engines_cache[db_conn]
def search(
query: str,
db_type: str,
db_conn: str,
top_k: int = 5,
filters: Optional[Dict[str, Any]] = None,
) -> List[Dict[str, Any]]:
"""Dokumentensuche mit Hilfe der hybriden RAGEngine-Engine (BM25 + Dense RRF in PostgreSQL)."""
try:
engine = _get_engine(db_conn)
results = engine.search_engine.search_hybrid_rrf([query], top_k=top_k)
formatted_docs = []
for r in results:
# Anwenden optionaler Metadatenfilter
if filters:
skip = False
for k, v in filters.items():
if r.metadata.get(k) != v:
skip = True
break
if skip:
continue
formatted_docs.append(
{
"id": r.chunk_id,
"doc_id": r.doc_id,
"text": r.text,
"score": r.rrf_score,
"source": r.source_file,
"bm25_score": r.bm25_score,
"dense_score": r.dense_score,
"metadata": r.metadata,
}
)
return formatted_docs
except Exception as e:
logger.error(f"Error in wrap.search: {e}", exc_info=True)
return []
def rerank(
query: str, documents: List[Dict[str, Any]], threshold: float = 0.7
) -> List[Dict[str, Any]]:
"""Gefundene Dokumente werden anhand eines Relevanzschwellenwerts (Score) gefiltert."""
return [doc for doc in documents if doc.get("score", 0.0) >= threshold]
def chunk_and_trim(
documents: List[Dict[str, Any]], max_tokens: int = 2000
) -> List[Dict[str, Any]]:
"""Die Liste der Chunks wird so gekürzt, dass sie dem Kontext-Token-Limit entspricht."""
trimmed = []
current_tokens = 0
for doc in documents:
text = doc.get("text", "")
estimated_tokens = len(text.split()) * 1.3
if current_tokens + estimated_tokens <= max_tokens:
trimmed.append(doc)
current_tokens += estimated_tokens
else:
break
return trimmed
def format_context(documents: List[Dict[str, Any]]) -> str:
"""Formatieren einer Liste von Dokumenten in eine einzelne Kontexttextzeichenfolge."""
if not documents:
return "Kontext nicht gefunden."
blocks = []
for doc in documents:
source = doc.get("source", "Unknown")
score = doc.get("score", "N/A")
text = doc.get("text", "")
blocks.append(f"--- Quelle: {source} (Score: {score}) ---\n{text}")
return "\n\n".join(blocks)
def build_system_prompt(
context: str, custom_instructions: Optional[str] = None
) -> str:
"""Erstellung einer Systemabfrage mit eingebettetem Wissensdatenbankkontext."""
base_instructions = custom_instructions or (
"Sie sind eine erfahrene Unternehmensassistentin. Beantworten Sie Benutzerfragen, "
"ausschließlich basierend auf dem bereitgestellten Wissensbasiskontext. "
"Wenn der Kontext keine Antwort liefert, machen Sie dies deutlich."
)
return f"{base_instructions}\n\n=== WISSENSBASISKONTEXT ===\n{context}"
def extract_sources(documents: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
"""Extrahieren von Metadaten über zitierte Quellen."""
return [
{
"id": doc.get("id"),
"source": doc.get("source"),
"score": doc.get("score"),
}
for doc in documents
]
def get_mcp_context(query: str, options: Dict[str, Any]) -> str:
"""Extraktion dynamischer Kontextinformationen aus dem Model Context Protocol (MCP)."""
return f"MCP Dynamic Context for query '{query}': PostgreSQL RAGEngine storage is active and operational."
def get_mcp_tools(options: Dict[str, Any]) -> List[Dict[str, Any]]:
"""Abrufen von Beschreibungen der MCP-Funktionen/Tools im OpenAI Tools-Format."""
return [
{
"type": "function",
"function": {
"name": "query_system_metrics",
"description": "Erhalten Sie aktuelle Kennzahlen zum Zustand des RAG-Systems.",
"parameters": {
"type": "object",
"properties": {
"metric_name": {
"type": "string",
"description": "Metrikname (Beispiel, total_chunks, db_status)",
}
},
"required": ["metric_name"],
},
},
}
]
def execute_mcp_tool(
name: str, arguments: Dict[str, Any], options: Dict[str, Any]
) -> Dict[str, Any]:
"""Führe einen MCP-Funktionsaufruf aus und gib das Ergebnis zurück."""
if name == "query_system_metrics":
metric = arguments.get("metric_name", "db_status")
db_conn = options.get("db_conn", "postgresql://postgres:postgres@localhost:5432/rag_db")
try:
engine = _get_engine(db_conn)
total_chunks = len(engine.search_engine.chunks_cache)
except Exception:
total_chunks = "unknown"
return {
"metric": metric,
"total_chunks_in_cache": total_chunks,
"status": "operational",
}
return {"error": f"Tool '{name}' not found in RAGEngine MCP registry."}
Erstellen Sie eine Datei app.py
import json
import logging
import os
import requests
from flask import Flask, Response, jsonify, request, stream_with_context
from openai import OpenAI
import wrap as rag
app = Flask(__name__)
logging.basicConfig(level=logging.INFO)
OPENAI_API_KEY = os.getenv("OPENAI_API_KEY")
if not OPENAI_API_KEY:
raise ValueError("OPENAI_API_KEY environment variable is required")
client = OpenAI(api_key=OPENAI_API_KEY)
DEFAULT_DB_TYPE = os.getenv("DB_TYPE")
DEFAULT_DB_CONN = os.getenv("DB_CONN")
DEFAULT_LLM_URL = os.getenv("LLM_URL")
MAX_TOOL_ITERATIONS = 5
def extract_user_query(messages: list[dict]) -> str:
for msg in reversed(messages):
if msg.get("role") == "user":
content = msg.get("content", "")
if isinstance(content, str):
return content
elif isinstance(content, list):
return " ".join([
part.get("text", "") for part in content
if isinstance(part, dict) and part.get("type") == "text"
])
return ""
@app.route("/health", methods=["GET"])
def health_check():
return jsonify({"status": "ok"}), 200
@app.route("/v1/chat/completions", methods=["POST"])
def chat_completions():
payload = request.get_json() or {}
messages = payload.get("messages", [])
model = payload.get("model", "gpt-4o-mini")
temperature = payload.get("temperature", 0.2)
stream = payload.get("stream", False)
rag_opts = payload.get("rag_options", {})
db_type = rag_opts.get("db_type", DEFAULT_DB_TYPE)
db_conn = rag_opts.get("db_conn", DEFAULT_DB_CONN)
llm_url = rag_opts.get("llm_url", DEFAULT_LLM_URL)
top_k = rag_opts.get("top_k", 5)
score_threshold = rag_opts.get("score_threshold", 0.7)
max_context_tokens = rag_opts.get("max_context_tokens", 2000)
filters = rag_opts.get("filters", {})
system_instructions = rag_opts.get("system_instructions", None)
enable_query_rewrite = rag_opts.get("rewrite_query", False)
rewrite_model = rag_opts.get("rewrite_model", "qwen")
use_mcp = rag_opts.get("use_mcp", False)
raw_user_query = extract_user_query(messages)
if not raw_user_query:
return jsonify({
"error": {
"message": "User query is missing in 'messages'",
"type": "invalid_request_error",
"code": 400
}
}), 400
try:
# 1. Query Rewrite
if enable_query_rewrite:
endpoint = f"{llm_url.rstrip('/')}/chat/completions"
rewrite_payload = {
"model": rewrite_model,
"temperature": 0.0,
"messages": [
{
"role": "system",
"content": "Rephrase the input query for optimal vector search. Return ONLY the rewritten query text."
},
{"role": "user", "content": raw_user_query}
]
}
resp = requests.post(endpoint, json=rewrite_payload, timeout=15)
resp.raise_for_status()
search_query = resp.json()["choices"][0]["message"]["content"].strip()
else:
search_query = raw_user_query
# 2. RAG Pipeline
raw_docs = rag.search(
query=search_query,
db_type=db_type,
db_conn=db_conn,
top_k=top_k,
filters=filters
)
reranked_docs = rag.rerank(
query=search_query,
documents=raw_docs,
threshold=score_threshold
)
trimmed_chunks = rag.chunk_and_trim(
documents=reranked_docs,
max_tokens=max_context_tokens
)
context_text = rag.format_context(trimmed_chunks)
system_prompt = rag.build_system_prompt(
context=context_text,
custom_instructions=system_instructions
)
sources_metadata = rag.extract_sources(trimmed_chunks)
# 3. MCP Integration
openai_tools = None
mcp_executed_tools = []
if use_mcp:
mcp_context = rag.get_mcp_context(query=search_query, options=rag_opts)
if mcp_context:
system_prompt = f"{system_prompt}\n\n--- Additional MCP Context ---\n{mcp_context}"
openai_tools = rag.get_mcp_tools(options=rag_opts)
# 4. OpenAI Completion Loop
conversation_history = [m for m in messages if m.get("role") != "system"]
final_messages = [{"role": "system", "content": system_prompt}] + conversation_history
completion_kwargs = {
"model": model,
"messages": final_messages,
"temperature": temperature
}
if openai_tools:
completion_kwargs["tools"] = openai_tools
iterations = 0
completion = None
while iterations < MAX_TOOL_ITERATIONS:
if stream and not (openai_tools and iterations == 0):
break
completion = client.chat.completions.create(**completion_kwargs)
response_message = completion.choices[0].message
if response_message.tool_calls:
final_messages.append(response_message.model_dump(exclude_none=True))
for tool_call in response_message.tool_calls:
tool_name = tool_call.function.name
tool_args = json.loads(tool_call.function.arguments or "{}")
tool_output = rag.execute_mcp_tool(
name=tool_name,
arguments=tool_args,
options=rag_opts
)
mcp_executed_tools.append({
"name": tool_name,
"args": tool_args,
"output": tool_output
})
final_messages.append({
"role": "tool",
"tool_call_id": tool_call.id,
"content": json.dumps(tool_output, ensure_ascii=False)
})
completion_kwargs["messages"] = final_messages
iterations += 1
else:
break
# 5. Output Response
if stream:
completion_kwargs["stream"] = True
response_stream = client.chat.completions.create(**completion_kwargs)
def generate_sse():
for chunk in response_stream:
yield f"data: {chunk.model_dump_json()}\n\n"
yield "data: [DONE]\n\n"
return Response(
stream_with_context(generate_sse()),
mimetype="text/event-stream"
)
if not completion:
raise RuntimeError("Failed to receive completion from OpenAI API.")
response_dict = completion.model_dump()
response_dict["rag_metadata"] = {
"db_type": db_type,
"query_rewritten": enable_query_rewrite,
"llm_url_used": llm_url if enable_query_rewrite else None,
"original_query": raw_user_query,
"search_query_used": search_query,
"retrieved_chunks_count": len(trimmed_chunks),
"sources": sources_metadata,
"mcp": {
"enabled": use_mcp,
"executed_tools": mcp_executed_tools
}
}
return jsonify(response_dict), 200
except Exception as e:
app.logger.error(f"Execution Error: {str(e)}", exc_info=True)
return jsonify({
"error": {
"message": f"Internal RAG Execution Error: {str(e)}",
"type": "api_error",
"code": 500
}
}), 500
Erstellen Sie eine Datei /etc/systemd/system/rag-api.service
[Unit]
Description=Flask RAG & MCP Proxy API Service
After=network.target
[Service]
User=www-data
Group=www-data
WorkingDirectory=/var/www/rag-api
RuntimeDirectory=rag-api
Environment="PATH=/var/www/rag-api/.venv/bin"
Environment="OPENAI_API_KEY=sk-proj-your-actual-api-key-here"
Environment="DB_TYPE=postgres"
Environment="DB_CONN=postgresql://postgres:pass@10.0.0.7:5432/rag"
Environment="LLM_URL=http://10.0.0.5/v1"
ExecStart=/var/www/rag-api/venv/bin/gunicorn \
--workers 4 \
--threads 2 \
--worker-class gthread \
--bind unix:/run/rag-api/rag-api.sock \
--access-logfile /var/log/rag-api/access.log \
--error-logfile /var/log/rag-api/error.log \
--umask 007 \
--timeout 120 \
app:app
Restart=always
RestartSec=5
StandardOutput=journal
StandardError=journal
[Install]
WantedBy=multi-user.target
Erstellen Sie eine Datei /etc/nginx/sites-available/rag-api
server {
listen 80;
server_name rag.domain.de;
client_max_body_size 20M;
location / {
proxy_pass http://unix:/run/rag-api/rag-api.sock;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
# Unterstützung SSE (Server-Sent Events / Streaming)
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_buffering off;
proxy_cache off;
chunked_transfer_encoding on;
# Timeouts für langfristige LLM- und RAG-Generationen
proxy_read_timeout 180s;
proxy_connect_timeout 60s;
proxy_send_timeout 180s;
}
}
Erstellen Sie Protokolldateien, starten Sie den systemd-Dienst neu und überprüfen Sie, ob die Socket-Datei erstellt wurde:
mkdir -p /var/log/rag-api
chown -R www-data:www-data /var/log/rag-api
systemctl daemon-reload
systemctl restart rag-api
systemctl status rag-api
ls -la /run/rag-api/rag-api.sock
Wir erheben den vhost
ln -s /etc/nginx/sites-available/rag-api /etc/nginx/sites-enabled/
nginx -t
systemctl reload nginx
Verbindung schützen
apt install certbot python3-certbot-nginx -y
certbot --nginx
crontab -e
0 3 * * * /usr/bin/certbot renew --quiet --deploy-hook "systemctl reload nginx"
certbot --nginx --redirect
nginx -t
/etc/init.d/nginx restart
Erstellen Sie eine Datei: /etc/logrotate.d/rag-api
/var/log/rag-api/*.log {
daily
rotate 14
missingok
notifempty
compress
delaycompress
create 0640 www-data www-data
sharedscripts
postrotate
# Senden eines USR1-Signals an Gunicorn, um Protokolldateien wieder zu öffnen, ohne Prozesse neu zu starten
/bin/systemctl kill -s USR1 rag-api.service > /dev/null 2>&1 || true
endscript
}
Wir lassen die root zurück
exit
4. Vervollständigung der Wissensdatenbank CLI
Grundlegende Indizierung einer Textdatei
python rag.py ingest-file --file /home/bogatyrev/perelman_poincare.txt
Indizieren einer Datei mithilfe einer benutzerdefinierten PostgreSQL-Verbindungszeichenfolge
python rag.py --db-conn "postgresql://rag_user:xxxxxxxxxxxxx@10.0.0.7:5432/rag_db" ingest-file --file /home/bogatyrev/postgres_tuning.md
Indizierung eines Markdown-Dokuments mit dem alternativen Sprachmodell spaCy (NLP)
python rag.py --spacy-model de_core_news_lg ingest-file --file /home/bogatyrev/kubectl_k3s.md
Indizierung des Quellcodes (Python-Skripte) in die Wissensbasis
python rag.py ingest-file --file /home/bogatyrev/anthologie.py
Indizieren der JSON-Konfigurationsdatei
python rag.py --db-conn "postgresql://rag_user:xxxxxxxxxxxxx@10.0.0.7:5432/rag_db" ingest-file --file config/cluster_topology.json
Indizieren eines PDF-Dokuments
python rag.py ingest-file --file /home/bogatyrev/reports/security_audit_2026.pdf
Indizierung eines unternehmensweiten MS Word-Dokuments (.docx)
python rag.py ingest-file --file /home/bogatyrev/sla_agreement.docx
Rekursives Scannen und Indizieren des gesamten Verzeichnisses
python rag.py ingest-dir --dir /home/bogatyrev/knowledge_base/
Rekursive Indizierung mit Filterung nach bestimmten Dateierweiterungen
python rag.py ingest-dir --dir /home/bogatyrev/projects/docs --exts ".md,.txt"
Rekursive Indizierung mit Filterung nach bestimmten Dateierweiterungen
python rag.py --max-workers 8 --vllm-url "https://llm.domain.de/v1/chat/completions" ingest-dir --dir ~/enterprise_docs/ --exts ".pdf,.docx,.md"
Liebe Leserinnen und Leser, ich freue mich sehr über Ihr Interesse an meinem Blog.
Lieber René aus Belgien (Ihrer IP-Adresse nach zu urteilen), ich beantworte Ihre Frage (psql: error: FATAL: password authentication failed for user "rene"): Ich gehe davon aus, dass Sie PostgreSQL nicht aus dem PostgreSQL-Repository, sondern aus dem Debian11-Repository installiert haben, wo aktuell Version 13 verfügbar ist. Diese Version verwendet einen anderen Passwortverschlüsselungsalgorithmus (scram-sha-256):
So lösen Sie das Problem
# Stellen wir den Benutzer und die Authentifizierungsmethode auf die altbewährte Methode md5 um.
sudo vi /etc/postgresql/13/main/pg_hba.conf
<ESC>i
host rag_db rene 127.0.0.1/32 md5
host rag_db rene ::1/128 md5
<ESC>:wq<ENTER>
sudo -u postgres psql
SET password_encryption = 'md5';
ALTER USER rene WITH ENCRYPTED PASSWORD 'xxxxxxxxxxxxxx';
\q
sudo systemctl restart postgresql
psql -h 127.0.0.1 -U rene -d rag_db
# Was könnte sonst noch passieren, wenn Sie scram-sha-256 weiterhin beibehalten möchten?
sudo vi /etc/postgresql/13/main/postgresql.conf
<ESC>i
password_encryption = scram-sha-256
# Falls die Einstellung abweicht oder auskommentiert ist, ändern Sie sie.
<ESC>:wq<ENTER>
# Falls dies der Fall ist, wechseln Sie zu MD5.
<ESC>:q!<ENTER>
sudo -u postgres psql
ALTER USER rene WITH ENCRYPTED PASSWORD 'xxxxxxxxxxxxxx';
\q
sudo systemctl restart postgresql
psql -h 127.0.0.1 -U rene -d rag_db
J'espère vous avoir aidé.