Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5a9bbe6e92 | ||
|
|
4983662514 |
@@ -8,7 +8,8 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen.
|
|||||||
| Datei | Zweck |
|
| Datei | Zweck |
|
||||||
|---|---|
|
|---|---|
|
||||||
| `sources.toml` | Quellenregister, einzige Wahrheit ueber die Quellen |
|
| `sources.toml` | Quellenregister, einzige Wahrheit ueber die Quellen |
|
||||||
| `schema.sql` | Postgres-Schema (`quellen`, `raw_items`, `codings`, `poll_laeufe`) |
|
| `schema.sql` | Postgres-Schema, vollstaendiger Stand fuer Neuinstallationen |
|
||||||
|
| `migrations/` | Aenderungen fuer bestehende Installationen, aufsteigend anwenden |
|
||||||
| `poller.py` | Poller, zwei Senken: Postgres und gzip-JSONL |
|
| `poller.py` | Poller, zwei Senken: Postgres und gzip-JSONL |
|
||||||
| `Containerfile` | Image fuer den Produktivbetrieb |
|
| `Containerfile` | Image fuer den Produktivbetrieb |
|
||||||
| `wurzelwerk-ingest.container`, `wurzelwerk.network` | Quadlet-Units, rootless |
|
| `wurzelwerk-ingest.container`, `wurzelwerk.network` | Quadlet-Units, rootless |
|
||||||
@@ -17,9 +18,13 @@ Kodierung — was hier landet, soll sich beliebig oft neu kodieren lassen.
|
|||||||
## Deploy mit Postgres + Podman/Quadlet (rootless)
|
## Deploy mit Postgres + Podman/Quadlet (rootless)
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
# 1. Schema einspielen
|
# 1. Schema einspielen (Neuinstallation)
|
||||||
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql
|
psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f schema.sql
|
||||||
|
|
||||||
|
# Bestehende Installation stattdessen migrieren, aufsteigend:
|
||||||
|
# psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/001-revisionen.sql
|
||||||
|
# psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/002-rohablage-und-felder.sql
|
||||||
|
|
||||||
# 2. Image bauen
|
# 2. Image bauen
|
||||||
podman build -t localhost/wurzelwerk-ingest:latest .
|
podman build -t localhost/wurzelwerk-ingest:latest .
|
||||||
|
|
||||||
@@ -53,7 +58,7 @@ python3 poller.py --sink jsonl:/pfad/zu/deinem/ablageordner
|
|||||||
```
|
```
|
||||||
|
|
||||||
Legt `<dir>/YYYYMMDD/YYYYMMDDHHMMSS.jsonl.gz` an, dedupliziert ueber eine
|
Legt `<dir>/YYYYMMDD/YYYYMMDDHHMMSS.jsonl.gz` an, dedupliziert ueber eine
|
||||||
SQLite daneben. Als `systemd --user`-Unit mit `Restart=always` einrichten.
|
SQLite daneben, und die rohen Feeds unter `<dir>/YYYYMMDD/feeds/`. Als `systemd --user`-Unit mit `Restart=always` einrichten.
|
||||||
|
|
||||||
Der Sinn ist nicht Redundanz um ihrer selbst willen: RSS-Feeds haben kein
|
Der Sinn ist nicht Redundanz um ihrer selbst willen: RSS-Feeds haben kein
|
||||||
Archiv. Steht der Produktiv-Poller einen Tag still, ist dieser Tag
|
Archiv. Steht der Produktiv-Poller einen Tag still, ist dieser Tag
|
||||||
@@ -78,6 +83,40 @@ gesperrte Quelle bleibt im Register stehen (`aktiv = false`), damit
|
|||||||
nachvollziehbar bleibt, warum sie fehlt. Der Poller verweigert den Start,
|
nachvollziehbar bleibt, warum sie fehlt. Der Poller verweigert den Start,
|
||||||
wenn eine als `gesperrt` gefuehrte Quelle aktiv geschaltet ist.
|
wenn eine als `gesperrt` gefuehrte Quelle aktiv geschaltet ist.
|
||||||
|
|
||||||
|
### Ueberschriften-Revisionen
|
||||||
|
|
||||||
|
```sql
|
||||||
|
-- welche Meldungen wurden nachtraeglich umbenannt
|
||||||
|
select quelle, fassungen, erste_fassung, letzte_fassung, zuletzt
|
||||||
|
from titel_revisionen order by zuletzt desc limit 20;
|
||||||
|
|
||||||
|
-- welche Redaktion revidiert am haeufigsten
|
||||||
|
select quelle, count(*) filter (where revisionen > 0) as revidiert,
|
||||||
|
count(*) as gesamt,
|
||||||
|
round(100.0 * count(*) filter (where revisionen > 0) / count(*), 1) as quote
|
||||||
|
from raw_items group by quelle order by quote desc;
|
||||||
|
```
|
||||||
|
|
||||||
|
### Was die Verlage selbst labeln
|
||||||
|
|
||||||
|
```sql
|
||||||
|
-- Themenverteilung je Quelle, ohne einen Kodierer bemueht zu haben
|
||||||
|
select quelle, k as kategorie, count(*)
|
||||||
|
from raw_items, unnest(kategorien) k
|
||||||
|
group by quelle, k order by count(*) desc limit 20;
|
||||||
|
```
|
||||||
|
|
||||||
|
### Speicherbedarf im Blick behalten
|
||||||
|
|
||||||
|
```sql
|
||||||
|
select pg_size_pretty(pg_total_relation_size('feed_abrufe')) as rohablage,
|
||||||
|
pg_size_pretty(pg_total_relation_size('raw_items')) as items;
|
||||||
|
|
||||||
|
-- Rohablage aelter als ein Jahr wegwerfen, falls es eng wird. raw_items
|
||||||
|
-- bleibt davon unberuehrt; nur die Moeglichkeit, neu zu parsen, entfaellt.
|
||||||
|
-- delete from feed_abrufe where abgerufen_am < now() - interval '1 year';
|
||||||
|
```
|
||||||
|
|
||||||
### Waechter
|
### Waechter
|
||||||
|
|
||||||
```sql
|
```sql
|
||||||
@@ -101,18 +140,56 @@ ueber die gesamte Historie, ohne dass ein Feed erneut abgerufen wird.
|
|||||||
|
|
||||||
## Verhalten
|
## Verhalten
|
||||||
|
|
||||||
- **2-Stunden-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz), um
|
- **15-Minuten-Takt**, an den Slotgrenzen ausgerichtet (+30 s Versatz).
|
||||||
unnoetigen Traffic bei den Quellen zu vermeiden.
|
Der Takt ist nicht frei waehlbar: ein RSS-Feed ist ein Fenster, kein Archiv.
|
||||||
|
Ist das Fenster kuerzer als der Poll-Abstand, fallen Artikel zwischen zwei
|
||||||
|
Abrufen heraus und sind unwiederbringlich weg. Gemessen am 2026-09-04:
|
||||||
|
zeit.de haelt 0,6 h vor, dlf gar nicht messbar, spiegel und welt rund 5 h --
|
||||||
|
bei 2 h Takt sah man von der Zeit etwa jeden dritten Artikel. Der bedingte
|
||||||
|
GET macht haeufiges Pollen billig: ein unveraenderter Abruf ist ein 304 ohne
|
||||||
|
Body.
|
||||||
|
- **Revisionen** werden mitgeschrieben. Redaktionen aendern Ueberschriften nach
|
||||||
|
der Veroeffentlichung; frueher lief das in ein `on conflict do nothing` und
|
||||||
|
blieb unsichtbar. `raw_items` haelt den aktuellen Stand, `raw_item_versionen`
|
||||||
|
jede Fassung. Der Datensatz existiert sonst nirgends -- GDELT sieht jede URL
|
||||||
|
genau einmal.
|
||||||
- **Bedingter GET** ueber ETag/Last-Modified, im Register gespeichert. Im
|
- **Bedingter GET** ueber ETag/Last-Modified, im Register gespeichert. Im
|
||||||
Test antworteten 4 von 12 Quellen im zweiten Durchlauf mit 304.
|
Test antworteten 4 von 12 Quellen im zweiten Durchlauf mit 304.
|
||||||
- **Dedup** ueber `unique (quelle, item_key)`; `item_key` ist guid, ersatzweise
|
- **Dedup** ueber `unique (quelle, item_key)`; `item_key` ist guid, ersatzweise
|
||||||
Link, ersatzweise ein Titel-Hash.
|
Link, ersatzweise ein Titel-Hash.
|
||||||
- **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der
|
- **Fehlertoleranz**: eine kaputte Quelle beendet den Durchlauf nicht, der
|
||||||
Fehler landet in `poll_laeufe`.
|
Fehler landet in `poll_laeufe`.
|
||||||
|
- **Rohe Feed-Antworten** landen gezippt in `feed_abrufe`, dedupliziert ueber
|
||||||
|
den sha256 der Rohbytes. Der Parser wird sich aendern; nur mit den Rohbytes
|
||||||
|
laesst sich eine spaetere Verbesserung rueckwirkend auf die Historie
|
||||||
|
anwenden. Kompression rund 22 %, etwa 20 kB je abgelegtem Abruf.
|
||||||
|
- **Der Parser nimmt mit, was im XML steht**: Kategorien (`<category>`,
|
||||||
|
`dc:subject`, Atom-`term`), `content:encoded` zusaetzlich zur kurzen
|
||||||
|
`description`, Medien-URLs aus `enclosure` und Media RSS, Verweise aus dem
|
||||||
|
Teaser-HTML, und die Position im Feed beim ersten Sehen. Die Kategorien sind
|
||||||
|
fuer Stufe 2 der wertvollste Posten: vom Verlag selbst vergebene Themenlabel,
|
||||||
|
brauchbar als Vorfilter und als Prompt-Kontext.
|
||||||
|
- **Items ohne Titel** werden weiterhin uebersprungen, aber gezaehlt
|
||||||
|
(`poll_laeufe.items_uebersprungen`) statt stillschweigend verworfen.
|
||||||
- **Kein Volltext-Abruf.** Nur der Feed selbst wird geholt. Titel und Teaser
|
- **Kein Volltext-Abruf.** Nur der Feed selbst wird geholt. Titel und Teaser
|
||||||
reichen als Kodierungsinput, und damit stellt sich die TDM-Frage nach
|
reichen als Kodierungsinput, und damit stellt sich die TDM-Frage nach
|
||||||
§ 44b UrhG gar nicht erst.
|
§ 44b UrhG gar nicht erst.
|
||||||
|
|
||||||
|
## Migration bestehender Installationen
|
||||||
|
|
||||||
|
`001` und `002` sind idempotent und koennen im laufenden Betrieb eingespielt
|
||||||
|
werden; danach den Poller neu starten.
|
||||||
|
|
||||||
|
Nach `002` sind die neuen Spalten auf Bestandszeilen zunaechst leer. Der
|
||||||
|
naechste Lauf traegt sie fuer alle Items nach, die die Quelle noch ausliefert
|
||||||
|
(in der Ausgabe als `nachgetr`), ohne sie als Revision zu zaehlen. Items, die
|
||||||
|
inzwischen aus dem Feed gefallen sind, behalten leere Zusatzfelder - die Daten
|
||||||
|
sind nicht mehr abrufbar.
|
||||||
|
|
||||||
|
Der Inhalts-Hash bleibt bewusst auf Titel, Teaser und URL beschraenkt. Naehme
|
||||||
|
er die neuen Felder auf, gaelte beim ersten Lauf nach dem Update jedes
|
||||||
|
Bestandsitem als revidiert.
|
||||||
|
|
||||||
## Gemessen (2026-09-04)
|
## Gemessen (2026-09-04)
|
||||||
|
|
||||||
Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter
|
Erster Lauf ueber 12 aktive Quellen: 820 Items gesehen, 810 neu, gesamter
|
||||||
|
|||||||
@@ -0,0 +1,80 @@
|
|||||||
|
-- Migration 001: Revisionshistorie fuer Items
|
||||||
|
--
|
||||||
|
-- Fuer bestehende Installationen. Bei einer Neuinstallation deckt schema.sql
|
||||||
|
-- das bereits ab; diese Datei ist dann ein No-op.
|
||||||
|
--
|
||||||
|
-- psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/001-revisionen.sql
|
||||||
|
|
||||||
|
begin;
|
||||||
|
|
||||||
|
alter table raw_items
|
||||||
|
add column if not exists inhalt_hash text,
|
||||||
|
add column if not exists revisionen integer not null default 0,
|
||||||
|
add column if not exists zuletzt_geaendert timestamptz;
|
||||||
|
|
||||||
|
-- Bestandsdaten nachziehen. Die Formel muss zeichengleich zu
|
||||||
|
-- poller.inhalt_hash() sein, sonst gilt beim naechsten Lauf jedes Altitem
|
||||||
|
-- als geaendert. Trennzeichen ist U+001F (unit separator).
|
||||||
|
update raw_items
|
||||||
|
set inhalt_hash = md5(coalesce(titel, '') || E'\x1f' ||
|
||||||
|
coalesce(teaser, '') || E'\x1f' ||
|
||||||
|
coalesce(url, ''))
|
||||||
|
where inhalt_hash is null;
|
||||||
|
|
||||||
|
alter table raw_items alter column inhalt_hash set not null;
|
||||||
|
|
||||||
|
create index if not exists raw_items_revidiert_idx
|
||||||
|
on raw_items (zuletzt_geaendert desc) where revisionen > 0;
|
||||||
|
|
||||||
|
create table if not exists raw_item_versionen (
|
||||||
|
id bigserial primary key,
|
||||||
|
raw_item_id bigint not null references raw_items(id) on delete cascade,
|
||||||
|
gesehen_am timestamptz not null default now(),
|
||||||
|
slot timestamptz not null,
|
||||||
|
titel text not null,
|
||||||
|
teaser text,
|
||||||
|
url text,
|
||||||
|
inhalt_hash text not null,
|
||||||
|
|
||||||
|
unique (raw_item_id, inhalt_hash)
|
||||||
|
);
|
||||||
|
|
||||||
|
create index if not exists raw_item_versionen_item_idx
|
||||||
|
on raw_item_versionen (raw_item_id, gesehen_am);
|
||||||
|
|
||||||
|
-- Erstfassung der Bestandsitems eintragen, damit die Historie vollstaendig ist.
|
||||||
|
-- gesehen_am auf abgerufen_am setzen, nicht auf now(): die Fassung ist alt.
|
||||||
|
insert into raw_item_versionen (raw_item_id, gesehen_am, slot, titel, teaser,
|
||||||
|
url, inhalt_hash)
|
||||||
|
select id, abgerufen_am, slot, titel, teaser, url, inhalt_hash
|
||||||
|
from raw_items
|
||||||
|
on conflict (raw_item_id, inhalt_hash) do nothing;
|
||||||
|
|
||||||
|
alter table poll_laeufe
|
||||||
|
add column if not exists items_geaendert integer not null default 0;
|
||||||
|
|
||||||
|
drop view if exists titel_revisionen;
|
||||||
|
create view titel_revisionen as
|
||||||
|
select r.id, r.quelle, r.url, r.revisionen,
|
||||||
|
count(*) as fassungen,
|
||||||
|
(array_agg(v.titel order by v.gesehen_am))[1] as erste_fassung,
|
||||||
|
(array_agg(v.titel order by v.gesehen_am desc))[1] as letzte_fassung,
|
||||||
|
min(v.gesehen_am) as erstmals,
|
||||||
|
max(v.gesehen_am) as zuletzt
|
||||||
|
from raw_items r
|
||||||
|
join raw_item_versionen v on v.raw_item_id = r.id
|
||||||
|
where r.revisionen > 0
|
||||||
|
group by r.id, r.quelle, r.url, r.revisionen;
|
||||||
|
|
||||||
|
create or replace view quellen_stillstand as
|
||||||
|
select q.id, q.name, max(r.abgerufen_am) as letztes_item,
|
||||||
|
now() - max(r.abgerufen_am) as stille
|
||||||
|
from quellen q
|
||||||
|
left join raw_items r on r.quelle = q.id
|
||||||
|
where q.aktiv
|
||||||
|
group by q.id, q.name
|
||||||
|
having max(r.abgerufen_am) is null
|
||||||
|
or now() - max(r.abgerufen_am) > interval '6 hours'
|
||||||
|
order by stille desc nulls first;
|
||||||
|
|
||||||
|
commit;
|
||||||
@@ -0,0 +1,47 @@
|
|||||||
|
-- Migration 002: Rohe Feed-Antworten ablegen, Parser-Felder ergaenzen
|
||||||
|
--
|
||||||
|
-- psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f migrations/002-rohablage-und-felder.sql
|
||||||
|
--
|
||||||
|
-- Setzt 001 voraus. Bestandszeilen bekommen die neuen Spalten leer; sie
|
||||||
|
-- fuellen sich, sobald die Quelle das Item das naechste Mal ausliefert.
|
||||||
|
-- Der Inhalts-Hash bleibt bewusst auf Titel/Teaser/URL beschraenkt, sonst
|
||||||
|
-- gaelte beim ersten Lauf nach dem Update jedes Bestandsitem als revidiert.
|
||||||
|
|
||||||
|
begin;
|
||||||
|
|
||||||
|
alter table raw_items
|
||||||
|
add column if not exists position integer,
|
||||||
|
add column if not exists volltext text,
|
||||||
|
add column if not exists kategorien text[],
|
||||||
|
add column if not exists medien jsonb,
|
||||||
|
add column if not exists links text[],
|
||||||
|
add column if not exists aktualisiert timestamptz;
|
||||||
|
|
||||||
|
comment on column raw_items.kategorien is
|
||||||
|
'Vom Verlag selbst vergebene Themenlabel - Vorfilter und Prompt-Kontext';
|
||||||
|
comment on column raw_items.position is
|
||||||
|
'Rang im Feed beim ersten Sehen. Spaeter nicht fortgeschrieben';
|
||||||
|
|
||||||
|
create table if not exists feed_abrufe (
|
||||||
|
id bigserial primary key,
|
||||||
|
quelle text not null references quellen(id),
|
||||||
|
slot timestamptz not null,
|
||||||
|
abgerufen_am timestamptz not null default now(),
|
||||||
|
http_status integer,
|
||||||
|
header jsonb,
|
||||||
|
koerper_hash text not null,
|
||||||
|
koerper bytea not null,
|
||||||
|
bytes_roh integer,
|
||||||
|
bytes_gz integer,
|
||||||
|
|
||||||
|
unique (quelle, koerper_hash)
|
||||||
|
);
|
||||||
|
|
||||||
|
create index if not exists feed_abrufe_quelle_idx
|
||||||
|
on feed_abrufe (quelle, abgerufen_am desc);
|
||||||
|
|
||||||
|
alter table poll_laeufe
|
||||||
|
add column if not exists items_uebersprungen integer not null default 0,
|
||||||
|
add column if not exists bytes_roh integer;
|
||||||
|
|
||||||
|
commit;
|
||||||
@@ -1,7 +1,7 @@
|
|||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
"""Ingest-Poller fuer den deutschen Ereignisdatensatz.
|
"""Ingest-Poller fuer den deutschen Ereignisdatensatz.
|
||||||
|
|
||||||
Liest die Quellen aus sources.toml, holt alle 2 Stunden die Feeds und legt
|
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
|
jedes neue Item roh ab. Bewusst ohne Kodierung: was hier landet, soll sich
|
||||||
beliebig oft neu kodieren lassen.
|
beliebig oft neu kodieren lassen.
|
||||||
|
|
||||||
@@ -10,7 +10,7 @@ Zwei Senken:
|
|||||||
--sink jsonl:/pfad/dir Schatten-Ingest als gzip-JSONL pro Slot
|
--sink jsonl:/pfad/dir Schatten-Ingest als gzip-JSONL pro Slot
|
||||||
|
|
||||||
Aufrufe:
|
Aufrufe:
|
||||||
poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 2-Stunden-Takt
|
poller.py --sink "$DATABASE_URL" # Dauerbetrieb, 15-Minuten-Takt
|
||||||
poller.py --sink "$DATABASE_URL" --once # ein Durchlauf
|
poller.py --sink "$DATABASE_URL" --once # ein Durchlauf
|
||||||
poller.py --check-robots # Lizenzstatus nachpruefen
|
poller.py --check-robots # Lizenzstatus nachpruefen
|
||||||
"""
|
"""
|
||||||
@@ -40,13 +40,18 @@ USER_AGENT = os.environ.get(
|
|||||||
"INGEST_USER_AGENT",
|
"INGEST_USER_AGENT",
|
||||||
"wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)",
|
"wurzelwerk-ingest/0.1 (+https://stinkwurzpresse.de/wurzelwerk; kontakt@stinkwurzpresse.de)",
|
||||||
)
|
)
|
||||||
TAKT = 2 * 60 * 60 # Slotlaenge in Sekunden
|
TAKT = 15 * 60 # Slotlaenge in Sekunden
|
||||||
VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf
|
VERSATZ = 30 # Sekunden nach der Slotgrenze, nicht exakt darauf
|
||||||
NS = {
|
NS = {
|
||||||
"a": "http://www.w3.org/2005/Atom",
|
"a": "http://www.w3.org/2005/Atom",
|
||||||
"dc": "http://purl.org/dc/elements/1.1/",
|
"dc": "http://purl.org/dc/elements/1.1/",
|
||||||
"content": "http://purl.org/rss/1.0/modules/content/",
|
"content": "http://purl.org/rss/1.0/modules/content/",
|
||||||
|
"media": "http://search.yahoo.com/mrss/",
|
||||||
}
|
}
|
||||||
|
# Header, die wir zum Abruf mitschreiben. Kosten nichts und beantworten
|
||||||
|
# spaeter Fragen zur Zwischenspeicherung, die man sonst nicht mehr stellen kann.
|
||||||
|
HEADER_MERKEN = ("Date", "Age", "Cache-Control", "Content-Type",
|
||||||
|
"Content-Length", "Server", "ETag", "Last-Modified")
|
||||||
|
|
||||||
_lauf = True
|
_lauf = True
|
||||||
|
|
||||||
@@ -66,7 +71,7 @@ def _jetzt():
|
|||||||
|
|
||||||
|
|
||||||
def slot_von(ts):
|
def slot_von(ts):
|
||||||
"""Auf TAKT-Sekunden seit Epoch abrunden (bei 2 h also volle gerade Stunden)."""
|
"""Auf TAKT-Sekunden seit Epoch abrunden (bei 15 min auf :00/:15/:30/:45)."""
|
||||||
sekunden = int(ts.timestamp())
|
sekunden = int(ts.timestamp())
|
||||||
return dt.datetime.fromtimestamp(
|
return dt.datetime.fromtimestamp(
|
||||||
(sekunden // TAKT) * TAKT, tz=dt.timezone.utc)
|
(sekunden // TAKT) * TAKT, tz=dt.timezone.utc)
|
||||||
@@ -107,6 +112,17 @@ def datum(s):
|
|||||||
return d.astimezone(dt.timezone.utc)
|
return d.astimezone(dt.timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def inhalt_hash(titel, teaser, url):
|
||||||
|
"""Fingerabdruck des Item-Inhalts.
|
||||||
|
|
||||||
|
Identisch zur SQL-Formel in der Migration, damit Bestandsdaten und neue
|
||||||
|
Zeilen denselben Hash tragen. md5 reicht - es geht um Gleichheit, nicht
|
||||||
|
um Faelschungssicherheit.
|
||||||
|
"""
|
||||||
|
roh = "\x1f".join((titel or "", teaser or "", url or ""))
|
||||||
|
return hashlib.md5(roh.encode("utf-8"), usedforsecurity=False).hexdigest()
|
||||||
|
|
||||||
|
|
||||||
def _text(el):
|
def _text(el):
|
||||||
return (el.text or "").strip() if el is not None and el.text else None
|
return (el.text or "").strip() if el is not None and el.text else None
|
||||||
|
|
||||||
@@ -124,7 +140,12 @@ def _erst(item, pfade):
|
|||||||
# Feed-Abruf und -Parsing
|
# Feed-Abruf und -Parsing
|
||||||
# --------------------------------------------------------------------------
|
# --------------------------------------------------------------------------
|
||||||
def hole(url, etag=None, last_modified=None, timeout=30):
|
def hole(url, etag=None, last_modified=None, timeout=30):
|
||||||
"""Bedingter GET. Gibt (status, bytes|None, etag, last_modified) zurueck."""
|
"""Bedingter GET.
|
||||||
|
|
||||||
|
Gibt (status, bytes|None, etag, last_modified, header) zurueck. Die
|
||||||
|
Rohbytes werden bewusst durchgereicht und nicht nur geparst: der Parser
|
||||||
|
wird sich aendern, die Historie soll dann rueckwirkend davon profitieren.
|
||||||
|
"""
|
||||||
kopf = {"User-Agent": USER_AGENT, "Accept-Encoding": "gzip"}
|
kopf = {"User-Agent": USER_AGENT, "Accept-Encoding": "gzip"}
|
||||||
if etag:
|
if etag:
|
||||||
kopf["If-None-Match"] = etag
|
kopf["If-None-Match"] = etag
|
||||||
@@ -136,23 +157,75 @@ def hole(url, etag=None, last_modified=None, timeout=30):
|
|||||||
roh = r.read()
|
roh = r.read()
|
||||||
if r.headers.get("Content-Encoding") == "gzip":
|
if r.headers.get("Content-Encoding") == "gzip":
|
||||||
roh = gzip.decompress(roh)
|
roh = gzip.decompress(roh)
|
||||||
return (r.status, roh,
|
header = {k: r.headers.get(k) for k in HEADER_MERKEN
|
||||||
r.headers.get("ETag"), r.headers.get("Last-Modified"))
|
if r.headers.get(k)}
|
||||||
|
return (r.status, roh, r.headers.get("ETag"),
|
||||||
|
r.headers.get("Last-Modified"), header)
|
||||||
except urllib.error.HTTPError as err:
|
except urllib.error.HTTPError as err:
|
||||||
if err.code == 304:
|
if err.code == 304:
|
||||||
return 304, None, etag, last_modified
|
return 304, None, etag, last_modified, {}
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
|
||||||
|
_HREF_RE = re.compile(r'href=["\']([^"\']+)["\']', re.I)
|
||||||
|
|
||||||
|
|
||||||
|
def _kategorien(it):
|
||||||
|
"""<category>, dc:subject, Atom-term. Vom Verlag selbst vergebene Label."""
|
||||||
|
aus = []
|
||||||
|
for el in it.findall("category") + it.findall("dc:subject", NS):
|
||||||
|
if _text(el):
|
||||||
|
aus.append(_text(el))
|
||||||
|
for el in it.findall("a:category", NS):
|
||||||
|
if el.get("term"):
|
||||||
|
aus.append(el.get("term"))
|
||||||
|
# Reihenfolge erhalten, Dubletten raus
|
||||||
|
return list(dict.fromkeys(aus))
|
||||||
|
|
||||||
|
|
||||||
|
def _medien(it):
|
||||||
|
"""enclosure und Media RSS. Acht von zwoelf Quellen liefern Bilder."""
|
||||||
|
aus = []
|
||||||
|
for el in it.findall("enclosure"):
|
||||||
|
if el.get("url"):
|
||||||
|
aus.append({"url": el.get("url"), "typ": el.get("type"),
|
||||||
|
"bytes": el.get("length"), "rolle": "enclosure"})
|
||||||
|
for tag, rolle in (("media:content", "content"),
|
||||||
|
("media:thumbnail", "thumbnail")):
|
||||||
|
for el in it.findall(tag, NS):
|
||||||
|
if el.get("url"):
|
||||||
|
aus.append({"url": el.get("url"), "typ": el.get("type"),
|
||||||
|
"bytes": el.get("fileSize"), "rolle": rolle})
|
||||||
|
return aus
|
||||||
|
|
||||||
|
|
||||||
|
def _links(*html_stuecke):
|
||||||
|
"""Verlinkungen aus dem Teaser-HTML, bevor klartext() die Tags entfernt.
|
||||||
|
|
||||||
|
GDELT haelt das Gegenstueck als PAGE_LINKS im GKG: welche Meldung auf
|
||||||
|
welche verweist, ist ein eigenes Signal.
|
||||||
|
"""
|
||||||
|
aus = []
|
||||||
|
for st in html_stuecke:
|
||||||
|
if st:
|
||||||
|
aus += _HREF_RE.findall(html.unescape(st))
|
||||||
|
return list(dict.fromkeys(aus))
|
||||||
|
|
||||||
|
|
||||||
def parse_feed(roh):
|
def parse_feed(roh):
|
||||||
"""RSS 2.0, RDF und Atom auf eine gemeinsame Form bringen."""
|
"""RSS 2.0, RDF und Atom auf eine gemeinsame Form bringen.
|
||||||
|
|
||||||
|
Gibt (items, uebersprungen) zurueck. Uebersprungen wird nur, was keinen
|
||||||
|
Titel hat - das wird gezaehlt statt stillschweigend verworfen.
|
||||||
|
"""
|
||||||
wurzel = ET.fromstring(roh)
|
wurzel = ET.fromstring(roh)
|
||||||
eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS)
|
eintraege = wurzel.findall(".//item") or wurzel.findall(".//a:entry", NS)
|
||||||
aus = []
|
aus, uebersprungen = [], 0
|
||||||
for it in eintraege:
|
for pos, it in enumerate(eintraege):
|
||||||
titel = _erst(it, ["title", "a:title"])
|
titel = _erst(it, ["title", "a:title"])
|
||||||
if not titel:
|
if not titel:
|
||||||
continue # ohne Titel ist ein Item fuer uns wertlos
|
uebersprungen += 1
|
||||||
|
continue
|
||||||
|
|
||||||
link = _erst(it, ["link", "a:id"])
|
link = _erst(it, ["link", "a:id"])
|
||||||
if not link:
|
if not link:
|
||||||
@@ -165,17 +238,36 @@ def parse_feed(roh):
|
|||||||
if not guid:
|
if not guid:
|
||||||
guid = hashlib.sha256(titel.encode()).hexdigest()[:32]
|
guid = hashlib.sha256(titel.encode()).hexdigest()[:32]
|
||||||
|
|
||||||
aus.append({
|
# Kurzfassung und Volltext getrennt halten: frueher gewann die kurze
|
||||||
|
# description und content:encoded fiel unter den Tisch.
|
||||||
|
teaser_html = _erst(it, ["description", "a:summary"])
|
||||||
|
voll_html = _erst(it, ["content:encoded", "a:content"])
|
||||||
|
if not teaser_html:
|
||||||
|
teaser_html = voll_html
|
||||||
|
|
||||||
|
eintrag = {
|
||||||
"titel": klartext(titel),
|
"titel": klartext(titel),
|
||||||
"url": link,
|
"url": link,
|
||||||
"guid": guid,
|
"guid": guid,
|
||||||
"teaser": klartext(_erst(it, ["description", "content:encoded",
|
"position": pos, # Reihenfolge im Feed = oft die Gewichtung
|
||||||
"a:summary", "a:content"])),
|
"teaser": klartext(teaser_html),
|
||||||
|
"volltext": klartext(voll_html) if voll_html is not teaser_html else None,
|
||||||
"autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]),
|
"autor": _erst(it, ["author", "dc:creator", "a:author/a:name"]),
|
||||||
"pubdate": datum(_erst(it, ["pubDate", "dc:date",
|
"kategorien": _kategorien(it),
|
||||||
"a:published", "a:updated"])),
|
"medien": _medien(it),
|
||||||
})
|
"links": _links(teaser_html, voll_html),
|
||||||
return aus
|
"pubdate": datum(_erst(it, ["pubDate", "dc:date", "a:published"])),
|
||||||
|
"aktualisiert": datum(_erst(it, ["a:updated", "dc:modified"])),
|
||||||
|
}
|
||||||
|
if eintrag["pubdate"] is None:
|
||||||
|
eintrag["pubdate"] = eintrag["aktualisiert"]
|
||||||
|
# Der Hash bleibt bewusst auf Titel/Teaser/URL beschraenkt. Naehme er
|
||||||
|
# die neuen Felder auf, gaelte beim ersten Lauf nach dem Update jedes
|
||||||
|
# Bestandsitem als revidiert.
|
||||||
|
eintrag["inhalt_hash"] = inhalt_hash(
|
||||||
|
eintrag["titel"], eintrag["teaser"], eintrag["url"])
|
||||||
|
aus.append(eintrag)
|
||||||
|
return aus, uebersprungen
|
||||||
|
|
||||||
|
|
||||||
# --------------------------------------------------------------------------
|
# --------------------------------------------------------------------------
|
||||||
@@ -195,7 +287,8 @@ class JsonlSenke:
|
|||||||
self.ziel.mkdir(parents=True, exist_ok=True)
|
self.ziel.mkdir(parents=True, exist_ok=True)
|
||||||
self.db = sqlite3.connect(self.ziel / "gesehen.sqlite")
|
self.db = sqlite3.connect(self.ziel / "gesehen.sqlite")
|
||||||
self.db.execute("create table if not exists gesehen ("
|
self.db.execute("create table if not exists gesehen ("
|
||||||
"quelle text, guid text, primary key (quelle, guid))")
|
"quelle text, guid text, hash text, "
|
||||||
|
"primary key (quelle, guid))")
|
||||||
self.db.commit()
|
self.db.commit()
|
||||||
|
|
||||||
def sync_quellen(self, quellen):
|
def sync_quellen(self, quellen):
|
||||||
@@ -206,16 +299,27 @@ class JsonlSenke:
|
|||||||
return None, None
|
return None, None
|
||||||
|
|
||||||
def schreibe(self, quelle, items, slot):
|
def schreibe(self, quelle, items, slot):
|
||||||
frisch = []
|
frisch, neu, geaendert = [], 0, 0
|
||||||
for it in items:
|
for it in items:
|
||||||
cur = self.db.execute(
|
row = self.db.execute(
|
||||||
"insert or ignore into gesehen (quelle, guid) values (?, ?)",
|
"select hash from gesehen where quelle = ? and guid = ?",
|
||||||
(quelle["id"], it["guid"]))
|
(quelle["id"], it["guid"])).fetchone()
|
||||||
if cur.rowcount:
|
if row is None:
|
||||||
|
self.db.execute("insert into gesehen (quelle, guid, hash) "
|
||||||
|
"values (?, ?, ?)",
|
||||||
|
(quelle["id"], it["guid"], it["inhalt_hash"]))
|
||||||
|
neu += 1
|
||||||
|
elif row[0] != it["inhalt_hash"]:
|
||||||
|
self.db.execute("update gesehen set hash = ? "
|
||||||
|
"where quelle = ? and guid = ?",
|
||||||
|
(it["inhalt_hash"], quelle["id"], it["guid"]))
|
||||||
|
geaendert += 1
|
||||||
|
else:
|
||||||
|
continue
|
||||||
frisch.append(it)
|
frisch.append(it)
|
||||||
self.db.commit()
|
self.db.commit()
|
||||||
if not frisch:
|
if not frisch:
|
||||||
return 0
|
return 0, 0, 0
|
||||||
tag = self.ziel / slot.strftime("%Y%m%d")
|
tag = self.ziel / slot.strftime("%Y%m%d")
|
||||||
tag.mkdir(exist_ok=True)
|
tag.mkdir(exist_ok=True)
|
||||||
pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz")
|
pfad = tag / (slot.strftime("%Y%m%d%H%M%S") + ".jsonl.gz")
|
||||||
@@ -226,7 +330,20 @@ class JsonlSenke:
|
|||||||
{"quelle": quelle["id"], "slot": slot.isoformat(),
|
{"quelle": quelle["id"], "slot": slot.isoformat(),
|
||||||
"abgerufen_am": _jetzt().isoformat(), **it},
|
"abgerufen_am": _jetzt().isoformat(), **it},
|
||||||
ensure_ascii=False, default=str) + "\n")
|
ensure_ascii=False, default=str) + "\n")
|
||||||
return len(frisch)
|
return neu, geaendert, 0
|
||||||
|
|
||||||
|
def speichere_feed(self, quelle, slot, status, header, koerper):
|
||||||
|
"""Rohes Feed-XML gezippt ablegen, dedupliziert ueber den Hash."""
|
||||||
|
if not koerper:
|
||||||
|
return False
|
||||||
|
h = hashlib.sha256(koerper).hexdigest()
|
||||||
|
tag = self.ziel / slot.strftime("%Y%m%d") / "feeds"
|
||||||
|
tag.mkdir(parents=True, exist_ok=True)
|
||||||
|
pfad = tag / f"{quelle['id']}-{h[:16]}.xml.gz"
|
||||||
|
if pfad.exists():
|
||||||
|
return False
|
||||||
|
pfad.write_bytes(gzip.compress(koerper))
|
||||||
|
return True
|
||||||
|
|
||||||
def notiere(self, *a, **kw):
|
def notiere(self, *a, **kw):
|
||||||
pass
|
pass
|
||||||
@@ -270,24 +387,84 @@ class PostgresSenke:
|
|||||||
return (r[0], r[1]) if r else (None, None)
|
return (r[0], r[1]) if r else (None, None)
|
||||||
|
|
||||||
def schreibe(self, quelle, items, slot):
|
def schreibe(self, quelle, items, slot):
|
||||||
neu = 0
|
"""Neue Items einfuegen, geaenderte fortschreiben.
|
||||||
|
|
||||||
|
Redaktionen aendern Ueberschriften nach der Veroeffentlichung. Frueher
|
||||||
|
lief das in ein `do nothing` und war unsichtbar. Jetzt schreibt der
|
||||||
|
Upsert den aktuellen Stand fort und legt jede Fassung zusaetzlich in
|
||||||
|
raw_item_versionen ab - auch die erste, damit die Historie vollstaendig
|
||||||
|
ist. Bleibt der Hash gleich, liefert das WHERE keine Zeile zurueck und
|
||||||
|
es passiert nichts.
|
||||||
|
"""
|
||||||
|
neu = geaendert = nachgetragen = 0
|
||||||
with self.conn.cursor() as cur:
|
with self.conn.cursor() as cur:
|
||||||
for it in items:
|
for it in items:
|
||||||
cur.execute("""
|
cur.execute("""
|
||||||
insert into raw_items
|
insert into raw_items
|
||||||
(quelle, item_key, url, titel, titel_norm, teaser,
|
(quelle, item_key, url, titel, titel_norm, teaser,
|
||||||
autor, pubdate, slot, roh)
|
autor, pubdate, slot, inhalt_hash, roh,
|
||||||
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
position, volltext, kategorien, medien, links,
|
||||||
on conflict (quelle, item_key) do nothing
|
aktualisiert)
|
||||||
returning id
|
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s,
|
||||||
|
%s, %s, %s, %s, %s, %s)
|
||||||
|
on conflict (quelle, item_key) do update set
|
||||||
|
url = excluded.url,
|
||||||
|
titel = excluded.titel,
|
||||||
|
titel_norm = excluded.titel_norm,
|
||||||
|
teaser = excluded.teaser,
|
||||||
|
autor = excluded.autor,
|
||||||
|
pubdate = excluded.pubdate,
|
||||||
|
position = coalesce(raw_items.position,
|
||||||
|
excluded.position),
|
||||||
|
volltext = excluded.volltext,
|
||||||
|
kategorien = excluded.kategorien,
|
||||||
|
medien = excluded.medien,
|
||||||
|
links = excluded.links,
|
||||||
|
aktualisiert = excluded.aktualisiert,
|
||||||
|
inhalt_hash = excluded.inhalt_hash,
|
||||||
|
roh = excluded.roh
|
||||||
|
where raw_items.inhalt_hash is distinct from excluded.inhalt_hash
|
||||||
|
or raw_items.kategorien is null
|
||||||
|
returning id, (xmax = 0) as ist_neu
|
||||||
""", (quelle["id"], it["guid"], it["url"], it["titel"],
|
""", (quelle["id"], it["guid"], it["url"], it["titel"],
|
||||||
norm_titel(it["titel"]), it["teaser"], it["autor"],
|
norm_titel(it["titel"]), it["teaser"], it["autor"],
|
||||||
it["pubdate"], slot,
|
it["pubdate"], slot, it["inhalt_hash"],
|
||||||
json.dumps(it, ensure_ascii=False, default=str)))
|
json.dumps(it, ensure_ascii=False, default=str),
|
||||||
if cur.fetchone():
|
it["position"], it["volltext"], it["kategorien"],
|
||||||
|
json.dumps(it["medien"], ensure_ascii=False),
|
||||||
|
it["links"], it["aktualisiert"]))
|
||||||
|
zeile = cur.fetchone()
|
||||||
|
if not zeile:
|
||||||
|
continue # unveraendert, nichts zu tun
|
||||||
|
rid, ist_neu = zeile
|
||||||
|
|
||||||
|
# Ob der Inhalt wirklich neu ist, sagt der Unique-Index der
|
||||||
|
# Fassungstabelle - zuverlaessiger als ein Vergleich in
|
||||||
|
# RETURNING, wo die Zeile schon fortgeschrieben ist.
|
||||||
|
cur.execute("""
|
||||||
|
insert into raw_item_versionen
|
||||||
|
(raw_item_id, slot, titel, teaser, url, inhalt_hash)
|
||||||
|
values (%s, %s, %s, %s, %s, %s)
|
||||||
|
on conflict (raw_item_id, inhalt_hash) do nothing
|
||||||
|
returning id
|
||||||
|
""", (rid, slot, it["titel"], it["teaser"], it["url"],
|
||||||
|
it["inhalt_hash"]))
|
||||||
|
inhalt_neu = cur.fetchone() is not None
|
||||||
|
|
||||||
|
if ist_neu:
|
||||||
neu += 1
|
neu += 1
|
||||||
|
elif inhalt_neu:
|
||||||
|
geaendert += 1
|
||||||
|
cur.execute("""
|
||||||
|
update raw_items
|
||||||
|
set revisionen = revisionen + 1,
|
||||||
|
zuletzt_geaendert = now()
|
||||||
|
where id = %s
|
||||||
|
""", (rid,))
|
||||||
|
else:
|
||||||
|
nachgetragen += 1 # nur die neuen Felder gefuellt
|
||||||
self.conn.commit()
|
self.conn.commit()
|
||||||
return neu
|
return neu, geaendert, nachgetragen
|
||||||
|
|
||||||
def merke_zustand(self, qid, etag, last_modified, geaendert):
|
def merke_zustand(self, qid, etag, last_modified, geaendert):
|
||||||
with self.conn.cursor() as cur:
|
with self.conn.cursor() as cur:
|
||||||
@@ -300,14 +477,43 @@ class PostgresSenke:
|
|||||||
""", (etag, last_modified, geaendert, qid))
|
""", (etag, last_modified, geaendert, qid))
|
||||||
self.conn.commit()
|
self.conn.commit()
|
||||||
|
|
||||||
def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu, fehler):
|
def speichere_feed(self, quelle, slot, status, header, koerper):
|
||||||
|
"""Rohes Feed-XML gezippt ablegen, dedupliziert ueber den Hash.
|
||||||
|
|
||||||
|
Der Parser wird sich aendern. Nur mit den Rohbytes laesst sich eine
|
||||||
|
spaetere Verbesserung rueckwirkend auf die Historie anwenden - dieselbe
|
||||||
|
Ueberlegung, aus der raw_items und codings getrennt sind.
|
||||||
|
"""
|
||||||
|
if not koerper:
|
||||||
|
return False
|
||||||
|
h = hashlib.sha256(koerper).hexdigest()
|
||||||
|
gz = gzip.compress(koerper)
|
||||||
|
with self.conn.cursor() as cur:
|
||||||
|
cur.execute("""
|
||||||
|
insert into feed_abrufe
|
||||||
|
(quelle, slot, http_status, header, koerper_hash,
|
||||||
|
koerper, bytes_roh, bytes_gz)
|
||||||
|
values (%s, %s, %s, %s, %s, %s, %s, %s)
|
||||||
|
on conflict (quelle, koerper_hash) do nothing
|
||||||
|
returning id
|
||||||
|
""", (quelle["id"], slot, status,
|
||||||
|
json.dumps(header, ensure_ascii=False), h,
|
||||||
|
gz, len(koerper), len(gz)))
|
||||||
|
neu = cur.fetchone() is not None
|
||||||
|
self.conn.commit()
|
||||||
|
return neu
|
||||||
|
|
||||||
|
def notiere(self, quelle, slot, dauer_ms, status, gesehen, neu,
|
||||||
|
geaendert, uebersprungen, bytes_roh, fehler):
|
||||||
with self.conn.cursor() as cur:
|
with self.conn.cursor() as cur:
|
||||||
cur.execute("""
|
cur.execute("""
|
||||||
insert into poll_laeufe
|
insert into poll_laeufe
|
||||||
(quelle, slot, dauer_ms, http_status,
|
(quelle, slot, dauer_ms, http_status, items_gesehen,
|
||||||
items_gesehen, items_neu, fehler)
|
items_neu, items_geaendert, items_uebersprungen,
|
||||||
values (%s, %s, %s, %s, %s, %s, %s)
|
bytes_roh, fehler)
|
||||||
""", (quelle, slot, dauer_ms, status, gesehen, neu, fehler))
|
values (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||||
|
""", (quelle, slot, dauer_ms, status, gesehen, neu, geaendert,
|
||||||
|
uebersprungen, bytes_roh, fehler))
|
||||||
self.conn.commit()
|
self.conn.commit()
|
||||||
|
|
||||||
def close(self):
|
def close(self):
|
||||||
@@ -323,36 +529,44 @@ def lade_quellen(pfad):
|
|||||||
|
|
||||||
|
|
||||||
def durchlauf(quellen, senke, slot):
|
def durchlauf(quellen, senke, slot):
|
||||||
ges_neu = ges_gesehen = 0
|
ges_neu = ges_gesehen = ges_geaendert = 0
|
||||||
for q in quellen:
|
for q in quellen:
|
||||||
if q["typ"] != "rss":
|
if q["typ"] != "rss":
|
||||||
continue
|
continue
|
||||||
t0 = time.monotonic()
|
t0 = time.monotonic()
|
||||||
etag, lm = senke.zustand(q["id"])
|
etag, lm = senke.zustand(q["id"])
|
||||||
status = gesehen = neu = 0
|
status = gesehen = neu = geaendert = uebersprungen = bytes_roh = 0
|
||||||
|
nachgetragen = 0
|
||||||
fehler = None
|
fehler = None
|
||||||
try:
|
try:
|
||||||
status, roh, n_etag, n_lm = hole(q["url"], etag, lm)
|
status, roh, n_etag, n_lm, header = hole(q["url"], etag, lm)
|
||||||
if status == 304:
|
if status == 304:
|
||||||
pass
|
pass
|
||||||
else:
|
else:
|
||||||
items = parse_feed(roh)
|
bytes_roh = len(roh)
|
||||||
|
senke.speichere_feed(q, slot, status, header, roh)
|
||||||
|
items, uebersprungen = parse_feed(roh)
|
||||||
gesehen = len(items)
|
gesehen = len(items)
|
||||||
neu = senke.schreibe(q, items, slot)
|
neu, geaendert, nachgetragen = senke.schreibe(q, items, slot)
|
||||||
if hasattr(senke, "merke_zustand"):
|
if hasattr(senke, "merke_zustand"):
|
||||||
senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0)
|
senke.merke_zustand(q["id"], n_etag, n_lm, neu > 0)
|
||||||
except Exception as e: # eine Quelle darf nicht
|
except Exception as e: # eine Quelle darf nicht
|
||||||
fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen
|
fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen
|
||||||
status = getattr(e, "code", None)
|
status = getattr(e, "code", None)
|
||||||
dauer = int((time.monotonic() - t0) * 1000)
|
dauer = int((time.monotonic() - t0) * 1000)
|
||||||
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, fehler)
|
senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert,
|
||||||
|
uebersprungen, bytes_roh, fehler)
|
||||||
ges_neu += neu
|
ges_neu += neu
|
||||||
|
ges_geaendert += geaendert
|
||||||
ges_gesehen += gesehen
|
ges_gesehen += gesehen
|
||||||
marke = "!" if fehler else (" " if status != 304 else ".")
|
marke = "!" if fehler else (" " if status != 304 else ".")
|
||||||
print(f" {marke} {q['id']:14} {str(status):>4} "
|
print(f" {marke} {q['id']:14} {str(status):>4} "
|
||||||
f"{gesehen:4} gesehen {neu:4} neu {dauer:5} ms"
|
f"{gesehen:4} gesehen {neu:4} neu {geaendert:3} rev "
|
||||||
|
f"{uebersprungen:2} uebrg"
|
||||||
|
+ (f" {nachgetragen:4} nachgetr" if nachgetragen else "")
|
||||||
|
+ f" {dauer:5} ms"
|
||||||
+ (f" {fehler}" if fehler else ""), flush=True)
|
+ (f" {fehler}" if fehler else ""), flush=True)
|
||||||
return ges_gesehen, ges_neu
|
return ges_gesehen, ges_neu, ges_geaendert
|
||||||
|
|
||||||
|
|
||||||
def check_robots(quellen):
|
def check_robots(quellen):
|
||||||
@@ -420,8 +634,9 @@ def main():
|
|||||||
while _lauf:
|
while _lauf:
|
||||||
slot = slot_von(_jetzt())
|
slot = slot_von(_jetzt())
|
||||||
print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True)
|
print(f"[{slot:%Y-%m-%d %H:%M}] Slot", flush=True)
|
||||||
gesehen, neu = durchlauf(aktiv, senke, slot)
|
gesehen, neu, geaendert = durchlauf(aktiv, senke, slot)
|
||||||
print(f" = {gesehen} gesehen, {neu} neu", flush=True)
|
print(f" = {gesehen} gesehen, {neu} neu, "
|
||||||
|
f"{geaendert} revidiert", flush=True)
|
||||||
if args.once:
|
if args.once:
|
||||||
break
|
break
|
||||||
ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ)
|
ziel = slot + dt.timedelta(seconds=TAKT + VERSATZ)
|
||||||
|
|||||||
+88
-3
@@ -50,17 +50,99 @@ create table if not exists raw_items (
|
|||||||
abgerufen_am timestamptz not null default now(),
|
abgerufen_am timestamptz not null default now(),
|
||||||
slot timestamptz not null, -- auf den Poll-Takt abgerundet
|
slot timestamptz not null, -- auf den Poll-Takt abgerundet
|
||||||
sprache text not null default 'de',
|
sprache text not null default 'de',
|
||||||
|
position integer, -- Rang beim ersten Sehen
|
||||||
|
volltext text, -- content:encoded, falls zusaetzlich da
|
||||||
|
kategorien text[], -- <category> / dc:subject / Atom-term
|
||||||
|
medien jsonb, -- enclosure und Media RSS
|
||||||
|
links text[], -- Verweise aus dem Teaser-HTML
|
||||||
|
aktualisiert timestamptz, -- Atom <updated>, neben pubdate
|
||||||
|
inhalt_hash text not null, -- md5 ueber titel|teaser|url
|
||||||
|
revisionen integer not null default 0,
|
||||||
|
zuletzt_geaendert timestamptz,
|
||||||
roh jsonb not null, -- alle geparsten Feldwerte, unveraendert
|
roh jsonb not null, -- alle geparsten Feldwerte, unveraendert
|
||||||
|
|
||||||
unique (quelle, item_key)
|
unique (quelle, item_key)
|
||||||
);
|
);
|
||||||
|
|
||||||
|
comment on column raw_items.kategorien is
|
||||||
|
'Vom Verlag selbst vergebene Themenlabel - Vorfilter und Prompt-Kontext';
|
||||||
|
comment on column raw_items.position is
|
||||||
|
'Rang im Feed beim ersten Sehen. Spaeter nicht fortgeschrieben';
|
||||||
|
comment on column raw_items.revisionen is
|
||||||
|
'Wie oft die Quelle das Item nach der Erstveroeffentlichung geaendert hat';
|
||||||
|
|
||||||
create index if not exists raw_items_slot_idx on raw_items (slot desc);
|
create index if not exists raw_items_slot_idx on raw_items (slot desc);
|
||||||
create index if not exists raw_items_pubdate_idx on raw_items (pubdate desc nulls last);
|
create index if not exists raw_items_pubdate_idx on raw_items (pubdate desc nulls last);
|
||||||
create index if not exists raw_items_quelle_idx on raw_items (quelle, abgerufen_am desc);
|
create index if not exists raw_items_quelle_idx on raw_items (quelle, abgerufen_am desc);
|
||||||
-- Titel-Clustering laeuft ueber Trigramm-Aehnlichkeit in SQL, nicht in Python
|
-- Titel-Clustering laeuft ueber Trigramm-Aehnlichkeit in SQL, nicht in Python
|
||||||
create index if not exists raw_items_titelnorm_trgm
|
create index if not exists raw_items_titelnorm_trgm
|
||||||
on raw_items using gin (titel_norm gin_trgm_ops);
|
on raw_items using gin (titel_norm gin_trgm_ops);
|
||||||
|
create index if not exists raw_items_revidiert_idx
|
||||||
|
on raw_items (zuletzt_geaendert desc) where revisionen > 0;
|
||||||
|
|
||||||
|
-- --------------------------------------------------------------------------
|
||||||
|
-- Fassungen. Redaktionen aendern Ueberschriften nach der Veroeffentlichung;
|
||||||
|
-- frueher lief das in ein `on conflict do nothing` und blieb unsichtbar.
|
||||||
|
-- raw_items haelt den aktuellen Stand, diese Tabelle jede Fassung - auch die
|
||||||
|
-- erste, damit die Historie vollstaendig ist.
|
||||||
|
--
|
||||||
|
-- Der Datensatz existiert sonst nirgends: GDELT sieht jede URL genau einmal.
|
||||||
|
-- --------------------------------------------------------------------------
|
||||||
|
create table if not exists raw_item_versionen (
|
||||||
|
id bigserial primary key,
|
||||||
|
raw_item_id bigint not null references raw_items(id) on delete cascade,
|
||||||
|
gesehen_am timestamptz not null default now(),
|
||||||
|
slot timestamptz not null,
|
||||||
|
titel text not null,
|
||||||
|
teaser text,
|
||||||
|
url text,
|
||||||
|
inhalt_hash text not null,
|
||||||
|
|
||||||
|
unique (raw_item_id, inhalt_hash)
|
||||||
|
);
|
||||||
|
|
||||||
|
create index if not exists raw_item_versionen_item_idx
|
||||||
|
on raw_item_versionen (raw_item_id, gesehen_am);
|
||||||
|
|
||||||
|
-- Items, deren Ueberschrift sich geaendert hat, mit erster und letzter Fassung
|
||||||
|
drop view if exists titel_revisionen;
|
||||||
|
create view titel_revisionen as
|
||||||
|
select r.id, r.quelle, r.url, r.revisionen,
|
||||||
|
count(*) as fassungen,
|
||||||
|
(array_agg(v.titel order by v.gesehen_am))[1] as erste_fassung,
|
||||||
|
(array_agg(v.titel order by v.gesehen_am desc))[1] as letzte_fassung,
|
||||||
|
min(v.gesehen_am) as erstmals,
|
||||||
|
max(v.gesehen_am) as zuletzt
|
||||||
|
from raw_items r
|
||||||
|
join raw_item_versionen v on v.raw_item_id = r.id
|
||||||
|
where r.revisionen > 0
|
||||||
|
group by r.id, r.quelle, r.url, r.revisionen;
|
||||||
|
|
||||||
|
-- --------------------------------------------------------------------------
|
||||||
|
-- Rohe Feed-Antworten. Der Parser wird sich aendern; nur mit den Rohbytes
|
||||||
|
-- laesst sich eine spaetere Verbesserung rueckwirkend auf die Historie
|
||||||
|
-- anwenden. Dieselbe Ueberlegung wie bei der Trennung raw_items / codings.
|
||||||
|
--
|
||||||
|
-- Dedupliziert ueber (quelle, koerper_hash): ein unveraenderter Feed wird
|
||||||
|
-- nicht zweimal abgelegt, und ein 304 hat ohnehin keinen Koerper.
|
||||||
|
-- --------------------------------------------------------------------------
|
||||||
|
create table if not exists feed_abrufe (
|
||||||
|
id bigserial primary key,
|
||||||
|
quelle text not null references quellen(id),
|
||||||
|
slot timestamptz not null,
|
||||||
|
abgerufen_am timestamptz not null default now(),
|
||||||
|
http_status integer,
|
||||||
|
header jsonb, -- Date, Age, Cache-Control, ...
|
||||||
|
koerper_hash text not null, -- sha256 ueber die Rohbytes
|
||||||
|
koerper bytea not null, -- gzip
|
||||||
|
bytes_roh integer,
|
||||||
|
bytes_gz integer,
|
||||||
|
|
||||||
|
unique (quelle, koerper_hash)
|
||||||
|
);
|
||||||
|
|
||||||
|
create index if not exists feed_abrufe_quelle_idx
|
||||||
|
on feed_abrufe (quelle, abgerufen_am desc);
|
||||||
|
|
||||||
-- --------------------------------------------------------------------------
|
-- --------------------------------------------------------------------------
|
||||||
-- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt
|
-- Kodierungen. Eine Zeile pro (Item, Kodierer-Version). Die Nutzlast bleibt
|
||||||
@@ -107,14 +189,17 @@ create table if not exists poll_laeufe (
|
|||||||
http_status integer, -- 304 = unveraendert
|
http_status integer, -- 304 = unveraendert
|
||||||
items_gesehen integer not null default 0,
|
items_gesehen integer not null default 0,
|
||||||
items_neu integer not null default 0,
|
items_neu integer not null default 0,
|
||||||
|
items_geaendert integer not null default 0,
|
||||||
|
items_uebersprungen integer not null default 0, -- ohne Titel, unbrauchbar
|
||||||
|
bytes_roh integer,
|
||||||
fehler text
|
fehler text
|
||||||
);
|
);
|
||||||
|
|
||||||
create index if not exists poll_laeufe_slot_idx on poll_laeufe (slot desc);
|
create index if not exists poll_laeufe_slot_idx on poll_laeufe (slot desc);
|
||||||
create index if not exists poll_laeufe_quelle_idx on poll_laeufe (quelle, begonnen_am desc);
|
create index if not exists poll_laeufe_quelle_idx on poll_laeufe (quelle, begonnen_am desc);
|
||||||
|
|
||||||
-- Welche aktive Quelle hat seit ueber 24 Stunden (12 verpasste Slots bei
|
-- Welche aktive Quelle hat seit ueber sechs Stunden (24 verpasste Slots bei
|
||||||
-- 2-Stunden-Takt) nichts Neues geliefert?
|
-- 15-Minuten-Takt) nichts Neues geliefert?
|
||||||
create or replace view quellen_stillstand as
|
create or replace view quellen_stillstand as
|
||||||
select q.id, q.name, max(r.abgerufen_am) as letztes_item,
|
select q.id, q.name, max(r.abgerufen_am) as letztes_item,
|
||||||
now() - max(r.abgerufen_am) as stille
|
now() - max(r.abgerufen_am) as stille
|
||||||
@@ -123,5 +208,5 @@ left join raw_items r on r.quelle = q.id
|
|||||||
where q.aktiv
|
where q.aktiv
|
||||||
group by q.id, q.name
|
group by q.id, q.name
|
||||||
having max(r.abgerufen_am) is null
|
having max(r.abgerufen_am) is null
|
||||||
or now() - max(r.abgerufen_am) > interval '24 hours'
|
or now() - max(r.abgerufen_am) > interval '6 hours'
|
||||||
order by stille desc nulls first;
|
order by stille desc nulls first;
|
||||||
|
|||||||
Reference in New Issue
Block a user