660 lines
26 KiB
Python
Executable File
660 lines
26 KiB
Python
Executable File
#!/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())
|