Initiales Ingest-Repo (aus gdelt-bot/ingest überführt)
Poller, Schema, Quadlet-Units und Quellenregister aus dem Monorepo gdelt-bot in ein eigenes Repo ausgelagert. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VzNVf1KaHPZ1fzJT2jd18j
This commit is contained in:
@@ -0,0 +1,435 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Ingest-Poller fuer den deutschen Ereignisdatensatz.
|
||||
|
||||
Liest die Quellen aus sources.toml, holt alle 15 Minuten die Feeds und legt
|
||||
jedes neue Item roh ab. Bewusst ohne Kodierung: was hier landet, soll sich
|
||||
beliebig oft neu kodieren lassen.
|
||||
|
||||
Zwei Senken:
|
||||
--sink postgres://... Bebop, die eigentliche Datenbank
|
||||
--sink jsonl:/pfad/dir Pi 5, Schatten-Ingest als gzip-JSONL pro Slot
|
||||
|
||||
Aufrufe:
|
||||
poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt
|
||||
poller.py --sink "$DATABASE_URL" --once # ein Durchlauf
|
||||
poller.py --check-robots # Lizenzstatus nachpruefen
|
||||
"""
|
||||
import argparse
|
||||
import datetime as dt
|
||||
import gzip
|
||||
import hashlib
|
||||
import html
|
||||
import json
|
||||
import os
|
||||
import pathlib
|
||||
import re
|
||||
import signal
|
||||
import sys
|
||||
import time
|
||||
import tomllib
|
||||
import sqlite3
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
import urllib.robotparser
|
||||
import xml.etree.ElementTree as ET
|
||||
from email.utils import parsedate_to_datetime
|
||||
|
||||
HIER = pathlib.Path(__file__).parent
|
||||
USER_AGENT = os.environ.get(
|
||||
"INGEST_USER_AGENT",
|
||||
"lagebild-ingest/0.1 (+https://example.org/lagebild; kontakt@example.org)",
|
||||
)
|
||||
TAKT = 15 * 60 # Slotlaenge in Sekunden, wie bei GDELT
|
||||
VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf
|
||||
NS = {
|
||||
"a": "http://www.w3.org/2005/Atom",
|
||||
"dc": "http://purl.org/dc/elements/1.1/",
|
||||
"content": "http://purl.org/rss/1.0/modules/content/",
|
||||
}
|
||||
|
||||
_lauf = True
|
||||
|
||||
|
||||
def _stop(signum, frame):
|
||||
global _lauf
|
||||
_lauf = False
|
||||
print(f"[{_jetzt():%H:%M:%S}] Signal {signum}, beende nach diesem Durchlauf",
|
||||
flush=True)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Hilfsfunktionen
|
||||
# --------------------------------------------------------------------------
|
||||
def _jetzt():
|
||||
return dt.datetime.now(dt.timezone.utc)
|
||||
|
||||
|
||||
def slot_von(ts):
|
||||
"""Auf 15 Minuten abrunden, wie GDELTs Paketgrenzen."""
|
||||
return ts.replace(second=0, microsecond=0,
|
||||
minute=(ts.minute // 15) * 15)
|
||||
|
||||
|
||||
def norm_titel(t):
|
||||
"""Wie build_site.norm_title: Basis fuer das Trigramm-Clustering."""
|
||||
t = (t or "").lower().split(" - ")[0].split(" | ")[0]
|
||||
# \w haelt Umlaute und ss, ohne die es im Deutschen nicht geht
|
||||
return re.sub(r"\W+", " ", t, flags=re.UNICODE).strip()
|
||||
|
||||
|
||||
_TAG_RE = re.compile(r"<[^>]+>")
|
||||
|
||||
|
||||
def klartext(s):
|
||||
"""HTML aus Teasern entfernen. Feeds kodieren teils doppelt."""
|
||||
if not s:
|
||||
return None
|
||||
s = html.unescape(html.unescape(s))
|
||||
s = _TAG_RE.sub(" ", s)
|
||||
return re.sub(r"\s+", " ", s).strip() or None
|
||||
|
||||
|
||||
def datum(s):
|
||||
s = (s or "").strip()
|
||||
if not s:
|
||||
return None
|
||||
try:
|
||||
d = parsedate_to_datetime(s)
|
||||
except Exception:
|
||||
try:
|
||||
d = dt.datetime.fromisoformat(s.replace("Z", "+00:00"))
|
||||
except Exception:
|
||||
return None
|
||||
if d.tzinfo is None:
|
||||
d = d.replace(tzinfo=dt.timezone.utc)
|
||||
return d.astimezone(dt.timezone.utc)
|
||||
|
||||
|
||||
def _text(el):
|
||||
return (el.text or "").strip() if el is not None and el.text else None
|
||||
|
||||
|
||||
def _erst(item, pfade):
|
||||
for p in pfade:
|
||||
el = item.find(p, NS) if ":" in p else item.find(p)
|
||||
v = _text(el)
|
||||
if v:
|
||||
return v
|
||||
return None
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Feed-Abruf und -Parsing
|
||||
# --------------------------------------------------------------------------
|
||||
def hole(url, etag=None, last_modified=None, timeout=30):
|
||||
"""Bedingter GET. Gibt (status, bytes|None, etag, last_modified) zurueck."""
|
||||
kopf = {"User-Agent": USER_AGENT, "Accept-Encoding": "gzip"}
|
||||
if etag:
|
||||
kopf["If-None-Match"] = etag
|
||||
if last_modified:
|
||||
kopf["If-Modified-Since"] = last_modified
|
||||
req = urllib.request.Request(url, headers=kopf)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=timeout) as r:
|
||||
roh = r.read()
|
||||
if r.headers.get("Content-Encoding") == "gzip":
|
||||
roh = gzip.decompress(roh)
|
||||
return (r.status, roh,
|
||||
r.headers.get("ETag"), r.headers.get("Last-Modified"))
|
||||
except urllib.error.HTTPError as err:
|
||||
if err.code == 304:
|
||||
return 304, None, etag, last_modified
|
||||
raise
|
||||
|
||||
|
||||
def parse_feed(roh):
|
||||
"""RSS 2.0, RDF und Atom auf eine gemeinsame Form bringen."""
|
||||
wurzel = ET.fromstring(roh)
|
||||
eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS)
|
||||
aus = []
|
||||
for it in eintraege:
|
||||
titel = _erst(it, ["title", "a:title"])
|
||||
if not titel:
|
||||
continue # ohne Titel ist ein Item fuer uns wertlos
|
||||
|
||||
link = _erst(it, ["link", "a:id"])
|
||||
if not link:
|
||||
for le in it.findall("a:link", NS):
|
||||
if le.get("rel") in (None, "alternate") and le.get("href"):
|
||||
link = le.get("href")
|
||||
break
|
||||
|
||||
guid = _erst(it, ["guid", "a:id"]) or link
|
||||
if not guid:
|
||||
guid = hashlib.sha256(titel.encode()).hexdigest()[:32]
|
||||
|
||||
aus.append({
|
||||
"titel": klartext(titel),
|
||||
"url": link,
|
||||
"guid": guid,
|
||||
"teaser": klartext(_erst(it, ["description", "content:encoded",
|
||||
"a:summary", "a:content"])),
|
||||
"autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]),
|
||||
"pubdate": datum(_erst(it, ["pubDate", "dc:date",
|
||||
"a:published", "a:updated"])),
|
||||
})
|
||||
return aus
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Senken
|
||||
# --------------------------------------------------------------------------
|
||||
class JsonlSenke:
|
||||
"""Schatten-Ingest fuer den Pi: gzip-JSONL, ein Verzeichnis pro Tag.
|
||||
|
||||
Kein Schema, keine Abfragen - nur die Garantie, dass die Rohdaten noch
|
||||
da sind, wenn Bebop einen Tag lang stillstand. Die SQLite daneben haelt
|
||||
nur die gesehenen Schluessel, damit nicht bei jedem Slot der komplette
|
||||
Feed erneut geschrieben wird.
|
||||
"""
|
||||
|
||||
def __init__(self, ziel):
|
||||
self.ziel = pathlib.Path(ziel)
|
||||
self.ziel.mkdir(parents=True, exist_ok=True)
|
||||
self.db = sqlite3.connect(self.ziel / "gesehen.sqlite")
|
||||
self.db.execute("create table if not exists gesehen ("
|
||||
"quelle text, guid text, primary key (quelle, guid))")
|
||||
self.db.commit()
|
||||
|
||||
def sync_quellen(self, quellen):
|
||||
(self.ziel / "sources.snapshot.json").write_text(
|
||||
json.dumps(quellen, indent=2, ensure_ascii=False, default=str))
|
||||
|
||||
def zustand(self, qid):
|
||||
return None, None
|
||||
|
||||
def schreibe(self, quelle, items, slot):
|
||||
frisch = []
|
||||
for it in items:
|
||||
cur = self.db.execute(
|
||||
"insert or ignore into gesehen (quelle, guid) values (?, ?)",
|
||||
(quelle["id"], it["guid"]))
|
||||
if cur.rowcount:
|
||||
frisch.append(it)
|
||||
self.db.commit()
|
||||
if not frisch:
|
||||
return 0
|
||||
tag = self.ziel / slot.strftime("%Y%m%d")
|
||||
tag.mkdir(exist_ok=True)
|
||||
pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz")
|
||||
# gzip-Member lassen sich anhaengen, das Ergebnis bleibt gueltiges gzip
|
||||
with gzip.open(pfad, "at", encoding="utf-8") as f:
|
||||
for it in frisch:
|
||||
f.write(json.dumps(
|
||||
{"quelle": quelle["id"], "slot": slot.isoformat(),
|
||||
"abgerufen_am": _jetzt().isoformat(), **it},
|
||||
ensure_ascii=False, default=str) + "\n")
|
||||
return len(frisch)
|
||||
|
||||
def notiere(self, *a, **kw):
|
||||
pass
|
||||
|
||||
def close(self):
|
||||
self.db.close()
|
||||
|
||||
|
||||
class PostgresSenke:
|
||||
def __init__(self, dsn):
|
||||
import psycopg
|
||||
self.conn = psycopg.connect(dsn, autocommit=False)
|
||||
|
||||
def sync_quellen(self, quellen):
|
||||
with self.conn.cursor() as cur:
|
||||
for q in quellen:
|
||||
cur.execute("""
|
||||
insert into quellen
|
||||
(id, name, typ, url, land, traeger, ressort,
|
||||
lizenz_status, robots_geprueft_am, aktiv, notiz)
|
||||
values (%(id)s, %(name)s, %(typ)s, %(url)s, %(land)s,
|
||||
%(traeger)s, %(ressort)s, %(lizenz_status)s,
|
||||
%(robots_geprueft_am)s, %(aktiv)s, %(notiz)s)
|
||||
on conflict (id) do update set
|
||||
name = excluded.name, typ = excluded.typ,
|
||||
url = excluded.url, land = excluded.land,
|
||||
traeger = excluded.traeger, ressort = excluded.ressort,
|
||||
lizenz_status = excluded.lizenz_status,
|
||||
robots_geprueft_am = excluded.robots_geprueft_am,
|
||||
aktiv = excluded.aktiv, notiz = excluded.notiz
|
||||
""", {k: q.get(k) for k in
|
||||
("id", "name", "typ", "url", "land", "traeger", "ressort",
|
||||
"lizenz_status", "robots_geprueft_am", "aktiv", "notiz")})
|
||||
self.conn.commit()
|
||||
|
||||
def zustand(self, qid):
|
||||
with self.conn.cursor() as cur:
|
||||
cur.execute("select etag, last_modified from quellen where id = %s",
|
||||
(qid,))
|
||||
r = cur.fetchone()
|
||||
return (r[0], r[1]) if r else (None, None)
|
||||
|
||||
def schreibe(self, quelle, items, slot):
|
||||
neu = 0
|
||||
with self.conn.cursor() as cur:
|
||||
for it in items:
|
||||
cur.execute("""
|
||||
insert into raw_items
|
||||
(quelle, item_key, url, titel, titel_norm, teaser,
|
||||
autor, pubdate, slot, roh)
|
||||
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
on conflict (quelle, item_key) do nothing
|
||||
returning id
|
||||
""", (quelle["id"], it["guid"], it["url"], it["titel"],
|
||||
norm_titel(it["titel"]), it["teaser"], it["autor"],
|
||||
it["pubdate"], slot,
|
||||
json.dumps(it, ensure_ascii=False, default=str)))
|
||||
if cur.fetchone():
|
||||
neu += 1
|
||||
self.conn.commit()
|
||||
return neu
|
||||
|
||||
def merke_zustand(self, qid, etag, last_modified, geaendert):
|
||||
with self.conn.cursor() as cur:
|
||||
cur.execute("""
|
||||
update quellen set etag = %s, last_modified = %s,
|
||||
zuletzt_geprueft = now(),
|
||||
zuletzt_geaendert = case when %s then now()
|
||||
else zuletzt_geaendert end
|
||||
where id = %s
|
||||
""", (etag, last_modified, geaendert, qid))
|
||||
self.conn.commit()
|
||||
|
||||
def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, fehler):
|
||||
with self.conn.cursor() as cur:
|
||||
cur.execute("""
|
||||
insert into poll_laeufe
|
||||
(quelle, slot, dauer_ms, http_status,
|
||||
items_gesehen, items_neu, fehler)
|
||||
values (%s, %s, %s, %s, %s, %s, %s)
|
||||
""", (quelle, slot, dauer_ms, status, gesehen, neu, fehler))
|
||||
self.conn.commit()
|
||||
|
||||
def close(self):
|
||||
self.conn.close()
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Durchlauf
|
||||
# --------------------------------------------------------------------------
|
||||
def lade_quellen(pfad):
|
||||
with open(pfad, "rb") as f:
|
||||
return tomllib.load(f)["quelle"]
|
||||
|
||||
|
||||
def durchlauf(quellen, senke, slot):
|
||||
ges_neu = ges_gesehen = 0
|
||||
for q in quellen:
|
||||
if q["typ"] != "rss":
|
||||
continue
|
||||
t0 = time.monotonic()
|
||||
etag, lm = senke.zustand(q["id"])
|
||||
status = gesehen = neu = 0
|
||||
fehler = None
|
||||
try:
|
||||
status, roh, n_etag, n_lm = hole(q["url"], etag, lm)
|
||||
if status == 304:
|
||||
pass
|
||||
else:
|
||||
items = parse_feed(roh)
|
||||
gesehen = len(items)
|
||||
neu = senke.schreibe(q, items, slot)
|
||||
if hasattr(senke, "merke_zustand"):
|
||||
senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0)
|
||||
except Exception as e: # eine Quelle darf nicht
|
||||
fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen
|
||||
status = getattr(e, "code", None)
|
||||
dauer = int((time.monotonic() - t0) * 1000)
|
||||
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, fehler)
|
||||
ges_neu += neu
|
||||
ges_gesehen += gesehen
|
||||
marke = "!" if fehler else (" " if status != 304 else ".")
|
||||
print(f" {marke} {q['id']:14} {str(status):>4} "
|
||||
f"{gesehen:4} gesehen {neu:4} neu {dauer:5} ms"
|
||||
+ (f" {fehler}" if fehler else ""), flush=True)
|
||||
return ges_gesehen, ges_neu
|
||||
|
||||
|
||||
def check_robots(quellen):
|
||||
"""Lizenzstatus gegen die aktuelle robots.txt nachpruefen."""
|
||||
print(f"{'quelle':14} {'im Register':12} {'robots.txt':12} Abweichung")
|
||||
abweichungen = 0
|
||||
for q in quellen:
|
||||
if q["typ"] != "rss":
|
||||
continue
|
||||
p = urllib.parse.urlparse(q["url"])
|
||||
rp = urllib.robotparser.RobotFileParser()
|
||||
rp.set_url(f"{p.scheme}://{p.netloc}/robots.txt")
|
||||
try:
|
||||
rp.read()
|
||||
ist = "offen" if rp.can_fetch(USER_AGENT, q["url"]) else "gesperrt"
|
||||
except Exception as e:
|
||||
ist = "fehler"
|
||||
soll = q["lizenz_status"]
|
||||
# 'ungeprueft' im Register ist keine Abweichung, sondern eine Aufgabe
|
||||
weicht = soll != "ungeprueft" and ist != soll
|
||||
abweichungen += weicht
|
||||
print(f"{q['id']:14} {soll:12} {ist:12} {'<<< PRUEFEN' if weicht else ''}")
|
||||
print(f"\n{abweichungen} Abweichung(en). "
|
||||
f"Bei Treffern sources.toml anpassen, nicht den Code.")
|
||||
return abweichungen
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description=__doc__,
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter)
|
||||
ap.add_argument("--sources", default=str(HIER / "sources.toml"))
|
||||
ap.add_argument("--sink", default=os.environ.get("DATABASE_URL"),
|
||||
help="postgres://... oder jsonl:/pfad/zum/verzeichnis")
|
||||
ap.add_argument("--once", action="store_true", help="ein Durchlauf, dann Ende")
|
||||
ap.add_argument("--check-robots", action="store_true",
|
||||
help="Lizenzstatus nachpruefen, nichts abrufen")
|
||||
args = ap.parse_args()
|
||||
|
||||
alle = lade_quellen(args.sources)
|
||||
|
||||
if args.check_robots:
|
||||
return 1 if check_robots(alle) else 0
|
||||
|
||||
if not args.sink:
|
||||
ap.error("--sink oder DATABASE_URL noetig")
|
||||
|
||||
aktiv = [q for q in alle if q.get("aktiv")]
|
||||
gesperrt = [q for q in alle if q.get("lizenz_status") == "gesperrt"]
|
||||
if any(q.get("aktiv") for q in gesperrt):
|
||||
sys.exit("FEHLER: gesperrte Quelle ist aktiv geschaltet")
|
||||
|
||||
if args.sink.startswith("jsonl:"):
|
||||
senke = JsonlSenke(args.sink[len("jsonl:"):])
|
||||
else:
|
||||
senke = PostgresSenke(args.sink)
|
||||
senke.sync_quellen(alle)
|
||||
|
||||
print(f"{len(aktiv)} aktive Quellen, {len(gesperrt)} gesperrt, "
|
||||
f"Senke: {args.sink.split('@')[-1]}", flush=True)
|
||||
|
||||
signal.signal(signal.SIGTERM, _stop)
|
||||
signal.signal(signal.SIGINT, _stop)
|
||||
|
||||
try:
|
||||
while _lauf:
|
||||
slot = slot_von(_jetzt())
|
||||
print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True)
|
||||
gesehen, neu = durchlauf(aktiv, senke, slot)
|
||||
print(f" = {gesehen} gesehen, {neu} neu", flush=True)
|
||||
if args.once:
|
||||
break
|
||||
ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ)
|
||||
while _lauf and _jetzt() < ziel:
|
||||
time.sleep(min(5, (ziel - _jetzt()).total_seconds()))
|
||||
finally:
|
||||
senke.close()
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user