This commit is contained in:
2026-07-10 10:34:50 -05:00
parent 158d2a16e5
commit 5650b7b3fa
5 changed files with 498 additions and 87 deletions

View File

@@ -0,0 +1,62 @@
# llama.cpp + Gemma-4-12B (Vision) - VULKAN-Backend statt ROCm.
#
# Grund fuer den Wechsel: Gemma-4-Vision unter ROCm geraet in eine
# Endlosschleife (<unused>-Flood, llama.cpp Issue #21416) - ein ROCm-Backend-Bug.
# Der Report bestaetigt: "With vulkan it works on the same hardware."
# Gleiche Karte (R9700), gleiches Modell, gleicher Client - nur Backend anders.
#
# Wichtige Unterschiede zu ROCm:
# - KEIN /dev/kfd noetig (Vulkan nutzt den Grafik-Stack, nicht ROCm-Compute)
# - Gruppe 'render' zusaetzlich (Zugriff auf /dev/dri/renderD*)
# - --device Vulkan0 explizit, sonst evtl. CPU-Fallback (llvmpipe)
#
# Laeuft ALLEIN auf der GPU (Qwen vorher stoppen).
#
# Start: docker compose -f docker-compose-llama-gemma12b.yml up -d
# Logs: docker compose -f docker-compose-llama-gemma12b.yml logs -f
services:
llamacpp-gemma12b:
image: ghcr.io/ggml-org/llama.cpp:server-vulkan
container_name: llamacpp-gemma12b
restart: unless-stopped
init: true
devices:
- /dev/dri:/dev/dri
group_add:
- video
- "991" # GID der render-Gruppe (getent group render -> 991)
security_opt:
- seccomp=unconfined
ipc: host
environment:
# Nur die erste (einzige) GPU sichtbar machen; CPU-Fallback vermeiden
- GGML_VK_VISIBLE_DEVICES=0
ports:
- "8000:8080"
volumes:
- ~/.cache/llama.cpp:/root/.cache/llama.cpp
command:
- -hf
- unsloth/gemma-4-12b-it-GGUF:UD-Q4_K_XL
- --host
- 0.0.0.0
- --port
- "8080"
- --device
- Vulkan0
- -ngl
- "99"
- --ctx-size
- "8192"
- --parallel
- "1"
- --jinja
- --reasoning
- "off"
- --cache-ram
- "0"
- --ctx-checkpoints
- "0"
- --alias
- gemma-4-12b

View File

@@ -1,6 +1,13 @@
# llama.cpp + Gemma-4-12B (Vision) fuer die Buyer-Sheet-Extraktion.
# Basiert auf dem bewaehrten Qwen-Run (gleiche ROCm-Flags), erweitert um
# Vision (mmproj laedt -hf automatisch) und --reasoning off (sonst leeres content).
# Laeuft ALLEIN auf der GPU (Qwen vorher stoppen: docker stop llamacpp-qwen36).
#
# Wichtige Stabilitaets-Parameter gegen das kumulative Verklemmen:
# --cache-ram 0 Prompt-Cache AUS. Die Vision-Requests teilen keinen
# sinnvollen Prefix; der Cache brachte nur Thrashing
# ("making room ... MiB") bis zum Slot-Deadlock.
# Ohne Cache ist jeder Request wirklich unabhaengig.
# --ctx-size 8192 reicht fuer ein Sheet + JSON-Antwort; halber KV-Speicher.
# --no-context-shift sauberes Abschneiden statt fragwuerdigem Shift.
#
# Start: docker compose -f docker-compose-llama-gemma12b.yml up -d
# Logs: docker compose -f docker-compose-llama-gemma12b.yml logs -f
@@ -11,7 +18,6 @@ services:
image: ghcr.io/ggml-org/llama.cpp:server-rocm
container_name: llamacpp-gemma12b
restart: unless-stopped
# init:true -> haengende Prozesse werden ordentlich eingesammelt (kein Zombie)
init: true
devices:
- /dev/kfd:/dev/kfd
@@ -22,8 +28,6 @@ services:
- seccomp=unconfined
ipc: host
ports:
# Achtung: Qwen laeuft schon auf 8080. Gemma bekommt 8000, damit beide
# parallel laufen koennen. Client entsprechend auf :8000 zeigen.
- "8000:8080"
volumes:
- ~/.cache/llama.cpp:/root/.cache/llama.cpp
@@ -37,11 +41,14 @@ services:
- -ngl
- "99"
- --ctx-size
- "16384"
- "8192"
- --parallel
- "1"
- --jinja
- --reasoning
- "off"
- --cache-ram
- "0"
- --no-context-shift
- --alias
- gemma-4-12b
- gemma-4-12b

