Dublettenprüfung
This commit is contained in:
Executable
+659
@@ -0,0 +1,659 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Phase 2, Modul 1: exakte Uebernahmen (ebene = 'dublette').
|
||||
|
||||
Agenturmeldungen erscheinen bei mehreren Haeusern nahezu gleichlautend.
|
||||
Dieses Modul fasst sie zusammen, damit Phase 3 nur den Leitartikel kodiert.
|
||||
|
||||
Kein thematisches Clustering. Nur near-exact: erkannt wird, dass zwei Haeuser
|
||||
denselben Text ausgespielt haben, nicht dass sie ueber dasselbe schreiben.
|
||||
|
||||
python3 dubletten.py --backfill # Gesamtbestand
|
||||
python3 dubletten.py --slot 2026-09-07T08:00
|
||||
python3 dubletten.py --backfill --trocken --stichprobe 200
|
||||
|
||||
Verfahren nach Abschnitt 3 der Spezifikation: Normalisierung, SimHash ueber
|
||||
Wort-Trigramme, Kandidaten ueber Hamming-Distanz, Bestaetigung ueber Jaccard,
|
||||
Union-Find.
|
||||
|
||||
Wo das Verfahren von der Spezifikation abweicht, steht die Begruendung an
|
||||
Ort und Stelle: BASIS_VORGABE, gruppen_id(), die Paarregel in paare() und
|
||||
_zerlege_ketten(). README.md fasst sie zusammen, NOTES.md haelt die Messungen
|
||||
fest, auf denen sie beruhen.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import datetime as dt
|
||||
import hashlib
|
||||
import itertools
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
import uuid
|
||||
|
||||
from psycopg.types.json import Jsonb
|
||||
|
||||
import normalisierung
|
||||
import version as versionsmodul
|
||||
from db import verbindung
|
||||
|
||||
EBENE = "dublette"
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Vergleichsbasis: "titel" oder "titel+teaser".
|
||||
#
|
||||
# Die Spezifikation sagt "Titel + Teaser, oder Volltext falls vorhanden".
|
||||
# Beides taugt hier nicht.
|
||||
#
|
||||
# Volltext nicht, weil ihn nur ein Teil der Quellen liefert - Ingest liest
|
||||
# RSS, und viele Haeuser spielen dort kein content:encoded aus (Stand
|
||||
# 7.9.2026: faz 311/489, heise 174/196, aber derstandard, welt, dlf, taz, sz
|
||||
# jeweils 0). Eine je Item verschiedene Basis macht ausgerechnet den
|
||||
# haeufigsten Fall unauffindbar: dieselbe Meldung bei einem Haus mit und bei
|
||||
# einem ohne Volltext.
|
||||
#
|
||||
# Titel + Teaser nicht, weil die Haeuser die Agenturueberschrift woertlich
|
||||
# uebernehmen, den Teaser aber selbst schreiben. Gemessen wurde beides: ueber
|
||||
# Titel + Teaser findet das Verfahren auf diesem Bestand keine einzige
|
||||
# hausuebergreifende Uebernahme, ueber den Titel allein 63 Gruppen. Der
|
||||
# Teaser verduennt genau das Signal, das traegt. Zahlen in NOTES.md.
|
||||
# --------------------------------------------------------------------------
|
||||
BASIS_VORGABE = "titel"
|
||||
|
||||
# Mindestwortzahl je Basis. Wenige Trigramme erreichen leicht einen hohen
|
||||
# Jaccard-Wert, ohne dass die Meldungen etwas miteinander zu tun haetten -
|
||||
# lieber eine Dublette verpassen als zwei Vorgaenge verschmelzen
|
||||
# (Abschnitt 3, "Erwartung"). Titel sind kuerzer, die Schwelle entsprechend.
|
||||
MINDESTWOERTER = {"titel": 5, "titel+teaser": 8}
|
||||
|
||||
# Reissleine fuer die Randerweiterung beim Slot-Lauf.
|
||||
RANDERWEITERUNGEN = 6
|
||||
|
||||
EINSTELLUNGEN = {
|
||||
"basis": BASIS_VORGABE,
|
||||
"simhash_bits": 64,
|
||||
"trigramm": "wort-3",
|
||||
"hamming_max": 3,
|
||||
"jaccard_min": 0.85,
|
||||
"fenster_tage": 1,
|
||||
"mindestwoerter": MINDESTWOERTER[BASIS_VORGABE],
|
||||
"ressort_abgetrennt": True,
|
||||
}
|
||||
|
||||
_AGENTUREN = [
|
||||
("dpa", re.compile(r"\((?:dpa|dpa-AFX)[^)]{0,20}\)|\bdpa\b")),
|
||||
("afp", re.compile(r"\(AFP[^)]{0,20}\)|\bAFP\b")),
|
||||
("reuters", re.compile(r"\bReuters\b")),
|
||||
("ap", re.compile(r"\(AP\)|\bAssociated Press\b")),
|
||||
("epd", re.compile(r"\bepd\b")),
|
||||
("kna", re.compile(r"\bKNA\b")),
|
||||
("sid", re.compile(r"\bSID\b")),
|
||||
("dts", re.compile(r"\bdts Nachrichtenagentur\b")),
|
||||
]
|
||||
|
||||
# Fester Namensraum, damit Gruppen-UUIDs ueber Laeufe und Maschinen hinweg
|
||||
# gleich bleiben. Willkuerlich gewaehlt, aber unveraenderlich.
|
||||
NAMENSRAUM = uuid.UUID("6f9b4d2e-1c53-4a7f-9e18-2b0d7c5a3e41")
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# SimHash
|
||||
# --------------------------------------------------------------------------
|
||||
def _merkmalshash(trigramm):
|
||||
"""64 Bit je Trigramm.
|
||||
|
||||
blake2b und nicht hash(): Pythons hash() fuer Zeichenketten ist je
|
||||
Prozess zufaellig gesalzen. Damit waere nichts reproduzierbar.
|
||||
"""
|
||||
roh = "\x1f".join(trigramm).encode("utf-8")
|
||||
return int.from_bytes(hashlib.blake2b(roh, digest_size=8).digest(), "big")
|
||||
|
||||
|
||||
def simhash(trigramme):
|
||||
if not trigramme:
|
||||
return 0
|
||||
gewichte = [0] * 64
|
||||
for tri in trigramme:
|
||||
h = _merkmalshash(tri)
|
||||
for bit in range(64):
|
||||
gewichte[bit] += 1 if (h >> bit) & 1 else -1
|
||||
wert = 0
|
||||
for bit in range(64):
|
||||
# Gleichstand faellt auf 0. Willkuerlich, aber deterministisch.
|
||||
if gewichte[bit] > 0:
|
||||
wert |= 1 << bit
|
||||
return wert
|
||||
|
||||
|
||||
def hamming(a, b):
|
||||
return (a ^ b).bit_count()
|
||||
|
||||
|
||||
BAENDER = 4
|
||||
|
||||
|
||||
def baender(wert):
|
||||
"""SimHash in vier Stuecke zu 16 Bit zerlegen.
|
||||
|
||||
Schubfachschluss: unterscheiden sich zwei Werte in hoechstens 3 Bits, so
|
||||
ist bei vier Baendern mindestens eines gleich. Nur Paare, die sich ein
|
||||
Band teilen, muessen ueberhaupt verglichen werden - das ersetzt den
|
||||
Vergleich jeder Meldung mit jeder. Die Zahl der Baender haengt an
|
||||
hamming_max und darf nicht ohne dieses geaendert werden.
|
||||
"""
|
||||
breite = 64 // BAENDER
|
||||
maske = (1 << breite) - 1
|
||||
return tuple((i, (wert >> (i * breite)) & maske) for i in range(BAENDER))
|
||||
|
||||
|
||||
def jaccard(a, b):
|
||||
if not a or not b:
|
||||
return 0.0
|
||||
schnitt = len(a & b)
|
||||
return schnitt / (len(a) + len(b) - schnitt)
|
||||
|
||||
|
||||
class UnionFind:
|
||||
def __init__(self):
|
||||
self.eltern = {}
|
||||
|
||||
def finde(self, x):
|
||||
self.eltern.setdefault(x, x)
|
||||
while self.eltern[x] != x:
|
||||
self.eltern[x] = self.eltern[self.eltern[x]]
|
||||
x = self.eltern[x]
|
||||
return x
|
||||
|
||||
def vereinige(self, a, b):
|
||||
wa, wb = self.finde(a), self.finde(b)
|
||||
if wa != wb:
|
||||
# Kleinere Wurzel gewinnt: macht das Ergebnis unabhaengig von der
|
||||
# Reihenfolge, in der die Paare hereinkommen.
|
||||
hoch, tief = max(wa, wb), min(wa, wb)
|
||||
self.eltern[hoch] = tief
|
||||
|
||||
|
||||
def agentur(*texte):
|
||||
"""Agentur nur mit Textbeleg. Keine Vermutung ins Blaue."""
|
||||
text = " ".join(t for t in texte if t)
|
||||
for name, muster in _AGENTUREN:
|
||||
if muster.search(text):
|
||||
return name
|
||||
return None
|
||||
|
||||
|
||||
def gruppen_id(leitartikel_id):
|
||||
"""UUID aus dem Leitartikel, nicht aus der Mitgliedermenge.
|
||||
|
||||
Die Spezifikation leitet sie aus der sortierten Menge der raw_item_id ab.
|
||||
Das haelt nicht: erscheint dieselbe Meldung drei Stunden spaeter bei einem
|
||||
weiteren Haus, aendert sich die Menge und damit die UUID - eine Gruppe,
|
||||
die Phase 3 bereits kodiert hat, hiesse ploetzlich anders.
|
||||
|
||||
Der Leitartikel ist das aeltestgesehene Mitglied. Kommt ein Mitglied
|
||||
hinzu, ist es zwangslaeufig juenger und aendert das Minimum nicht. Die
|
||||
Kennung bleibt damit stabil, solange die Gruppe waechst. Sie wechselt
|
||||
nur, wenn zwei bestehende Gruppen verschmelzen - dann ist es richtig,
|
||||
dass sich etwas aendert.
|
||||
"""
|
||||
return str(uuid.uuid5(NAMENSRAUM, f"dublette:{leitartikel_id}"))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Laden
|
||||
# --------------------------------------------------------------------------
|
||||
class Item:
|
||||
__slots__ = ("id", "quelle", "titel", "teaser", "gesehen", "datum",
|
||||
"inhalt_hash", "basis", "text", "trigramme", "simhash", "woerter")
|
||||
|
||||
def __init__(self, id, quelle, titel, teaser, gesehen, datum, inhalt_hash,
|
||||
basis=None):
|
||||
self.id = id
|
||||
self.quelle = quelle
|
||||
self.titel = titel
|
||||
self.teaser = teaser
|
||||
self.gesehen = gesehen
|
||||
self.datum = datum
|
||||
self.inhalt_hash = inhalt_hash
|
||||
self.basis = basis or BASIS_VORGABE
|
||||
if self.basis == "titel":
|
||||
# Ressortkuerzel abtrennen: dieselbe Agenturzeile bekommt bei
|
||||
# jedem Haus ein anderes vorangestellt, oder gar keins.
|
||||
self.text = normalisierung.normalisiere(titel, ressort_abtrennen=True)
|
||||
else:
|
||||
self.text = normalisierung.normalisiere(titel, teaser)
|
||||
self.woerter = len(normalisierung.woerter(self.text))
|
||||
self.trigramme = normalisierung.trigramme(self.text)
|
||||
self.simhash = simhash(self.trigramme)
|
||||
|
||||
@property
|
||||
def vergleichbar(self):
|
||||
return self.woerter >= MINDESTWOERTER[self.basis]
|
||||
|
||||
|
||||
_SPALTEN = """id, quelle, titel, teaser, abgerufen_am,
|
||||
coalesce(pubdate, abgerufen_am)::date as datum, inhalt_hash"""
|
||||
|
||||
|
||||
def lade_bereich(conn, basis, von=None, bis=None):
|
||||
"""Items eines Datumsbereichs, ohne Grenzen der Gesamtbestand.
|
||||
|
||||
Das Datum kommt aus pubdate mit Rueckfall auf abgerufen_am: Phase 1
|
||||
vermerkt ausdruecklich, dass pubdate fehlen oder luegen kann.
|
||||
"""
|
||||
if von is None:
|
||||
zeilen = conn.execute(f"select {_SPALTEN} from raw_items order by id")
|
||||
else:
|
||||
zeilen = conn.execute(
|
||||
f"""select {_SPALTEN} from raw_items
|
||||
where coalesce(pubdate, abgerufen_am)::date between %s and %s
|
||||
order by id""", (von, bis))
|
||||
return [Item(*z, basis=basis) for z in zeilen]
|
||||
|
||||
|
||||
def _beruehrt_rand(ergebnis, lade_von, lade_bis):
|
||||
"""Reicht eine Gruppe bis an den Rand des geladenen Bereichs?
|
||||
|
||||
Dann kann jenseits davon noch ein Mitglied liegen, und die Gruppe ist
|
||||
unvollstaendig gerechnet.
|
||||
"""
|
||||
for gruppe in ergebnis.values():
|
||||
if len(gruppe["mitglieder"]) < 2:
|
||||
continue
|
||||
for m in gruppe["mitglieder"]:
|
||||
if m.datum <= lade_von or m.datum >= lade_bis:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def lade_und_gruppiere(conn, slot=None, fenster_tage=None, basis=None, **grenzen):
|
||||
"""Liefert (items, uf, ergebnis, schreibbereich).
|
||||
|
||||
Laden und Gruppieren haengen beim Slot-Lauf zusammen und sind darum eine
|
||||
Einheit: wieviel geladen werden muss, zeigt sich erst an den Gruppen.
|
||||
|
||||
Ohne slot der Gesamtbestand, dann ist alles beschreibbar.
|
||||
|
||||
Mit slot wird ein Rand mitgeladen und nur der innere Bereich beschrieben -
|
||||
dort, wo jedes Item seine vollstaendige Nachbarschaft gesehen hat. Eine
|
||||
unvollstaendig gerechnete Zeile darf eine vollstaendige nicht
|
||||
ueberschreiben.
|
||||
|
||||
Ein fester Rand genuegt dabei nicht. Uebernahmen bilden Ketten: dieselbe
|
||||
Schlagzeile erscheint am Montag bei einem, am Dienstag beim zweiten, am
|
||||
Mittwoch beim dritten Haus. Wo eine Kette am geladenen Rand abgeschnitten
|
||||
wird, faellt die Gruppe anders aus als im Backfill - und genau das darf
|
||||
nicht sein (Invariante 2 und 4). Also wird geladen, gruppiert, und wenn
|
||||
eine Gruppe bis an den Rand reicht, weiter geladen. In der Praxis laeuft
|
||||
das ein- bis zweimal.
|
||||
"""
|
||||
fenster_tage = EINSTELLUNGEN["fenster_tage"] if fenster_tage is None else fenster_tage
|
||||
|
||||
if slot is None:
|
||||
items = lade_bereich(conn, basis)
|
||||
uf, ergebnis = gruppiere(items, fenster_tage=fenster_tage, **grenzen)
|
||||
return items, uf, ergebnis, None
|
||||
|
||||
von, bis = conn.execute(
|
||||
"""select min(coalesce(pubdate, abgerufen_am)::date),
|
||||
max(coalesce(pubdate, abgerufen_am)::date)
|
||||
from raw_items where slot = %s""", (slot,)).fetchone()
|
||||
if von is None:
|
||||
return [], None, None, None
|
||||
|
||||
fenster = dt.timedelta(days=fenster_tage)
|
||||
schreibbereich = (von - fenster, bis + fenster)
|
||||
|
||||
rand = 2 * fenster_tage
|
||||
for versuch in range(RANDERWEITERUNGEN):
|
||||
lade_von = von - dt.timedelta(days=rand)
|
||||
lade_bis = bis + dt.timedelta(days=rand)
|
||||
items = lade_bereich(conn, basis, lade_von, lade_bis)
|
||||
uf, ergebnis = gruppiere(items, fenster_tage=fenster_tage, **grenzen)
|
||||
if not _beruehrt_rand(ergebnis, lade_von, lade_bis):
|
||||
return items, uf, ergebnis, schreibbereich
|
||||
rand += fenster_tage
|
||||
|
||||
# Reissleine: lieber den ganzen Bestand rechnen als etwas Falsches
|
||||
# schreiben. Bisher nie erreicht.
|
||||
items = lade_bereich(conn, basis)
|
||||
uf, ergebnis = gruppiere(items, fenster_tage=fenster_tage, **grenzen)
|
||||
return items, uf, ergebnis, schreibbereich
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Gruppieren
|
||||
# --------------------------------------------------------------------------
|
||||
def paare(items, fenster_tage=EINSTELLUNGEN["fenster_tage"],
|
||||
hamming_max=EINSTELLUNGEN["hamming_max"],
|
||||
jaccard_min=EINSTELLUNGEN["jaccard_min"]):
|
||||
"""Bestaetigte Paare finden.
|
||||
|
||||
Erst Kandidaten ueber gemeinsame SimHash-Baender, dann die drei Pruefungen
|
||||
aus der Spezifikation: Datumsnaehe, Hamming-Distanz, Jaccard-Wert.
|
||||
"""
|
||||
nach_band = {}
|
||||
brauchbar = [i for i in items if i.vergleichbar]
|
||||
for item in brauchbar:
|
||||
for band in baender(item.simhash):
|
||||
nach_band.setdefault(band, []).append(item)
|
||||
|
||||
gesehen = set()
|
||||
fenster = dt.timedelta(days=fenster_tage)
|
||||
for eimer in nach_band.values():
|
||||
if len(eimer) < 2:
|
||||
continue
|
||||
for a, b in itertools.combinations(eimer, 2):
|
||||
schluessel = (a.id, b.id) if a.id < b.id else (b.id, a.id)
|
||||
if schluessel in gesehen:
|
||||
continue
|
||||
gesehen.add(schluessel)
|
||||
if abs(a.datum - b.datum) > fenster:
|
||||
continue
|
||||
# Dasselbe Haus an verschiedenen Tagen ist keine Uebernahme,
|
||||
# sondern ein wiederkehrendes Format: "tagesschau in 100
|
||||
# Sekunden", "Wetter", "Die Nachrichten". Gleicher Text, anderer
|
||||
# Vorgang - kein Textverfahren kann die auseinanderhalten, also
|
||||
# muss die Regel es tun.
|
||||
#
|
||||
# Die Regel sitzt bewusst hier auf Paarebene und nicht auf der
|
||||
# fertigen Gruppe: sie haengt allein an den beiden Items und
|
||||
# faellt damit gleich aus, egal wieviel Bestand geladen ist.
|
||||
# Eine Regel auf Gruppenebene waere fensterabhaengig - und dann
|
||||
# lieferte ein Slot-Lauf andere Gruppen als ein Backfill.
|
||||
if a.quelle == b.quelle and a.datum != b.datum:
|
||||
continue
|
||||
if hamming(a.simhash, b.simhash) > hamming_max:
|
||||
continue
|
||||
# Bestaetigung. Auf dem Bestand vom 4.-7.9.2026 verwirft sie
|
||||
# kein einziges Paar mehr, das die Hamming-Grenze passiert hat -
|
||||
# sie ist trotzdem die eigentliche Praezisionszusage. SimHash
|
||||
# kann kollidieren, Jaccard nicht.
|
||||
j = jaccard(a.trigramme, b.trigramme)
|
||||
if j < jaccard_min:
|
||||
continue
|
||||
yield a, b, j
|
||||
|
||||
|
||||
def _zerlege_ketten(gruppe, fenster_tage):
|
||||
"""Gruppen aufbrechen, die nur ueber eine Kette zusammenhaengen.
|
||||
|
||||
Union-Find ist transitiv, das Datumsfenster ist es nicht: A passt zu B,
|
||||
B zu C, und schon haengt A mit C zusammen, obwohl mehr als ein Tag
|
||||
dazwischen liegt. Eine Uebernahme ist ein Vorgang, der bei mehreren
|
||||
Haeusern gleichzeitig erscheint - keine Schlagzeile, die drei Tage lang
|
||||
weitergereicht wird. Eine Gruppe darf darum nicht mehr Tage umspannen
|
||||
als das Fenster.
|
||||
|
||||
Auf dem Bestand vom 4.-7.9.2026 trifft das genau eine von 135
|
||||
Komponenten: dieselbe Newsblog-Zeile bei handelsblatt (4.9.), sz (5.9.)
|
||||
und tagesschau (6.9.). Der Schnitt ist also kein Massenphaenomen,
|
||||
sondern verhindert, dass Phase 3 den Stand vom Mittwoch mit der
|
||||
Kodierung vom Montag versieht.
|
||||
|
||||
Getrennt wird nach Datum, damit das Ergebnis nicht von der
|
||||
Ladereihenfolge abhaengt.
|
||||
"""
|
||||
spanne = dt.timedelta(days=fenster_tage)
|
||||
geordnet = sorted(gruppe, key=lambda i: (i.datum, i.gesehen, i.id))
|
||||
teile, aktuell = [], [geordnet[0]]
|
||||
for item in geordnet[1:]:
|
||||
if item.datum - aktuell[0].datum > spanne:
|
||||
teile.append(aktuell)
|
||||
aktuell = [item]
|
||||
else:
|
||||
aktuell.append(item)
|
||||
teile.append(aktuell)
|
||||
return teile
|
||||
|
||||
|
||||
def gruppiere(items, fenster_tage=EINSTELLUNGEN["fenster_tage"], **grenzen):
|
||||
"""Union-Find ueber die bestaetigten Paare.
|
||||
|
||||
Liefert je Item die Gruppe, ihre Mitglieder, den Leitartikel und den
|
||||
kleinsten Jaccard-Wert, der die Gruppe zusammenhaelt.
|
||||
"""
|
||||
uf = UnionFind()
|
||||
schwaechste = {}
|
||||
for a, b, j in paare(items, fenster_tage=fenster_tage, **grenzen):
|
||||
uf.vereinige(a.id, b.id)
|
||||
schwaechste[(a.id, b.id)] = j
|
||||
|
||||
roh = {}
|
||||
for item in items:
|
||||
roh.setdefault(uf.finde(item.id), []).append(item)
|
||||
|
||||
# Ketten aufbrechen und die Zugehoerigkeit neu festschreiben. Danach ist
|
||||
# uf.finde() wieder die Wahrheit ueber die Gruppen.
|
||||
mitglieder = {}
|
||||
for gruppe in roh.values():
|
||||
for teil in _zerlege_ketten(gruppe, fenster_tage):
|
||||
wurzel = min(i.id for i in teil)
|
||||
for item in teil:
|
||||
uf.eltern[item.id] = wurzel
|
||||
mitglieder[wurzel] = teil
|
||||
|
||||
ergebnis = {}
|
||||
for wurzel, gruppe in mitglieder.items():
|
||||
# Leitartikel: aeltestes gesehenes Item, bei Gleichstand die kleinste
|
||||
# id. Rein deterministisch, unabhaengig von der Ladereihenfolge.
|
||||
leit = min(gruppe, key=lambda i: (i.gesehen, i.id))
|
||||
ids = {i.id for i in gruppe}
|
||||
werte = [j for (x, y), j in schwaechste.items() if x in ids and y in ids]
|
||||
ergebnis[wurzel] = {
|
||||
"mitglieder": gruppe,
|
||||
"leitartikel": leit,
|
||||
"jaccard_min": round(min(werte), 4) if werte else None,
|
||||
}
|
||||
return uf, ergebnis
|
||||
|
||||
|
||||
def nutzlast(item, gruppe):
|
||||
leit = gruppe["leitartikel"]
|
||||
groesse = len(gruppe["mitglieder"])
|
||||
return {
|
||||
"gruppe": gruppen_id(leit.id),
|
||||
"gruppengroesse": groesse,
|
||||
"leitartikel": leit.id,
|
||||
"simhash": f"{item.simhash:016x}",
|
||||
"jaccard_min": gruppe["jaccard_min"],
|
||||
"agentur_vermutet": agentur(item.titel, item.teaser),
|
||||
# Basis und Stand des Texts, gegen den kodiert wurde. Ohne das laesst
|
||||
# sich spaeter nicht sagen, ob eine Kodierung noch zum Item passt:
|
||||
# raw_items haelt den *aktuellen* Titel, Redaktionen aendern ihn nach
|
||||
# (in diesem Bestand 3 von 10 Items), und dann ist der Vergleich von
|
||||
# gestern gegen einen Text gelaufen, den es so nicht mehr gibt.
|
||||
"basis": item.basis,
|
||||
"basis_hash": item.inhalt_hash,
|
||||
"woerter": item.woerter,
|
||||
"vergleichbar": item.vergleichbar,
|
||||
# Eine Gruppe aus einem einzigen Haus ist keine Uebernahme, sondern
|
||||
# eine Wiederholung im selben Feed. Fuer Phase 3 ist der Unterschied
|
||||
# wesentlich: bei einer Uebernahme erben die Mitglieder die Kodierung
|
||||
# des Leitartikels, eine Wiederholung ist derselbe Vorgang zweimal.
|
||||
"haeuser": len({i.quelle for i in gruppe["mitglieder"]}),
|
||||
}
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Schreiben
|
||||
# --------------------------------------------------------------------------
|
||||
def schreibe(conn, items, ergebnis, uf, coder_version, trocken=False, schreibbereich=None):
|
||||
"""Codings ablegen, aber nur wo sich etwas geaendert hat.
|
||||
|
||||
Jedes verarbeitete Item bekommt eine Zeile, auch ein Einzelstueck: sonst
|
||||
liesse sich "keine Dublette" nicht von "nicht verarbeitet" unterscheiden.
|
||||
|
||||
Unveraenderte Zeilen werden nicht angefasst. Ein Slot-Lauf rechnet ueber
|
||||
das ganze Datumsfenster, und ohne diesen Vergleich schriebe er bei jedem
|
||||
Durchgang einige tausend Zeilen neu - mit neuem kodiert_am, ohne dass
|
||||
sich etwas geaendert haette.
|
||||
"""
|
||||
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))
|
||||
}
|
||||
|
||||
if schreibbereich is None:
|
||||
zu_schreiben = items
|
||||
else:
|
||||
von, bis = schreibbereich
|
||||
zu_schreiben = [i for i in items if von <= i.datum <= bis]
|
||||
|
||||
neu = geaendert = unveraendert = 0
|
||||
stapel = []
|
||||
for item in zu_schreiben:
|
||||
gruppe = ergebnis[uf.finde(item.id)]
|
||||
last = nutzlast(item, gruppe)
|
||||
alt = vorhanden.get(item.id)
|
||||
if alt == last:
|
||||
unveraendert += 1
|
||||
continue
|
||||
if alt is None:
|
||||
neu += 1
|
||||
else:
|
||||
geaendert += 1
|
||||
stapel.append((item.id, coder_version, EBENE, last["jaccard_min"], Jsonb(last)))
|
||||
|
||||
if stapel and not trocken:
|
||||
with conn.cursor() as cur:
|
||||
cur.executemany(
|
||||
"""insert into codings (raw_item_id, coder_version, ebene, konfidenz, nutzlast)
|
||||
values (%s, %s, %s, %s, %s)
|
||||
on conflict (raw_item_id, coder_version, ebene)
|
||||
do update set nutzlast = excluded.nutzlast,
|
||||
konfidenz = excluded.konfidenz,
|
||||
kodiert_am = now()""",
|
||||
stapel)
|
||||
return {"neu": neu, "geaendert": geaendert, "unveraendert": unveraendert,
|
||||
"ausserhalb": len(items) - len(zu_schreiben)}
|
||||
|
||||
|
||||
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))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Bericht und Stichprobe
|
||||
# --------------------------------------------------------------------------
|
||||
def bericht(items, ergebnis):
|
||||
gruppen = [g for g in ergebnis.values() if len(g["mitglieder"]) > 1]
|
||||
in_gruppen = sum(len(g["mitglieder"]) for g in gruppen)
|
||||
ueber = [g for g in gruppen if len({i.quelle for i in g["mitglieder"]}) > 1]
|
||||
kurz = sum(1 for i in items if not i.vergleichbar)
|
||||
mindest = MINDESTWOERTER[items[0].basis] if items else 0
|
||||
print(f" Items: {len(items)}")
|
||||
print(f" zu kurz: {kurz} (< {mindest} Woerter, nicht verglichen)")
|
||||
print(f" Uebernahmegruppen:{len(gruppen):>4} mit {in_gruppen} Mitgliedern")
|
||||
print(f" davon ueber mehrere Haeuser: {len(ueber)} <- der eigentliche Zweck")
|
||||
if gruppen:
|
||||
groessen = sorted((len(g["mitglieder"]) for g in gruppen), reverse=True)
|
||||
print(f" groesste Gruppe: {groessen[0]}")
|
||||
for g in sorted(gruppen, key=lambda g: -len(g["mitglieder"]))[:5]:
|
||||
haeuser = sorted({i.quelle for i in g["mitglieder"]})
|
||||
leit = g["leitartikel"]
|
||||
print(f" [{len(g['mitglieder'])}] j={g['jaccard_min']} {', '.join(haeuser)}")
|
||||
print(f" {leit.titel[:88]}")
|
||||
|
||||
|
||||
def stichprobe(ergebnis, anzahl, pfad):
|
||||
"""Paare zum Nachprüfen von Hand ausschreiben.
|
||||
|
||||
Die Spezifikation verlangt Praezision >= 0.98 auf 200 handgeprueften
|
||||
Paaren. Pruefen kann das nur ein Mensch; diese Datei ist die Vorlage
|
||||
dafuer. Sie enthaelt Titel im Klartext und gehoert deshalb nicht in die
|
||||
Datenbank, sondern neben sie.
|
||||
"""
|
||||
zeilen = []
|
||||
for gruppe in ergebnis.values():
|
||||
m = gruppe["mitglieder"]
|
||||
if len(m) < 2:
|
||||
continue
|
||||
for a, b in itertools.combinations(sorted(m, key=lambda i: i.id), 2):
|
||||
zeilen.append((jaccard(a.trigramme, b.trigramme),
|
||||
hamming(a.simhash, b.simhash), a, b))
|
||||
zeilen.sort(key=lambda z: (z[0], z[3].id)) # schwaechste zuerst - dort sitzen die Fehler
|
||||
with open(pfad, "w", encoding="utf-8") as f:
|
||||
f.write("# Dubletten-Stichprobe. Spalte 'urteil' von Hand fuellen: j = Dublette, n = keine.\n")
|
||||
f.write("# Absteigend nach Zweifel sortiert: die obersten Paare sind die knappsten.\n\n")
|
||||
for j, h, a, b in zeilen[:anzahl]:
|
||||
f.write(f"urteil=_ jaccard={j:.3f} hamming={h}\n")
|
||||
f.write(f" {a.id:>7} {a.quelle:<14}{a.titel}\n")
|
||||
f.write(f" {b.id:>7} {b.quelle:<14}{b.titel}\n\n")
|
||||
return min(len(zeilen), anzahl), len(zeilen)
|
||||
|
||||
|
||||
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("--basis", choices=("titel", "titel+teaser"), default=BASIS_VORGABE,
|
||||
help=f"Vergleichsbasis (Vorgabe {BASIS_VORGABE})")
|
||||
p.add_argument("--jaccard", type=float, default=EINSTELLUNGEN["jaccard_min"],
|
||||
help="Bestaetigungsschwelle (Vorgabe %(default)s)")
|
||||
p.add_argument("--hamming", type=int, default=EINSTELLUNGEN["hamming_max"],
|
||||
help="hoechste Hamming-Distanz fuer Kandidaten (Vorgabe %(default)s)")
|
||||
p.add_argument("--fenster", type=int, default=1, help="Datumsfenster in Tagen (Vorgabe 1)")
|
||||
p.add_argument("--trocken", action="store_true", help="rechnen, nichts schreiben")
|
||||
p.add_argument("--stichprobe", type=int, metavar="N", help="N Paare zum Nachpruefen ausschreiben")
|
||||
p.add_argument("--stichprobe-datei", default="stichprobe-dubletten.txt")
|
||||
a = p.parse_args(argv)
|
||||
|
||||
if not a.backfill and not a.slot:
|
||||
p.error("--backfill oder --slot angeben")
|
||||
|
||||
slot = dt.datetime.fromisoformat(a.slot) if a.slot else None
|
||||
begonnen = time.monotonic()
|
||||
|
||||
einstellungen = dict(EINSTELLUNGEN,
|
||||
basis=a.basis,
|
||||
jaccard_min=a.jaccard,
|
||||
hamming_max=a.hamming,
|
||||
fenster_tage=a.fenster,
|
||||
mindestwoerter=MINDESTWOERTER[a.basis])
|
||||
|
||||
with verbindung() as conn:
|
||||
# Die Einstellungen gehen in den Hash ein: eine andere Basis oder eine
|
||||
# andere Schwelle ist eine andere Kodierung und darf die vorhandene
|
||||
# nicht ueberschreiben.
|
||||
v = versionsmodul.eintragen(
|
||||
conn, komponenten=versionsmodul.komponenten(dubletten=einstellungen))
|
||||
print(f"coder_version {v} basis={a.basis} j>={a.jaccard} h<={a.hamming}"
|
||||
f"{' (trocken)' if a.trocken else ''}")
|
||||
|
||||
items, uf, ergebnis, schreibbereich = lade_und_gruppiere(
|
||||
conn, slot=slot, fenster_tage=a.fenster, basis=a.basis,
|
||||
jaccard_min=a.jaccard, hamming_max=a.hamming)
|
||||
if not items:
|
||||
print(" nichts zu tun")
|
||||
return 0
|
||||
bericht(items, ergebnis)
|
||||
|
||||
zahlen = schreibe(conn, items, ergebnis, uf, v, trocken=a.trocken,
|
||||
schreibbereich=schreibbereich)
|
||||
dauer = int((time.monotonic() - begonnen) * 1000)
|
||||
print(f" codings: neu {zahlen['neu']}, geaendert {zahlen['geaendert']},"
|
||||
f" unveraendert {zahlen['unveraendert']}"
|
||||
+ (f", Kranz uebergangen {zahlen['ausserhalb']}" if zahlen["ausserhalb"] else ""))
|
||||
print(f" Dauer: {dauer} ms")
|
||||
|
||||
if a.stichprobe:
|
||||
n, gesamt = stichprobe(ergebnis, a.stichprobe, a.stichprobe_datei)
|
||||
print(f" Stichprobe: {n} von {gesamt} Paaren nach {a.stichprobe_datei}")
|
||||
|
||||
if not a.trocken:
|
||||
protokolliere(conn, slot, v, dauer, len(items),
|
||||
zahlen["neu"] + zahlen["geaendert"])
|
||||
conn.commit()
|
||||
else:
|
||||
conn.rollback()
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user