diff --git a/bayarea-ai-server/docker-compose-llama-gemma12b-vulkan.yml b/bayarea-ai-server/docker-compose-llama-gemma12b-vulkan.yml new file mode 100644 index 0000000..3486a32 --- /dev/null +++ b/bayarea-ai-server/docker-compose-llama-gemma12b-vulkan.yml @@ -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 (-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 \ No newline at end of file diff --git a/bayarea-ai-server/docker-compose-llama-gemma12b.yml b/bayarea-ai-server/docker-compose-llama-gemma12b.yml index 7d74972..d3186e7 100644 --- a/bayarea-ai-server/docker-compose-llama-gemma12b.yml +++ b/bayarea-ai-server/docker-compose-llama-gemma12b.yml @@ -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 \ No newline at end of file diff --git a/extract_buyers_llamacpp.py b/extract_buyers_llamacpp.py index 733dcec..038c08c 100644 --- a/extract_buyers_llamacpp.py +++ b/extract_buyers_llamacpp.py @@ -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 "-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() \ No newline at end of file diff --git a/extract_buyers_poc.py b/extract_buyers_poc.py index 1bffdd9..6003ac5 100644 --- a/extract_buyers_poc.py +++ b/extract_buyers_poc.py @@ -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() \ No newline at end of file diff --git a/merge_buyers.py b/merge_buyers.py new file mode 100644 index 0000000..942e125 --- /dev/null +++ b/merge_buyers.py @@ -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 [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 ") + 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() \ No newline at end of file