View File

@@ -67,7 +67,12 @@ JSON_SCHEMA = {
"background_experience": {"type": ["string", "null"]},
"date_of_introduction": {"type": ["string", "null"]},
},
"required": ["is_buyer_sheet", "types_of_business"],
"required": [
"is_buyer_sheet", "name_company", "prospective_buyer", "company",
"phone", "cell", "email", "address", "state", "how_did_you_hear",
"interested_in_updates", "types_of_business_raw", "types_of_business",
"background_experience", "date_of_introduction"
],
}
SYSTEM_PROMPT = """Du bist ein praezises Datenextraktions-System fuer Formulare der Firma "BizMatch Business Brokerage".
@@ -169,12 +174,12 @@ def extract_text(pdf_path, max_pages=4):
return ""
def pdf_to_base64_images(pdf_path, max_pages=2, max_dim=2200, quality=85, dpi=200):
def pdf_to_base64_images(pdf_path, max_pages=2, max_dim=1600, quality=80, dpi=150):
"""
Erste max_pages Seiten -> Base64-JPEG.
Hoehere Aufloesung als zuvor (dpi=200, max_dim=2200), damit die kleinen
YES/NO-Kaestchen erkennbar bleiben. max_pages=2 reicht: Buyer-Felder
stehen auf den ersten Seiten, und es haelt die Vision-Token im Rahmen.
Moderate Aufloesung (dpi=150, max_dim=1600): reicht fuer Checkbox/Handschrift
(im Test bei dpi=150 erkannt), haelt aber die Vision-Payload klein genug,
dass der Prompt-Cache nicht ueberlaeuft (die "making room ... MiB"-Warnungen).
"""
try:
images = convert_from_path(pdf_path, first_page=1, last_page=max_pages, dpi=dpi)
@@ -199,7 +204,10 @@ def pdf_to_base64_images(pdf_path, max_pages=2, max_dim=2200, quality=85, dpi=20
def call_model(client, model, content_payload, vision_budget=None):
# llama.cpp erzwingt JSON ueber response_format mit json_schema.
# (vLLM-spezifisches guided_json / mm_processor_kwargs gibt es hier nicht.)
# WICHTIG gegen Gemma-4 "<unused49>-Flood"/Runaway:
# - enable_thinking:false per Request (zuverlaessiger als nur Server-Flag)
# - kleines max_tokens (ein Datensatz braucht ~300-400 Tokens)
# - repeat_penalty gegen Wiederholschleifen
resp = client.chat.completions.create(
model=model,
messages=[
@@ -207,13 +215,22 @@ def call_model(client, model, content_payload, vision_budget=None):
{"role": "user", "content": content_payload},
],
temperature=0.0,
max_tokens=1200,
max_tokens=512,
response_format={
"type": "json_schema",
"json_schema": {"name": "buyer_sheet", "schema": JSON_SCHEMA},
},
extra_body={
"chat_template_kwargs": {"enable_thinking": False},
"repeat_penalty": 1.05,
},
)
return resp.choices[0].message.content
choice = resp.choices[0]
content = choice.message.content
# finish_reason "length" = Runaway (max_tokens erreicht) -> als Warnung markieren
if getattr(choice, "finish_reason", None) == "length":
return content, "length"
return content, "stop"
def parse_json_loose(txt):
@@ -234,45 +251,111 @@ def parse_json_loose(txt):
TEXT_SCAN_THRESHOLD = 120 # < so viele Textzeichen -> als Scan behandeln
MAX_SCAN_PAGES_TOTAL = 12 # PDFs mit mehr Seiten gelten als "kein reines Sheet"
# Felder, bei denen wir Vision bevorzugen (handschriftlich/visuell auf dem Formular)
_PREFER_VISION = {
"prospective_buyer", "name_company", "company", "phone", "cell", "email",
"address", "state", "how_did_you_hear", "interested_in_updates",
"date_of_introduction",
}
# Felder, bei denen Text meist die sauberere (getippte) Quelle ist
_PREFER_TEXT = {"background_experience", "types_of_business", "types_of_business_raw"}
def process_pdf(client, model, pdf_path, vision_budget):
def _empty(v):
return v in (None, "", [], {})
def merge_records(text_rec, vision_rec):
"""
Fuehrt Text- und Vision-Extraktion desselben PDF zusammen.
Grundregel: nicht-leerer Wert schlaegt leeren; bei Konflikt entscheidet
die bevorzugte Quelle je Feld (Vision fuer Handschrift, Text fuer Getipptes).
"""
if text_rec is None:
return vision_rec
if vision_rec is None:
return text_rec
merged = {}
keys = set(text_rec) | set(vision_rec)
for k in keys:
tv, vv = text_rec.get(k), vision_rec.get(k)
if _empty(tv) and _empty(vv):
merged[k] = tv if k in text_rec else vv
elif _empty(tv):
merged[k] = vv
elif _empty(vv):
merged[k] = tv
else:
# beide gefuellt -> bevorzugte Quelle
if k in _PREFER_VISION:
merged[k] = vv
elif k in _PREFER_TEXT:
merged[k] = tv
else:
merged[k] = vv # Default: Vision
return merged
def _extract_via_text(client, model, text):
payload = USER_TEXT_INTRO + text[:12000]
raw, finish = call_model(client, model, payload)
data = parse_json_loose(raw)
if data is not None and finish == "length":
data["_runaway"] = True # Antwort war abgeschnitten -> unvollstaendig moeglich
return data
def _extract_via_vision(client, model, pdf_path, max_pages):
imgs = pdf_to_base64_images(pdf_path, max_pages=max_pages)
if not imgs:
return None
payload = [{"type": "text", "text": USER_IMG_INTRO}]
for b64 in imgs:
payload.append({"type": "image_url",
"image_url": {"url": f"data:image/jpeg;base64,{b64}"}})
raw, finish = call_model(client, model, payload)
data = parse_json_loose(raw)
if data is not None and finish == "length":
data["_runaway"] = True
return data
def process_pdf(client, model, pdf_path, vision_budget=None):
"""
Hybrid-Strategie:
- reiner Scan (keine Textebene): nur Vision
- grosser Scan (> MAX_SCAN_PAGES_TOTAL): nur erste Seite Vision
- getipptes PDF mit Textebene: Text UND Vision, dann mergen
(Text bringt getippten Background, Vision die handschriftlichen Felder)
"""
n_pages = pdf_page_count(pdf_path)
text = extract_text(pdf_path)
is_scan = len(text.strip()) < TEXT_SCAN_THRESHOLD
has_text = len(text.strip()) >= TEXT_SCAN_THRESHOLD
# Riesige, voll gescannte Stapel: nicht komplett verarbeiten.
# Wir schauen nur auf die erste Seite, ob ueberhaupt ein Sheet vorliegt.
truncated_note = None
if is_scan and n_pages > MAX_SCAN_PAGES_TOTAL:
truncated_note = f"grosser Scan ({n_pages} Seiten), nur erste Seite geprueft"
imgs = pdf_to_base64_images(pdf_path, max_pages=1)
elif is_scan:
imgs = pdf_to_base64_images(pdf_path, max_pages=2)
else:
imgs = None
if is_scan:
if not imgs:
return {"_error": "Scan, aber Bildkonvertierung fehlgeschlagen",
"n_pages": n_pages}, "vision"
payload = [{"type": "text", "text": USER_IMG_INTRO}]
for b64 in imgs:
payload.append({"type": "image_url",
"image_url": {"url": f"data:image/jpeg;base64,{b64}"}})
mode = "vision"
else:
payload = USER_TEXT_INTRO + text[:12000]
mode = "text"
try:
raw = call_model(client, model, payload, vision_budget if mode == "vision" else None)
if not has_text and n_pages > MAX_SCAN_PAGES_TOTAL:
# grosser reiner Scan: nur erste Seite pruefen
truncated_note = f"grosser Scan ({n_pages} Seiten), nur erste Seite geprueft"
data = _extract_via_vision(client, model, pdf_path, max_pages=1)
mode = "vision"
elif not has_text:
# reiner Scan
data = _extract_via_vision(client, model, pdf_path, max_pages=2)
mode = "vision"
else:
# Textebene vorhanden -> HYBRID: beide Pfade, dann mergen
text_rec = _extract_via_text(client, model, text)
vision_rec = None
if n_pages <= MAX_SCAN_PAGES_TOTAL:
vision_rec = _extract_via_vision(client, model, pdf_path, max_pages=2)
data = merge_records(text_rec, vision_rec)
mode = "hybrid" if vision_rec is not None else "text"
except Exception as e:
return {"_error": f"API: {type(e).__name__}: {e}", "n_pages": n_pages}, mode
return {"_error": f"API: {type(e).__name__}: {e}", "n_pages": n_pages}, "?"
data = parse_json_loose(raw)
if data is None:
return {"_error": "JSON-Parse fehlgeschlagen", "_raw": (raw or "")[:400],
"n_pages": n_pages}, mode
return {"_error": "JSON-Parse fehlgeschlagen", "n_pages": n_pages}, mode
data["n_pages"] = n_pages
if truncated_note:
data["_note"] = truncated_note
@@ -292,7 +375,7 @@ def main():
ap.add_argument("--limit", type=int, default=150)
ap.add_argument("--recursive", action="store_true",
help="auch Unterordner durchsuchen (Default: nur oberste Ebene)")
ap.add_argument("--timeout", type=float, default=120,
ap.add_argument("--timeout", type=float, default=60,
help="HTTP-Timeout je Anfrage in Sekunden")
ap.add_argument("--vision-budget", type=int, default=0,
help="Gemma Vision-Token-Budget (0=aus/Default 280; sonst 70/140/280/560/1120). "
@@ -319,7 +402,7 @@ def main():
f"{'rekursiv' if args.recursive else 'nur oberste Ebene'}).\n")
all_categories = {}
n_ok = n_err = n_text = n_vision = n_nosheet = 0
n_ok = n_err = n_text = n_vision = n_hybrid = n_nosheet = 0
with open(jsonl_path, "w", encoding="utf-8") as jf:
for i, pdf in enumerate(pdfs, 1):
@@ -336,6 +419,7 @@ def main():
n_text += (mode == "text")
n_vision += (mode == "vision")
n_hybrid += (mode == "hybrid")
record = {"source_file": rel,
"file_date": fdate.isoformat() if fdate else None,
"extraction_mode": mode, **data}
@@ -364,6 +448,7 @@ def main():
print(f" Fehler: {n_err}")
print(f" Text-Pfad: {n_text}")
print(f" Vision-Pfad: {n_vision}")
print(f" Hybrid-Pfad: {n_hybrid}")
print(f" kein Buyer-Sheet: {n_nosheet}")
print(f" distinkte Kategorien (roh): {len(all_categories)}")
print(f"\n JSONL: {jsonl_path}")
@@ -371,4 +456,4 @@ def main():
if __name__ == "__main__":
main()
main()

View File

@@ -67,7 +67,12 @@ JSON_SCHEMA = {
"background_experience": {"type": ["string", "null"]},
"date_of_introduction": {"type": ["string", "null"]},
},
"required": ["is_buyer_sheet", "types_of_business"],
"required": [
"is_buyer_sheet", "name_company", "prospective_buyer", "company",
"phone", "cell", "email", "address", "state", "how_did_you_hear",
"interested_in_updates", "types_of_business_raw", "types_of_business",
"background_experience", "date_of_introduction"
],
}
SYSTEM_PROMPT = """Du bist ein praezises Datenextraktions-System fuer Formulare der Firma "BizMatch Business Brokerage".
@@ -198,10 +203,8 @@ def pdf_to_base64_images(pdf_path, max_pages=2, max_dim=2200, quality=85, dpi=20
# ---------------------------------------------------------------------------
def call_model(client, model, content_payload, vision_budget=None):
extra_body = {"guided_json": JSON_SCHEMA}
# Gemma-4 Vision-Budget hochsetzen fuer feine Checkbox-Erkennung
if vision_budget:
extra_body["mm_processor_kwargs"] = {"max_soft_tokens": vision_budget}
# llama.cpp erzwingt JSON ueber response_format mit json_schema.
# (vLLM-spezifisches guided_json / mm_processor_kwargs gibt es hier nicht.)
resp = client.chat.completions.create(
model=model,
messages=[
@@ -210,7 +213,10 @@ def call_model(client, model, content_payload, vision_budget=None):
],
temperature=0.0,
max_tokens=1200,
extra_body=extra_body,
response_format={
"type": "json_schema",
"json_schema": {"name": "buyer_sheet", "schema": JSON_SCHEMA},
},
)
return resp.choices[0].message.content
@@ -233,45 +239,105 @@ def parse_json_loose(txt):
TEXT_SCAN_THRESHOLD = 120 # < so viele Textzeichen -> als Scan behandeln
MAX_SCAN_PAGES_TOTAL = 12 # PDFs mit mehr Seiten gelten als "kein reines Sheet"
# Felder, bei denen wir Vision bevorzugen (handschriftlich/visuell auf dem Formular)
_PREFER_VISION = {
"prospective_buyer", "name_company", "company", "phone", "cell", "email",
"address", "state", "how_did_you_hear", "interested_in_updates",
"date_of_introduction",
}
# Felder, bei denen Text meist die sauberere (getippte) Quelle ist
_PREFER_TEXT = {"background_experience", "types_of_business", "types_of_business_raw"}
def process_pdf(client, model, pdf_path, vision_budget):
def _empty(v):
return v in (None, "", [], {})
def merge_records(text_rec, vision_rec):
"""
Fuehrt Text- und Vision-Extraktion desselben PDF zusammen.
Grundregel: nicht-leerer Wert schlaegt leeren; bei Konflikt entscheidet
die bevorzugte Quelle je Feld (Vision fuer Handschrift, Text fuer Getipptes).
"""
if text_rec is None:
return vision_rec
if vision_rec is None:
return text_rec
merged = {}
keys = set(text_rec) | set(vision_rec)
for k in keys:
tv, vv = text_rec.get(k), vision_rec.get(k)
if _empty(tv) and _empty(vv):
merged[k] = tv if k in text_rec else vv
elif _empty(tv):
merged[k] = vv
elif _empty(vv):
merged[k] = tv
else:
# beide gefuellt -> bevorzugte Quelle
if k in _PREFER_VISION:
merged[k] = vv
elif k in _PREFER_TEXT:
merged[k] = tv
else:
merged[k] = vv # Default: Vision
return merged
def _extract_via_text(client, model, text):
payload = USER_TEXT_INTRO + text[:12000]
raw = call_model(client, model, payload)
return parse_json_loose(raw)
def _extract_via_vision(client, model, pdf_path, max_pages):
imgs = pdf_to_base64_images(pdf_path, max_pages=max_pages)
if not imgs:
return None
payload = [{"type": "text", "text": USER_IMG_INTRO}]
for b64 in imgs:
payload.append({"type": "image_url",
"image_url": {"url": f"data:image/jpeg;base64,{b64}"}})
raw = call_model(client, model, payload)
return parse_json_loose(raw)
def process_pdf(client, model, pdf_path, vision_budget=None):
"""
Hybrid-Strategie:
- reiner Scan (keine Textebene): nur Vision
- grosser Scan (> MAX_SCAN_PAGES_TOTAL): nur erste Seite Vision
- getipptes PDF mit Textebene: Text UND Vision, dann mergen
(Text bringt getippten Background, Vision die handschriftlichen Felder)
"""
n_pages = pdf_page_count(pdf_path)
text = extract_text(pdf_path)
is_scan = len(text.strip()) < TEXT_SCAN_THRESHOLD
has_text = len(text.strip()) >= TEXT_SCAN_THRESHOLD
# Riesige, voll gescannte Stapel: nicht komplett verarbeiten.
# Wir schauen nur auf die erste Seite, ob ueberhaupt ein Sheet vorliegt.
truncated_note = None
if is_scan and n_pages > MAX_SCAN_PAGES_TOTAL:
truncated_note = f"grosser Scan ({n_pages} Seiten), nur erste Seite geprueft"
imgs = pdf_to_base64_images(pdf_path, max_pages=1)
elif is_scan:
imgs = pdf_to_base64_images(pdf_path, max_pages=2)
else:
imgs = None
if is_scan:
if not imgs:
return {"_error": "Scan, aber Bildkonvertierung fehlgeschlagen",
"n_pages": n_pages}, "vision"
payload = [{"type": "text", "text": USER_IMG_INTRO}]
for b64 in imgs:
payload.append({"type": "image_url",
"image_url": {"url": f"data:image/jpeg;base64,{b64}"}})
mode = "vision"
else:
payload = USER_TEXT_INTRO + text[:12000]
mode = "text"
try:
raw = call_model(client, model, payload, vision_budget if mode == "vision" else None)
if not has_text and n_pages > MAX_SCAN_PAGES_TOTAL:
# grosser reiner Scan: nur erste Seite pruefen
truncated_note = f"grosser Scan ({n_pages} Seiten), nur erste Seite geprueft"
data = _extract_via_vision(client, model, pdf_path, max_pages=1)
mode = "vision"
elif not has_text:
# reiner Scan
data = _extract_via_vision(client, model, pdf_path, max_pages=2)
mode = "vision"
else:
# Textebene vorhanden -> HYBRID: beide Pfade, dann mergen
text_rec = _extract_via_text(client, model, text)
vision_rec = None
if n_pages <= MAX_SCAN_PAGES_TOTAL:
vision_rec = _extract_via_vision(client, model, pdf_path, max_pages=2)
data = merge_records(text_rec, vision_rec)
mode = "hybrid" if vision_rec is not None else "text"
except Exception as e:
return {"_error": f"API: {type(e).__name__}: {e}", "n_pages": n_pages}, mode
return {"_error": f"API: {type(e).__name__}: {e}", "n_pages": n_pages}, "?"
data = parse_json_loose(raw)
if data is None:
return {"_error": "JSON-Parse fehlgeschlagen", "_raw": (raw or "")[:400],
"n_pages": n_pages}, mode
return {"_error": "JSON-Parse fehlgeschlagen", "n_pages": n_pages}, mode
data["n_pages"] = n_pages
if truncated_note:
data["_note"] = truncated_note
@@ -287,7 +353,7 @@ def main():
ap.add_argument("--src", required=True)
ap.add_argument("--out", default="./poc_out")
ap.add_argument("--api", default="http://192.168.100.160:8000/v1")
ap.add_argument("--model", default="google/gemma-4-12b-it")
ap.add_argument("--model", default="gemma-4-12b") # --alias des llama-server
ap.add_argument("--limit", type=int, default=150)
ap.add_argument("--recursive", action="store_true",
help="auch Unterordner durchsuchen (Default: nur oberste Ebene)")
@@ -318,7 +384,7 @@ def main():
f"{'rekursiv' if args.recursive else 'nur oberste Ebene'}).\n")
all_categories = {}
n_ok = n_err = n_text = n_vision = n_nosheet = 0
n_ok = n_err = n_text = n_vision = n_hybrid = n_nosheet = 0
with open(jsonl_path, "w", encoding="utf-8") as jf:
for i, pdf in enumerate(pdfs, 1):
@@ -335,6 +401,7 @@ def main():
n_text += (mode == "text")
n_vision += (mode == "vision")
n_hybrid += (mode == "hybrid")
record = {"source_file": rel,
"file_date": fdate.isoformat() if fdate else None,
"extraction_mode": mode, **data}
@@ -363,6 +430,7 @@ def main():
print(f" Fehler: {n_err}")
print(f" Text-Pfad: {n_text}")
print(f" Vision-Pfad: {n_vision}")
print(f" Hybrid-Pfad: {n_hybrid}")
print(f" kein Buyer-Sheet: {n_nosheet}")
print(f" distinkte Kategorien (roh): {len(all_categories)}")
print(f"\n JSONL: {jsonl_path}")
@@ -370,4 +438,4 @@ def main():
if __name__ == "__main__":
main()
main()

