#!/usr/bin/env python3 """Phase 2, Modul 5: Satzvektoren (ebene = 'embedding'). Ein Encoder bildet Titel und Teaser auf einen Vektor ab. Zwei Meldungen ueber denselben Vorgang liegen dann nah beieinander, auch wenn kein Wort uebereinstimmt - genau das, was das Dublettenmodul konstruktionsbedingt nicht kann. python3 embeddings.py --backfill python3 embeddings.py --slot 2026-09-07T08:00 python3 embeddings.py --backfill --index # HNSW-Index danach anlegen python3 embeddings.py --nachweis # Determinismus messen Phase 2 clustert nicht. Der Vektor ist Vorfilter fuer die Kandidatengenerierung in Phase 3b, kein Selbstzweck (Abschnitt 7 der Spezifikation). Kein LLM: multilingual-e5-large ist ein Encoder. Text hinein, Vektor heraus - kein Prompt, keine Sampling-Temperatur, keine Textausgabe, feste Gewichte, kein Netzzugriff nach dem ersten Laden. """ import argparse import hashlib import os import struct import sys import time import normalisierung import version as versionsmodul from db import verbindung EBENE = "embedding" MODELL = "intfloat/multilingual-e5-large" DIMENSIONEN = 1024 # e5 erwartet ein Praefix. Die Modellkarte unterscheidet "query: " und # "passage: " fuer asymmetrische Suche - Frage gegen Dokument. Hier werden # Meldungen mit Meldungen verglichen, also symmetrisch, und dafuer nennt # dieselbe Karte "query: " auf beiden Seiten. PRAEFIX = "query: " # Feste Laenge fuer *jede* Zeile, nicht nur je Stapel. # # Das ist der Kern der Reproduzierbarkeit. Ein Encoder rechnet stapelweise; # fuellt man nur bis zur laengsten Zeile des Stapels auf, aendert sich die # Form der Matrizen mit der Zusammensetzung des Stapels, und damit die # Reihenfolge der Gleitkommaadditionen. Ein Slot-Lauf ueber 31 Items lieferte # dann andere letzte Stellen als ein Backfill ueber 3407 - und die Zusage # "Live-Lauf gleich Neuaufbau" waere hin. # # Bei fester Laenge hat jeder Stapel dieselbe Form, und eine Zeile haengt # nicht mehr von ihren Nachbarn ab. # # 128 Token: Titel und Teaser zusammen liegen bei etwa 31 Woertern, im # Bestand ueberschreitet keine Zeile die Grenze (nachgeprueft beim Lauf, # siehe `gekuerzt` im Bericht). MAX_TOKEN = 128 STAPEL = 32 # Threadzahl. Sie bestimmt, wie die Summen aufgeteilt werden, und gehoert # darum in den Versions-Hash. Der Container setzt sie ueber OMP_NUM_THREADS. THREADS = int(os.environ.get("OMP_NUM_THREADS", "4")) def einstellungen(revision=None, torch_version=None): """Was in die coder_version eingeht. Auch die Bibliotheksversion: ein anderer torch kann andere Kernel waehlen und damit andere letzte Stellen liefern. Das entwertet die Vektoren nicht, aber es macht sie zu einem anderen Bestand, und genau das soll der Hash sagen. """ return { "modell": MODELL, "modell_revision": revision, "dimensionen": DIMENSIONEN, "praefix": PRAEFIX, "basis": "titel+teaser", "text": "bereinigt", "max_token": MAX_TOKEN, "fuellung": "feste_laenge", "pooling": "mittel_maskiert", "l2_normiert": True, "dtype": "float32", "stapel": STAPEL, "threads": THREADS, "torch": torch_version, } # -------------------------------------------------------------------------- # Modell # -------------------------------------------------------------------------- def lade_modell(): """(tokenizer, modell, revision, torch_version). Der Import steht hier und nicht oben: die uebrigen Module dieses Repos laufen ohne torch, und `python3 -m unittest` soll nicht 2 GB laden muessen, um die Normalisierung zu pruefen. """ import torch from transformers import AutoModel, AutoTokenizer torch.set_num_threads(THREADS) torch.set_grad_enabled(False) tokenizer = AutoTokenizer.from_pretrained(MODELL) modell = AutoModel.from_pretrained(MODELL, dtype=torch.float32) modell.eval() # Die aufgeloeste Revision aus dem Cache. Ohne sie sagt "e5-large" nur, # welches Modell gemeint war, nicht welche Gewichte gerechnet haben. revision = getattr(getattr(modell, "config", None), "_commit_hash", None) return tokenizer, modell, revision, torch.__version__ def kodiere(tokenizer, modell, texte, stapel=STAPEL): """Vektoren zu den Texten, in derselben Reihenfolge. Mittelwert ueber die nicht aufgefuellten Positionen, dann L2-normiert - das Verfahren der e5-Modellkarte. Nach der Normierung ist das Skalarprodukt die Kosinus-Aehnlichkeit, und pgvector kann mit vector_cosine_ops direkt darauf suchen. """ import torch alle = [] for anfang in range(0, len(texte), stapel): teil = texte[anfang:anfang + stapel] eingabe = tokenizer(teil, padding="max_length", truncation=True, max_length=MAX_TOKEN, return_tensors="pt") ausgabe = modell(**eingabe).last_hidden_state maske = eingabe["attention_mask"].unsqueeze(-1).to(ausgabe.dtype) mittel = (ausgabe * maske).sum(1) / maske.sum(1).clamp(min=1e-9) alle.append(torch.nn.functional.normalize(mittel, p=2, dim=1)) return torch.cat(alle) if alle else torch.empty(0, DIMENSIONEN) def zu_lang(tokenizer, texte): """Wieviele Texte die Tokengrenze reissen. Sollte 0 sein.""" return sum(1 for t in texte if len(tokenizer(t, truncation=False)["input_ids"]) > MAX_TOKEN) # -------------------------------------------------------------------------- # Laden und Schreiben # -------------------------------------------------------------------------- _SPALTEN = "select id, quelle, titel, teaser from raw_items" _ORDNUNG = " order by id" def lade(conn, slot=None): if slot is None: zeilen = conn.execute(_SPALTEN + _ORDNUNG).fetchall() else: zeilen = conn.execute(_SPALTEN + " where slot = %s" + _ORDNUNG, (slot,)).fetchall() return zeilen def text(titel, teaser): """Was kodiert wird. bereinige() statt normalisiere(): der Encoder ist auf natuerlichem Text trainiert. Grossschreibung unterscheidet im Deutschen Wortarten, und die Anfuehrungszeichen um ein Zitat sind Bedeutung, kein Rauschen. """ return PRAEFIX + normalisierung.bereinige(titel, teaser) def vektor_hash(werte): """Kennung des Vektors, damit sich ohne Fliesskommavergleich sagen laesst, ob sich etwas geaendert hat - und damit der Determinismus pruefbar ist. struct statt numpy: dieselben Bytes (float32, little endian), aber ohne Abhaengigkeit. Die Funktion laesst sich damit auch dort pruefen, wo kein torch installiert ist. """ liste = werte.tolist() if hasattr(werte, "tolist") else list(werte) roh = struct.pack(f"<{len(liste)}f", *liste) return hashlib.blake2b(roh, digest_size=8).hexdigest() def als_vector(werte): """pgvector-Literal. Spart die Abhaengigkeit auf das Adapterpaket.""" return "[" + ",".join(repr(float(w)) for w in werte) + "]" def schreibe(conn, zeilen, vektoren, coder_version, revision, trocken=False): """Vektor nach item_embeddings, Kennung nach codings. Zwei Tabellen, weil ein Vektor in JSONB weder indizierbar noch platzsparend waere (siehe p2-002). Die Coding-Zeile ist der Nachweis, dass ein Item verarbeitet wurde, und traegt, womit. """ from psycopg.types.json import Jsonb vorhanden = { r[0]: r[1] for r in conn.execute( "select raw_item_id, nutzlast from codings" " where coder_version = %s and ebene = %s", (coder_version, EBENE)) } neu = geaendert = unveraendert = 0 vektorstapel, codingstapel = [], [] for (item_id, _quelle, titel, teaser), v in zip(zeilen, vektoren): last = { "modell": MODELL, "modell_revision": revision, "dimensionen": DIMENSIONEN, "basis": "titel+teaser", "vektor_hash": vektor_hash(v), "woerter": len(normalisierung.bereinige(titel, teaser).split()), } alt = vorhanden.get(item_id) if alt == last: unveraendert += 1 continue if alt is None: neu += 1 else: geaendert += 1 vektorstapel.append((item_id, coder_version, als_vector(v))) codingstapel.append((item_id, coder_version, EBENE, Jsonb(last))) if codingstapel and not trocken: with conn.cursor() as cur: cur.executemany( """insert into item_embeddings (raw_item_id, coder_version, vektor) values (%s, %s, %s::vector) on conflict (raw_item_id, coder_version) do update set vektor = excluded.vektor, erzeugt_am = now()""", vektorstapel) cur.executemany( """insert into codings (raw_item_id, coder_version, ebene, nutzlast) values (%s, %s, %s, %s) on conflict (raw_item_id, coder_version, ebene) do update set nutzlast = excluded.nutzlast, kodiert_am = now()""", codingstapel) return {"neu": neu, "geaendert": geaendert, "unveraendert": unveraendert} def protokolliere(conn, slot, coder_version, dauer_ms, gesehen, kodiert, fehler=None): conn.execute( """insert into p2_laeufe (slot, modul, coder_version, dauer_ms, items_gesehen, items_kodiert, items_fehler, fehler) values (%s, %s, %s, %s, %s, %s, %s, %s)""", (slot, EBENE, coder_version, dauer_ms, gesehen, kodiert, 0 if not fehler else 1, fehler)) # -------------------------------------------------------------------------- # Nachweis # -------------------------------------------------------------------------- def nachweis(tokenizer, modell, zeilen, stapel=STAPEL): """Haengt ein Vektor von der Zusammensetzung seines Stapels ab? Die Zusage aus Invariante 2 lautet, dass ein Slot-Lauf dasselbe liefert wie ein Neuaufbau. Bei den regelbasierten Modulen folgt das aus dem Verfahren; hier muss es gemessen werden. Geprueft wird das Aeusserste, was im Betrieb vorkommt: dieselben Zeilen einmal einzeln, einmal im vollen Stapel, einmal in umgekehrter Reihenfolge. """ texte = [text(t, s) for _, _, t, s in zeilen] voll = kodiere(tokenizer, modell, texte, stapel=stapel) einzeln = kodiere(tokenizer, modell, texte, stapel=1) rueckwaerts = kodiere(tokenizer, modell, texte[::-1], stapel=stapel).flip(0) gleich_einzeln = sum(vektor_hash(a) == vektor_hash(b) for a, b in zip(voll, einzeln)) gleich_rueck = sum(vektor_hash(a) == vektor_hash(b) for a, b in zip(voll, rueckwaerts)) abstand_einzeln = float((voll - einzeln).abs().max()) abstand_rueck = float((voll - rueckwaerts).abs().max()) # Bei bitgleichen Vektoren liegt dieser Wert knapp unter 1 - nicht weil # sie sich unterscheiden, sondern weil die Summe von 1024 Quadraten in # float32 nicht exakt auf 1 faellt. Aussagekraft hat er nur, wenn die # Zeilen darueber eine Abweichung melden. kosinus = float((voll * einzeln).sum(1).min()) print(f" Zeilen geprueft: {len(texte)}") print(f" bitgleich einzeln: {gleich_einzeln}/{len(texte)}" f" groesste Abweichung {abstand_einzeln:.2e}") print(f" bitgleich rueckwaerts:{gleich_rueck}/{len(texte)}" f" groesste Abweichung {abstand_rueck:.2e}") print(f" kleinster Kosinus zwischen den Fassungen: {kosinus:.9f}" + (" (Rundung der Quadratsumme, keine Abweichung)" if abstand_einzeln == 0.0 and abstand_rueck == 0.0 else "")) return gleich_einzeln == len(texte) and gleich_rueck == len(texte) def main(argv=None): p = argparse.ArgumentParser(description=__doc__.splitlines()[0]) p.add_argument("--slot", help="ein Slot, z.B. 2026-09-07T08:00") p.add_argument("--backfill", action="store_true", help="Gesamtbestand") p.add_argument("--nachweis", type=int, nargs="?", const=64, metavar="N", help="Determinismus an N Zeilen messen, nichts schreiben") p.add_argument("--index", action="store_true", help="HNSW-Index fuer diese Version anlegen") p.add_argument("--trocken", action="store_true", help="rechnen, nichts schreiben") a = p.parse_args(argv) if not (a.backfill or a.slot or a.nachweis): p.error("--backfill, --slot oder --nachweis angeben") import datetime as dt slot = dt.datetime.fromisoformat(a.slot) if a.slot else None begonnen = time.monotonic() print(f"Modell {MODELL} laedt, {THREADS} Threads ...") tokenizer, modell, revision, torch_version = lade_modell() print(f" Revision {revision} torch {torch_version}" f" ({int(time.monotonic() - begonnen)} s)") with verbindung() as conn: if a.nachweis: zeilen = conn.execute(_SPALTEN + _ORDNUNG + " limit %s", (a.nachweis,)).fetchall() return 0 if nachweis(tokenizer, modell, zeilen) else 1 v = versionsmodul.eintragen( conn, komponenten=versionsmodul.komponenten( embedding=einstellungen(revision, torch_version))) print(f"coder_version {v}{' (trocken)' if a.trocken else ''}") zeilen = lade(conn, slot=slot) if not zeilen: print(" nichts zu tun") return 0 texte = [text(t, s) for _, _, t, s in zeilen] lang = zu_lang(tokenizer, texte) gestartet = time.monotonic() vektoren = kodiere(tokenizer, modell, texte) dauer_kodierung = time.monotonic() - gestartet print(f" Items: {len(zeilen)}") # Ein paar gekuerzte Zeilen sind erwartbar (Sendungsablaeufe wie # "tagesschau 20:00 Uhr" listen alle Themen auf). Gewarnt wird erst, # wenn es mehr als ein Prozent trifft - dann stimmt die Grenze nicht. print(f" ueber {MAX_TOKEN} Token: {lang}" + (" <- Grenze pruefen" if lang > len(zeilen) / 100 else "")) print(f" Kodierung: {dauer_kodierung:.1f} s" f" ({len(zeilen) / max(dauer_kodierung, 1e-9):.0f} Items/s)") zahlen = schreibe(conn, zeilen, vektoren, v, revision, trocken=a.trocken) dauer = int((time.monotonic() - begonnen) * 1000) print(f" Vektoren: neu {zahlen['neu']}," f" geaendert {zahlen['geaendert']}," f" unveraendert {zahlen['unveraendert']}") if a.index and not a.trocken: name = conn.execute("select p2_embedding_index(%s)", (v,)).fetchone()[0] print(f" Index: {name}") print(f" Dauer gesamt: {dauer} ms") if not a.trocken: protokolliere(conn, slot, v, dauer, len(zeilen), zahlen["neu"] + zahlen["geaendert"]) conn.commit() else: conn.rollback() return 0 if __name__ == "__main__": sys.exit(main())