189
merge_buyers.py Normal file
View File

@@ -0,0 +1,189 @@
#!/usr/bin/env python3
"""
merge_buyers.py - Fuehrt Buyer-Datensaetze auf PERSONENEBENE zusammen.
Problem: Eine Person hat oft mehrere Buyer Information Sheets (verschiedene
Daten, Text- und Notes-Varianten). Beispiel Sudduth: eine Text-Datei + eine
Notes-Scan-Datei = dieselbe Person. Oder Sahota mit 6 Dateien.
Dieses Script liest die rohe buyers_raw.jsonl (ein Datensatz je DATEI) und
erzeugt buyers_merged.jsonl (ein Datensatz je PERSON), mit:
- allen Business-Kategorien ueber alle Sheets vereinigt
- fruehestem und spaetestem Date-of-Introduction
- je Kontaktfeld dem besten (nicht-leeren) Wert
- Liste der Quelldateien zur Nachvollziehbarkeit
Die rohe Extraktion bleibt unangetastet (nachvollziehbar). Merge ist ein
separater, pruefbarer Schritt.
Aufruf:
python merge_buyers.py ./poc_out/buyers_raw.jsonl ./poc_out/buyers_merged.jsonl
"""
import sys
import json
import re
from collections import defaultdict
def _empty(v):
return v in (None, "", [], {})
def name_from_filename(source_file):
"""
Extrahiert 'Nachname, Vorname' aus dem Dateinamen als robusten Fallback.
Die Dateien folgen konsistent dem Schema 'Nachname, Vorname <datum> [Notes].pdf'.
"""
if not source_file:
return ""
base = source_file.rsplit("/", 1)[-1] # nur Dateiname
base = re.sub(r"\.pdf$", "", base, flags=re.I)
# Datum, 'Notes', Zahlen und Zusaetze abschneiden
base = re.sub(r"\b\d{6,8}\b.*$", "", base) # ab erstem Datum abschneiden
base = re.sub(r"(?i)\bnotes\b.*$", "", base)
base = base.strip(" -_")
return base
def normalize_name(rec):
"""
Personen-Schluessel. Bevorzugt den echten Namen (prospective_buyer),
faellt auf name_company zurueck, dann auf den DATEINAMEN (robust, da
konsistentes Schema). Normalisiert "Nachname, Vorname" und
"Vorname Nachname" auf eine vergleichbare Form.
"""
name = rec.get("prospective_buyer") or rec.get("name_company") or ""
if not name.strip():
# Fallback: aus Dateiname (gerade beim Text-Pfad oft noetig)
name = name_from_filename(rec.get("source_file", ""))
name = name.strip().lower()
name = re.sub(r"\s+", " ", name)
# "sudduth, henry" -> "henry sudduth" (Komma-Form angleichen)
if "," in name:
parts = [p.strip() for p in name.split(",", 1)]
if len(parts) == 2 and parts[1]:
name = f"{parts[1]} {parts[0]}"
# Satzzeichen weg
name = re.sub(r"[^\w\s]", "", name)
return name.strip()
def pick_best(values):
"""Ersten nicht-leeren Wert aus einer Liste waehlen."""
for v in values:
if not _empty(v):
return v
return None
def merge_person(records):
"""Fuehrt alle Sheets EINER Person zu einem Datensatz zusammen."""
# Kontaktfelder: bester nicht-leerer Wert (spaetere Sheets zuerst,
# da meist aktueller - wir sortieren unten nach Datum absteigend)
single_fields = ["name_company", "prospective_buyer", "company", "phone",
"cell", "email", "address", "state", "how_did_you_hear",
"background_experience"]
out = {}
for f in single_fields:
out[f] = pick_best([r.get(f) for r in records])
# interested_in_updates: wenn IRGENDEIN Sheet true/false sagt, nimm das
# (bevorzugt das neueste eindeutige)
upd = pick_best([r.get("interested_in_updates") for r in records
if r.get("interested_in_updates") is not None])
out["interested_in_updates"] = upd
# types_of_business: Vereinigung ueber alle Sheets
cats = []
for r in records:
for c in (r.get("types_of_business") or []):
c = c.strip()
if c and c not in cats:
cats.append(c)
out["types_of_business"] = cats
# Daten: alle gueltigen ISO-Daten sammeln
dates = sorted(d for d in (r.get("date_of_introduction") for r in records)
if isinstance(d, str) and re.match(r"\d{4}-\d{2}-\d{2}", d))
out["date_first_introduction"] = dates[0] if dates else None
out["date_last_introduction"] = dates[-1] if dates else None
out["all_introduction_dates"] = dates
# Nachvollziehbarkeit
out["source_files"] = [r.get("source_file") for r in records]
out["n_sheets"] = len(records)
return out
def main():
if len(sys.argv) < 3:
print("Aufruf: python merge_buyers.py <input.jsonl> <output.jsonl>")
sys.exit(1)
inp, outp = sys.argv[1], sys.argv[2]
records = []
with open(inp, encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
r = json.loads(line)
# Fehler und Nicht-Sheets ueberspringen
if "_error" in r:
continue
if r.get("is_buyer_sheet") is False:
continue
records.append(r)
# nach Person gruppieren
groups = defaultdict(list)
unkeyed = []
for r in records:
key = normalize_name(r)
if key:
groups[key].append(r)
else:
unkeyed.append(r) # ohne erkennbaren Namen: einzeln behalten
merged = []
for key, recs in groups.items():
# innerhalb der Person nach Datum absteigend (neuestes zuerst)
recs_sorted = sorted(
recs,
key=lambda r: (r.get("date_of_introduction") or ""),
reverse=True,
)
m = merge_person(recs_sorted)
m["person_key"] = key
merged.append(m)
# namenlose einzeln anhaengen
for r in unkeyed:
m = merge_person([r])
m["person_key"] = "(kein Name erkannt)"
merged.append(m)
# nach neuestem Datum sortieren
merged.sort(key=lambda m: (m.get("date_last_introduction") or ""), reverse=True)
with open(outp, "w", encoding="utf-8") as f:
for m in merged:
f.write(json.dumps(m, ensure_ascii=False) + "\n")
# Statistik
multi = [m for m in merged if m["n_sheets"] > 1]
print(f"Eingelesen: {len(records)} Datensaetze (Sheets)")
print(f"Distinkte Personen: {len(merged)}")
print(f"davon mit >1 Sheet: {len(multi)}")
print(f"Ausgabe: {outp}")
if multi:
print("\nBeispiele (Personen mit mehreren Sheets):")
for m in sorted(multi, key=lambda x: -x["n_sheets"])[:8]:
print(f" {m['person_key']!r}: {m['n_sheets']} Sheets, "
f"Kategorien={m['types_of_business']}, "
f"Daten {m['date_first_introduction']}..{m['date_last_introduction']}")
if __name__ == "__main__":
main()