Große Datenmengen über Proxys auslagern und nach einem Abbruch nicht von vorn beginnen
Inhalt des Artikels
- Einführung: warum eine lange auslagerung fast immer abbricht und das normal ist
- Vorbereitung: tools, zugänge und umgebung
- Grundbegriffe: wörterbuch des robusten auslagerns in einfachen worten
- Schritt 1: http-wiederaufnahme über den range-header
- Schritt 2: checkpoints für seitenweise auslagerungen
- Schritt 3: idempotenz, damit wiederholung keine duplikate erzeugt
- Schritt 4: deduplizierung von ergebnissen ohne speicher aufzublähen
- Schritt 5: parallelität ohne verluste
- Schritt 6: wiederaufnahme nach langer pause
- Schritt 7: fertiges grundgerüst eines robusten downloaders in python
- Ergebnisprüfung: checkliste für robustes auslagern
- Typische fehler und ihre lösungen
- Zusätzliche möglichkeiten und optimierung
- Faq: häufige fragen zum robusten auslagern
- Fazit: was sie jetzt können und wohin sie weitergehen
Einführung: Warum eine lange Auslagerung fast immer abbricht und das normal ist
Wenn Sie schon einmal eine große Datenauslagerung gestartet haben, kennen Sie dieses Gefühl. Der Prozess lief mehrere Stunden, erreichte neunzig Prozent und brach ab. Die Verbindung fiel aus, der Server gab einen Fehler zurück, das Notebook schlief ein. Und alles muss von vorn beginnen. Diese Anleitung wurde geschrieben, damit das nicht mehr passiert.
Was Sie am Ende erhalten. Sie lernen, einen Downloader zu bauen, der Abbrüche übersteht. Er setzt Dateien ab der Mitte fort, merkt sich, auf welcher Seite er stehen geblieben ist, erzeugt bei Wiederholung keine Duplikate und kann selbst nach einer langen Pause wieder aufnehmen. Sie erhalten ein fertiges Python-Grundgerüst, das Sie an Ihre Aufgabe anpassen können.
Für wen diese Anleitung ist. Für Ingenieure, Analysten und Entwickler, die Daten aus APIs auslagern, große Dateien herunterladen oder seitenweise Ergebnisse über Proxys sammeln. Mittleres Niveau. Sie sollten die Grundlagen von HTTP verstehen und Python-Code lesen können. Tiefe Kenntnisse der Netzwerkprogrammierung sind nicht erforderlich.
Was Sie vorher wissen sollten. Basis-Python, das Konzept von HTTP-Anfrage und -Antwort, was Header und Statuscodes sind. Wenn Sie mit der requests-Bibliothek gearbeitet haben, reicht das aus.
Wie viel Zeit Sie brauchen. Für das Lesen und Verstehen der Konzepte etwa vierzig Minuten. Für den Zusammenbau eines funktionierenden Downloaders nach unserem Grundgerüst ein bis drei Stunden, je nach Ihrer Datenquelle.
Wichtige Klarstellung zum Thema. Wir werden keine Statuscodes wie 429 und keine Wiederholungsstrategien mit Verzögerung (Backoff) behandeln. Dazu gibt es separates Material. Hier liegt der Fokus auf einer Sache: dem Zustand des Prozesses und seiner Wiederaufnahme. Wie man den Fortschritt speichert, wie man Daten nicht verliert und nicht dupliziert, wie man dort weitermacht, wo man aufgehört hat.
Tipp: Halten Sie einen Notizblock oder eine separate Datei bereit, in der Sie die Parameter Ihrer Quelle notieren: ob sie Wiederaufnahme unterstützt, ob es Seiten-Navigation gibt, welches Cursor-Format sie verwendet. Diese Notizen werden Ihnen bei jedem Schritt nützlich sein.
Vorbereitung: Tools, Zugänge und Umgebung
Bevor wir Code schreiben, richten wir die Arbeitsumgebung ein. Das dauert zehn Minuten, spart aber Stunden beim Debuggen.
Was installiert werden muss
- Installieren Sie Python Version 3.10 oder neuer. Prüfen Sie die Version mit dem Befehl
python --versionim Terminal. - Erstellen Sie eine virtuelle Umgebung mit dem Befehl
python -m venv venv, damit die Projektabhängigkeiten nicht mit den Systemabhängigkeiten vermischt werden. - Aktivieren Sie die Umgebung. Unter Windows ist das
venv\Scripts\activate, unter macOS und Linuxsource venv/bin/activate. - Installieren Sie die Bibliothek für HTTP-Anfragen mit dem Befehl
pip install requests. - Für eine schnellere Arbeit mit der Zustandsdatenbank müssen Sie nichts zusätzlich installieren: Das Modul
sqlite3ist bereits in der Python-Standardbibliothek enthalten.
Welche Zugänge nötig sind
- Zugang zu Ihrer Datenquelle: URL, Token oder API-Schlüssel, falls erforderlich.
- Ein Proxeon-Proxy mit Adresse, Port und Authentifizierungsdaten. Ohne einen stabilen Proxy verliert robustes Auslagern seinen Sinn, denn genau der Proxy verteilt die Last und macht Verbindungen vorhersehbar.
- Speicherplatz für die Zustandsdatei und für die auszulagernden Daten selbst.
Proxeon-Proxy prüfen
- Nehmen Sie eine Verbindungszeichenfolge der Form
http://login:passwort@adresse:port. - Testen Sie sie mit einer einfachen Anfrage. Führen Sie im Terminal
curl -x http://login:passwort@adresse:port https://api.ipify.orgaus und stellen Sie sicher, dass die IP-Adresse des Proxys zurückgegeben wird, nicht Ihre eigene.
Tipp: Speichern Sie die Verbindungszeichenfolge zum Proxy in einer Umgebungsvariable, nicht im Code. So senden Sie das Passwort nicht versehentlich ins Versionskontrollsystem. Im Code lesen Sie sie über os.environ.
⚠️ Achtung: Arbeiten Sie immer nur mit Datenquellen, zu denen Sie rechtmäßigen Zugang haben. Beachten Sie die Nutzungsbedingungen des Dienstes und die vom Eigentümer angegebenen Limits. Der Proxeon-Proxy ist für legale Ingenieursarbeit gedacht: Lastverteilung, Verbindungsstabilität und korrektes Auslagern von Daten.
✅ Prüfung: Die Umgebung ist bereit, wenn der Befehl python -c "import requests, sqlite3" fehlerfrei ausgeführt wurde und die Anfrage über den Proxy die Adresse des Proxy-Servers zurückgegeben hat.
Grundbegriffe: Wörterbuch des robusten Auslagerns in einfachen Worten
Bevor wir Code schreiben, klären wir die Schlüsselbegriffe. Ohne sie klingen die weiteren Schritte wie Zauberformeln.
Wiederaufnahme
Wiederaufnahme ist die Fortsetzung des Downloads ab dem Byte, an dem er unterbrochen wurde. Statt die Datei von vorn zu laden, bitten Sie den Server, nur das fehlende Stück zu senden. Das funktioniert über den HTTP-Header Range.
Checkpoint
Checkpoint ist ein gespeicherter Fortschrittspunkt. Stellen Sie sich das Speichern in einem Computerspiel vor. Wenn etwas schiefgeht, kehren Sie zum letzten Speicherstand zurück, nicht zum Spielanfang. Beim Auslagern speichert der Checkpoint, auf welcher Seite oder bei welchem Datensatz Sie stehen geblieben sind.
Cursor
Cursor ist eine Markierung, die Ihnen die API gibt, damit Sie die nächste Datenportion anfordern können. Oft ist es eine Zeichenkette wie eyJvZmZzZXQiOjEwMH0. Sie senden sie zurück, und der Server weiß, wo er fortfahren soll.
Idempotenz
Idempotenz ist die Eigenschaft einer Operation, bei der ihre Wiederholung das Ergebnis nicht ändert. Wenn Sie zweimal dieselbe Zeile mit demselben Schlüssel schreiben, bleibt am Ende eine Zeile, nicht zwei. Das ist der Schutz vor Duplikaten bei Wiederholungen.
Deduplizierung
Deduplizierung ist das Aussortieren wiederholter Datensätze. Selbst bei sorgfältiger Arbeit kann dasselbe Objekt zweimal ankommen. Die Deduplizierung stellt sicher, dass es in Ihrem Endergebnis nur einmal bleibt.
Grundprinzip
Ein robuster Downloader basiert auf einer Idee: Fortschritt muss ständig gespeichert werden, nicht nur am Ende. Jeder Schritt kann der letzte vor einem Abbruch sein. Also muss nach jedem erfolgreichen Arbeitsstück der Zustand auf die Festplatte geschrieben werden. Dann ist Wiederaufnahme einfach das Lesen des Zustands und das Fortsetzen.
Tipp: Merken Sie sich die Regel der drei Fragen für jede Auslagerung. Erstens: Wo habe ich aufgehört? Zweitens: Wie verhindere ich Duplikate des bereits Erhaltenen? Drittens: Was veraltet, während ich weg war? Die Antworten darauf machen die Robustheit aus.
Schritt 1: HTTP-Wiederaufnahme über den Range-Header
Ziel dieses Schritts. Lernen, eine große Datei so herunterzuladen, dass nach einem Abbruch ab dem fehlenden Byte fortgesetzt wird, nicht von vorn.
Wie es funktioniert
HTTP erlaubt es, nicht die ganze Datei anzufordern, sondern nur einen Teil. Dazu wird der Header Range zur Anfrage hinzugefügt. Zum Beispiel bedeutet Range: bytes=1048576-: Gib mir alles ab Byte-Nummer 1048576. Aber zuerst muss geprüft werden, ob der Server das kann.
- Senden Sie eine HEAD-Anfrage oder eine normale GET-Anfrage an die Datei und sehen Sie sich die Antwort-Header an.
- Suchen Sie den Header
Accept-Ranges. Wenn dortbytessteht, unterstützt der Server die Wiederaufnahme. - Wenn der Header fehlt oder
nonelautet, ist Wiederaufnahme nicht möglich. Dann muss die Datei in einem Durchgang vollständig geladen werden oder eine alternative Quelle gesucht werden.
Unterstützung der Wiederaufnahme prüfen
Hier ist Code, der prüft, ob der Server Teile der Datei senden kann.
import requests
def supports_resume(url, proxies):
resp = requests.head(url, proxies=proxies, timeout=30, allow_redirects=True)
accept = resp.headers.get("Accept-Ranges", "none")
total = resp.headers.get("Content-Length")
return accept.lower() == "bytes", totalDatei ab der Mitte fortsetzen
Jetzt der Hauptcode. Er prüft, wie viele Bytes bereits lokal heruntergeladen wurden, und fordert vom Server nur den Rest an.
import os
import requests
def download_resumable(url, dest, proxies):
already = 0
if os.path.exists(dest):
already = os.path.getsize(dest)
headers = {}
if already > 0:
headers["Range"] = f"bytes={already}-"
mode = "ab" if already > 0 else "wb"
with requests.get(url, headers=headers, proxies=proxies,
stream=True, timeout=60) as r:
if already > 0 and r.status_code == 200:
mode = "wb"
already = 0
with open(dest, mode) as f:
for chunk in r.iter_content(chunk_size=65536):
if chunk:
f.write(chunk)
return os.path.getsize(dest)Wichtige Punkte. Wenn der Server Status 206 zurückgibt, hat er ehrlich einen Teil der Datei gesendet, und das Anhängen funktioniert korrekt. Wenn der Server trotz Range-Header 200 zurückgibt, ignoriert er die Wiederaufnahme und sendet die ganze Datei. In diesem Fall wechseln wir in den Modus des vollständigen Überschreibens, damit wir nicht das alte Stück mit dem neuen zusammenkleben und die Datei beschädigen.
⚠️ Achtung: Hängen Sie niemals Daten im Modus ab an, wenn Sie nicht sicher sind, dass der Server mit Code 206 geantwortet hat. Sonst erhalten Sie eine kaputte Datei, bei der der Anfang der Rest des vorherigen Versuchs ist und die Fortsetzung die neue vollständige Datei. Eine solche Datei öffnet sich mit Fehler, und Sie verlieren Zeit bei der Ursachensuche.
Tipp: Laden Sie nicht direkt in die Zieldatei, sondern in eine temporäre Datei mit der Endung .part. Wenn der Download vollständig abgeschlossen ist, benennen Sie sie in den endgültigen Namen um. So verwechseln Sie nie eine fertige Datei mit einer unvollständigen.
Integritätskontrolle
Nach dem vollständigen Download ist es gut zu prüfen, ob die Datei nicht beschädigt ist. Wenn der Server den Header Content-Length gesendet hat, vergleichen Sie ihn mit der tatsächlichen Dateigröße auf der Festplatte. Wenn die Größen übereinstimmen, ist die Datei vollständig angekommen.
def verify_size(dest, expected):
if expected is None:
return True
return os.path.getsize(dest) == int(expected)✅ Prüfung: Unterbrechen Sie den Download in der Mitte, indem Sie das Programm schließen. Starten Sie es erneut. In den Logs sollten Sie sehen, dass die Anfrage mit dem Range-Header gesendet wurde und die Datei fortgesetzt wurde, nicht neu begonnen hat. Die endgültige Größe stimmt mit der erwarteten überein.
Schritt 2: Checkpoints für seitenweise Auslagerungen
Ziel dieses Schritts. Das Speichern des Fortschritts für APIs einrichten, die Daten seitenweise liefern, damit nach einem Abbruch mit der richtigen Seite fortgesetzt wird.
Was genau gespeichert werden soll
Eine Datei wird byteweise fortgesetzt, eine seitenweise Auslagerung folgt einer ganz anderen Logik. Hier gibt es keine Bytes, sondern Seiten und Datensätze. Also muss im Checkpoint etwas anderes gespeichert werden.
- Cursor wenn die API mit Cursorn arbeitet. Das ist die zuverlässigste Variante, weil der Cursor selbst weiß, wo fortgesetzt werden soll.
- Seitennummer oder Offset wenn die API mit offset und limit arbeitet. Speichern Sie die Nummer der letzten erfolgreich verarbeiteten Seite.
- ID des letzten Datensatzes wenn nach aufsteigender ID oder Datum sortiert werden kann. Dann fordert die nächste Anfrage Datensätze mit einer ID größer als die gespeicherte an.
- Zähler der verarbeiteten Datensätze zur Kontrolle und Berichterstattung.
Wo der Zustand gespeichert wird
Es gibt drei Hauptvarianten, von einfach bis robust.
- JSON-Datei. Am einfachsten. Sie schreiben ein Wörterbuch mit Cursor und Zähler nach jeder Seite in die Datei. Geeignet für einzelne, nicht parallele Auslagerungen.
- SQLite-Datenbank. Robuster. Bietet Transaktionen, sodass der Zustand bei einem Abbruch während des Schreibens nicht beschädigt wird. Gut, wenn viele Daten anfallen und Deduplizierung nötig ist.
- Externe Datenbank. Für große verteilte Auslagerungen, wenn mehrere Prozesse eine Arbeit teilen.
Checkpoint in JSON speichern
import json
import os
def save_checkpoint(path, cursor, page, last_id, count):
tmp = path + ".tmp"
data = {
"cursor": cursor,
"page": page,
"last_id": last_id,
"count": count,
}
with open(tmp, "w") as f:
json.dump(data, f)
os.replace(tmp, path)
def load_checkpoint(path):
if not os.path.exists(path):
return {"cursor": None, "page": 0, "last_id": None, "count": 0}
with open(path) as f:
return json.load(f)Achten Sie auf den Trick mit der temporären Datei. Wir schreiben in eine Datei mit dem Suffix .tmp und benennen sie dann atomar über os.replace um. Das schützt vor der Situation, dass das Programm genau während des Schreibens des Checkpoints abstürzt. Der alte Checkpoint bleibt dabei intakt und verwandelt sich nicht in ein halbes JSON, das nicht gelesen werden kann.
⚠️ Achtung: Schreiben Sie den Checkpoint niemals direkt in dieselbe Datei über den alten Stand ohne temporäre Datei. Ein Abbruch mitten im Schreiben hinterlässt einen beschädigten Checkpoint, und die Wiederaufnahme wird unmöglich. Die atomare Ersetzung löst dieses Problem vollständig.
Tipp: Speichern Sie den Checkpoint erst, nachdem die Daten der Seite tatsächlich im Speicher geschrieben wurden. Die Reihenfolge ist: Seite erhalten, Daten geschrieben, dann Checkpoint aktualisiert. Wenn Sie die Reihenfolge vertauschen, überspringen Sie bei einem Abbruch eine Seite und verlieren Daten.
Hauptschleife mit Checkpoints
def paginate(fetch_page, save_data, cp_path, proxies):
cp = load_checkpoint(cp_path)
cursor = cp["cursor"]
count = cp["count"]
while True:
items, next_cursor = fetch_page(cursor, proxies)
if not items:
break
save_data(items)
count += len(items)
last_id = items[-1].get("id")
save_checkpoint(cp_path, next_cursor, cp["page"] + 1,
last_id, count)
cursor = next_cursor
if next_cursor is None:
break
return count✅ Prüfung: Starten Sie die Auslagerung, lassen Sie sie ein paar Seiten verarbeiten, unterbrechen Sie. Öffnen Sie die Checkpoint-Datei und stellen Sie sicher, dass der aktuelle Cursor und Zähler gespeichert sind. Starten Sie erneut: Die Auslagerung sollte mit dem gespeicherten Cursor fortgesetzt werden, nicht mit der ersten Seite.
Schritt 3: Idempotenz, damit Wiederholung keine Duplikate erzeugt
Ziel dieses Schritts. Dafür sorgen, dass ein erneuter Lauf oder die Wiederholung einer bestimmten Anfrage nicht zu identischen Datensätzen in Ihrem Speicher führt.
Warum Duplikate entstehen
Stellen Sie sich vor: Sie haben eine Datenseite erhalten, sie in die Datei geschrieben, aber das Programm ist abgestürzt, bevor der Checkpoint aktualisiert wurde. Beim nächsten Lauf fordern Sie dieselbe Seite erneut an. Die Daten kommen wieder und werden ein zweites Mal geschrieben. So entstehen Duplikate. Das ist eine unvermeidliche Folge von Abbrüchen, und man muss auf Architekturebene dagegen ankämpfen.
Deduplizierungsschlüssel
Das Hauptwerkzeug der Idempotenz ist der Deduplizierungsschlüssel. Das ist ein Feld oder eine Kombination von Feldern, die einen Datensatz eindeutig bestimmen. Die richtige Wahl des Schlüssels löst die Hälfte des Problems.
- Natürliche ID. Wenn der Datensatz eine eindeutige Kennung von der Quelle hat, verwenden Sie sie. Das ist der ideale Schlüssel.
- Feldkombination. Wenn es keine einheitliche ID gibt, bauen Sie den Schlüssel aus mehreren stabilen Feldern. Zum Beispiel E-Mail plus Registrierungsdatum.
- Inhalts-Hash. Wenn es gar keine stabilen Felder gibt, berechnen Sie einen Hash des gesamten Datensatzes. Das ist die letzte Option, weil jede Feldänderung einen neuen Schlüssel erzeugt.
Schreiben ohne Duplikate über UPSERT
Wenn Sie das Ergebnis in SQLite oder einer anderen Datenbank speichern, verwenden Sie eine Einfügung mit Ignorieren des Konflikts. Dann bewirkt ein erneutes Schreiben mit demselben Schlüssel einfach nichts.
import sqlite3
def init_db(path):
conn = sqlite3.connect(path)
conn.execute(
"CREATE TABLE IF NOT EXISTS records ("
"dedup_key TEXT PRIMARY KEY, payload TEXT)"
)
conn.commit()
return conn
def save_records(conn, items):
rows = [(item["id"], json.dumps(item)) for item in items]
conn.executemany(
"INSERT OR IGNORE INTO records (dedup_key, payload) "
"VALUES (?, ?)", rows
)
conn.commit()Das entscheidende Detail hier ist der PRIMARY KEY auf dem Feld dedup_key. Die Datenbank lehnt eine erneute Einfügung mit demselben Schlüssel von selbst ab, weil INSERT OR IGNORE den Konflikt still schluckt. Sie müssen nicht manuell prüfen, ob ein solcher Datensatz bereits existiert. Die Datenbank erledigt das für Sie und macht es schnell.
Tipp: Wählen Sie den Deduplizierungsschlüssel einmal am Projektanfang und dokumentieren Sie ihn. Ein Wechsel des Schlüssels mitten in der Auslagerung bedeutet, dass alte und neue Datensätze nicht mehr zugeordnet werden können und doch Duplikate entstehen. Stabilität des Schlüssels ist wichtiger als seine Schönheit.
✅ Prüfung: Starten Sie die Auslagerung zweimal hintereinander auf demselben Datenbereich. Zählen Sie die Anzahl der Zeilen in der Datenbank mit dem Befehl SELECT COUNT(*) FROM records. Die Zahl muss nach dem ersten und nach dem zweiten Lauf gleich sein.
Schritt 4: Deduplizierung von Ergebnissen ohne Speicher aufzublähen
Ziel dieses Schritts. Wiederholte Datensätze bei Millionen von Zeilen aussortieren, ohne den gesamten Arbeitsspeicher des Computers mit einer Menge bereits gesehener Schlüssel zu belasten.
Naiver Ansatz und sein Problem
Die einfachste Deduplizierung: ein set im Speicher mit allen gesehenen Schlüsseln halten. Für jeden neuen Datensatz prüfen, ob der Schlüssel im Set ist. Funktioniert hervorragend bei Hunderttausenden Zeilen. Aber bei Millionen und Dutzenden Millionen wächst das Set und frisst Gigabyte an Arbeitsspeicher. Das Programm wird langsamer oder stürzt ab.
Lösung eins: auf die Datenbank verlassen
Die einfachste und zuverlässigste Methode bei großen Volumina ist, das Gesehene gar nicht im Speicher zu halten, sondern die Prüfung der Datenbank über PRIMARY KEY zu überlassen, wie wir es im vorherigen Schritt gemacht haben. Die Datenbank speichert den Index auf der Festplatte, nicht im Speicher Ihres Prozesses. Sie bewältigt Dutzende Millionen Schlüssel ohne Belastung Ihres Arbeitsspeichers.
Lösung zwei: Hash des Datensatzes
Wenn es keinen natürlichen Schlüssel gibt, berechnen Sie einen kompakten Hash des Datensatzes. Der Hash belegt einen festen und kleinen Umfang unabhängig von der Größe des Datensatzes selbst.
import hashlib
import json
def record_hash(item):
raw = json.dumps(item, sort_keys=True, ensure_ascii=False)
return hashlib.sha256(raw.encode("utf-8")).hexdigest()Der Parameter sort_keys=True ist hier entscheidend. Er stellt sicher, dass inhaltlich gleiche Datensätze denselben Hash ergeben, auch wenn die Felder darin in unterschiedlicher Reihenfolge standen. Ohne diese Sortierung können zwei identische Objekte unterschiedliche Hashes erhalten und als verschiedene Datensätze durchrutschen.
Lösung drei: Bloom-Filter zur Speichereinsparung
Wenn Sie dennoch eine schnelle Prüfung im Speicher bei riesigen Volumina benötigen, verwendet man einen Bloom-Filter. Das ist eine Struktur, die wenig Platz einnimmt und schnell antwortet, ob wir einen Schlüssel gesehen haben oder definitiv nicht. Sie hat eine Besonderheit: Sie kann gelegentlich fälschlich sagen, dass ein Schlüssel schon da war, obwohl er nicht da war. Deshalb wird der Bloom-Filter als schneller Vorfilter verwendet, und die endgültige Prüfung überlässt man der Datenbank.
- Wir prüfen den Schlüssel mit dem Bloom-Filter.
- Wenn der Filter sagt, dass wir ihn definitiv nicht gesehen haben, schreiben wir sofort in die Datenbank.
- Wenn der Filter sagt, dass wir ihn möglicherweise gesehen haben, machen wir eine genaue Prüfung in der Datenbank.
⚠️ Achtung: Versuchen Sie nicht, Dutzende Millionen Zeilen mit einem gewöhnlichen Set im Speicher zu deduplizieren. Auf einem typischen Laptop führt das zur Erschöpfung des Speichers und zum Absturz des Prozesses mitten in der Auslagerung. Verlagern Sie die Last auf die Festplatte über die Datenbank oder verwenden Sie einen Bloom-Filter.
Tipp: Wenn Sie Daten in Portionen auslagern und innerhalb einer Portion Duplikate möglich sind, deduplizieren Sie die Portion im Speicher mit einem gewöhnlichen Set vor dem Schreiben in die Datenbank. Die Portion ist klein, der Speicher leidet nicht, und es gehen weniger überflüssige Einfügungen an die Datenbank.
✅ Prüfung: Starten Sie die Deduplizierung auf einem großen Testsatz mit absichtlichen Wiederholungen. Prüfen Sie, dass die endgültige Anzahl eindeutiger Datensätze korrekt ist und der Speicherverbrauch des Prozesses stabil bleibt und nicht linear mit der Zeilenzahl wächst.
Schritt 5: Parallelität ohne Verluste
Ziel dieses Schritts. Die Auslagerung durch parallele Anfragen beschleunigen, ohne dabei eine Aufgabe zu verlieren und fehlgeschlagene korrekt zu wiederholen.
Aufgabenwarteschlange
Die Grundlage sicherer Parallelität ist eine Aufgabenwarteschlange. Sie zerlegen die Arbeit im Voraus in unabhängige Stücke. Zum Beispiel eine Liste von Seiten oder Bereichen. Sie legen sie in eine Warteschlange. Mehrere Worker nehmen Aufgaben aus der Warteschlange, führen sie aus und legen das Ergebnis ab. Wenn ein Worker abstürzt, kann seine Aufgabe zurück in die Warteschlange gelegt und einem anderen übergeben werden.
Begrenzung der Gleichzeitigkeit
Man kann nicht unendlich viele parallele Anfragen starten. Das überlastet die Quelle und Ihren Proxy. Der richtige Ansatz ist, die Anzahl gleichzeitiger Worker auf einen vernünftigen Wert zu begrenzen. Beginnen Sie mit einer kleinen Anzahl und erhöhen Sie sie, während Sie die Stabilität beobachten.
from concurrent.futures import ThreadPoolExecutor, as_completed
def run_parallel(tasks, worker, proxies, max_workers=5):
results = []
failed = []
with ThreadPoolExecutor(max_workers=max_workers) as pool:
future_map = {
pool.submit(worker, t, proxies): t for t in tasks
}
for future in as_completed(future_map):
task = future_map[future]
try:
results.append(future.result())
except Exception:
failed.append(task)
return results, failedWiederholung fehlgeschlagener Aufgaben
Die gesammelte Liste failed sind keine verlorenen Daten, sondern eine Liste dessen, was wiederholt werden muss. Nach dem ersten Durchlauf lassen Sie die fehlgeschlagenen Aufgaben erneut laufen. Normalerweise reicht das aus, um den Rest zu erledigen.
def run_with_retry(tasks, worker, proxies, rounds=3):
remaining = tasks
for _ in range(rounds):
done, remaining = run_parallel(remaining, worker, proxies)
if not remaining:
break
return remainingDie Rolle des Proxeon-Proxys bei der Parallelität. Bei paralleler Arbeit verteilt der Proxy die Verbindungen, was die Auslagerung stabiler und vorhersehbarer macht. Jeder Worker arbeitet über seine eigene Verbindung, und die Last konzentriert sich nicht an einem Punkt.
⚠️ Achtung: Bei parallelem Schreiben in eine Datei oder einen Checkpoint entstehen Datenrennen. Zwei Worker können den Zustand des anderen überschreiben. Schreiben Sie Ergebnisse nur in eine Datenbank mit Transaktionen oder verwenden Sie eine separate Datei pro Worker und sammeln Sie den Gesamt-Checkpoint in einem separaten Thread.
Tipp: Machen Sie Aufgaben klein und unabhängig. Wenn eine Aufgabe einen zu großen Bereich abdeckt, wirft ihr Abbruch viel Arbeit weg. Kleine Aufgaben lassen sich günstig und fast unbemerkt wiederholen.
✅ Prüfung: Starten Sie eine parallele Auslagerung und lassen Sie einen Teil der Worker absichtlich abstürzen. Nach den Wiederholungsrunden sollte die Liste remaining leer sein und der endgültige Datensatz vollständig. Vergleichen Sie die Anzahl der erhaltenen Datensätze mit der erwarteten.
Schritt 6: Wiederaufnahme nach langer Pause
Ziel dieses Schritts. Die Auslagerung korrekt fortsetzen, wenn zwischen den Versuchen viel Zeit vergangen ist, und verstehen, was in dieser Zeit veraltet sein kann.
Was mit der Zeit veraltet
Ein Abbruch für eine Minute und eine Pause für einen Tag sind unterschiedliche Situationen. Nach einer langen Pause kann ein Teil Ihres Zustands ungültig werden.
- Sitzung. Viele Dienste halten eine Sitzung nur begrenzte Zeit. Nach einer langen Pause vergisst der Server sie, und Anfragen beginnen, Autorisierungsfehler zurückzugeben.
- Zugriffstoken. API-Token haben oft eine Lebensdauer von Minuten oder Stunden. Ein abgelaufenes Token muss vor der Fortsetzung erneuert werden.
- Cursor. Manche Cursor leben nicht lange. Wenn der Cursor veraltet ist, muss man vom nächstgelegenen stabilen Punkt beginnen, zum Beispiel anhand der ID des letzten Datensatzes.
- Die Daten selbst. Während der Pause können in der Quelle neue Datensätze erschienen oder alte geändert worden sein. Das beeinflusst Offsets bei der seitenweisen Navigation per offset.
Strategie für sichere Wiederaufnahme
- Prüfen Sie beim Start das Alter des Checkpoints. Wenn er alt ist, seien Sie darauf vorbereitet, dass ein Teil des Zustands veraltet ist.
- Erneuern Sie das Zugriffstoken und erstellen Sie eine neue Sitzung vor der ersten Anfrage. Verlassen Sie sich nicht auf die alten.
- Bevorzugen Sie die Wiederaufnahme anhand der ID des letzten Datensatzes, nicht anhand der Seitennummer. Die ID ist stabil, während sich die Seitennummer verschiebt, wenn sich Daten ändern.
- Machen Sie eine Testanfrage mit dem gespeicherten Cursor. Wenn sie einen Fehler wegen ungültigen Cursors zurückgibt, wechseln Sie zur Wiederaufnahme anhand von last_id.
def resume(cp, fetch_by_id, fetch_by_cursor, proxies):
if cp["cursor"]:
try:
return fetch_by_cursor(cp["cursor"], proxies)
except CursorExpired:
pass
return fetch_by_id(cp["last_id"], proxies)Warum die Wiederaufnahme per ID zuverlässiger ist. Stellen Sie sich vor, Sie haben auf Seite 50 bei einer Sortierung nach Datum gestoppt. Während Ihrer Abwesenheit wurden neue Datensätze am Anfang hinzugefügt. Jetzt enthält Seite 50 ganz andere Daten, und Sie werden einen Teil der Datensätze überspringen. Die Wiederaufnahme anhand der ID des letzten Datensatzes leidet darunter nicht: Sie fordern einfach alles an, was größer als die gespeicherte ID ist.
Tipp: Speichern Sie im Checkpoint immer sowohl den Cursor als auch die ID des letzten Datensatzes. Der Cursor ist schneller, aber die ID ist Ihr Sicherungsseil für den Fall, dass der Cursor während einer langen Pause veraltet.
✅ Prüfung: Stoppen Sie die Auslagerung, warten Sie lange genug, damit Token oder Cursor veralten, und starten Sie erneut. Der Downloader sollte das Token erneuern, den veralteten Cursor erkennen und anhand der ID ohne Verlust und ohne Duplizierung von Datensätzen fortfahren.
Schritt 7: Fertiges Grundgerüst eines robusten Downloaders in Python
Ziel dieses Schritts. Alles Gelernte in ein einheitliches funktionierendes Gerüst zusammenzufügen, das Sie an Ihre Datenquelle anpassen.
Unten steht ein Gerüst, das Checkpoints, Deduplizierung über die Datenbank, Token-Erneuerung und Wiederaufnahme vereint. Die Funktionen zum Abrufen einer Seite ersetzen Sie durch Ihre eigenen, passend zur jeweiligen API.
import os
import json
import sqlite3
import requests
class ResilientLoader:
def __init__(self, cp_path, db_path, proxies):
self.cp_path = cp_path
self.proxies = proxies
self.conn = sqlite3.connect(db_path)
self.conn.execute(
"CREATE TABLE IF NOT EXISTS records ("
"dedup_key TEXT PRIMARY KEY, payload TEXT)"
)
self.conn.commit()
def load_cp(self):
if not os.path.exists(self.cp_path):
return {"cursor": None, "last_id": None, "count": 0}
with open(self.cp_path) as f:
return json.load(f)
def save_cp(self, cp):
tmp = self.cp_path + ".tmp"
with open(tmp, "w") as f:
json.dump(cp, f)
os.replace(tmp, self.cp_path)
def save_records(self, items):
rows = [(str(i["id"]), json.dumps(i)) for i in items]
self.conn.executemany(
"INSERT OR IGNORE INTO records "
"(dedup_key, payload) VALUES (?, ?)", rows
)
self.conn.commit()
def run(self, fetch_page):
cp = self.load_cp()
while True:
items, next_cursor = fetch_page(
cp["cursor"], cp["last_id"], self.proxies
)
if not items:
break
self.save_records(items)
cp["count"] += len(items)
cp["last_id"] = items[-1]["id"]
cp["cursor"] = next_cursor
self.save_cp(cp)
if next_cursor is None:
break
return cp["count"]Beispiel einer Funktion zum Abrufen einer Seite für Ihre Quelle. Hier implementieren Sie die Logik der Anfrage und das Parsen der Antwort.
def fetch_page(cursor, last_id, proxies):
params = {"limit": 100}
if cursor:
params["cursor"] = cursor
elif last_id:
params["after_id"] = last_id
r = requests.get(
"https://example-source/api/records",
params=params, proxies=proxies, timeout=60
)
r.raise_for_status()
data = r.json()
return data["items"], data.get("next_cursor")Der Start des gesamten Mechanismus sieht einfach aus.
proxies = {
"http": os.environ["PROXEON_URL"],
"https": os.environ["PROXEON_URL"],
}
loader = ResilientLoader("state.json", "out.db", proxies)
total = loader.run(fetch_page)
print("Gesamtanzahl Datensätze:", total)Tipp: Fügen Sie in die Schleife eine Protokollierung jeder hundert Datensätze ein: Zeit, Zähler, aktueller Cursor. So sehen Sie den Fortschritt und verstehen leicht, wenn die Auslagerung an einer Stelle hängt.
✅ Prüfung: Starten Sie das Gerüst auf einer echten Quelle, unterbrechen Sie in der Mitte, starten Sie erneut. Die endgültige Anzahl der Datensätze nach der Fortsetzung stimmt mit der vollständigen Anzahl der Quelldatensätze überein, und ein erneuter Lauf erhöht den Zähler der eindeutigen Zeilen nicht.
Ergebnisprüfung: Checkliste für robustes Auslagern
Gehen Sie diese Liste durch. Wenn alle Punkte erfüllt sind, ist Ihr Downloader wirklich robust.
- Die Datei-Wiederaufnahme wird ab dem fehlenden Byte fortgesetzt, nicht von vorn.
- Der Downloader behandelt korrekt den Fall, dass der Server den Range-Header ignoriert.
- Der Checkpoint wird nach jeder verarbeiteten Seite gespeichert, nicht nur am Ende.
- Der Checkpoint wird atomar über eine temporäre Datei und Umbenennen geschrieben.
- Ein erneuter Lauf erzeugt keine Duplikate im Endergebnis.
- Die Deduplizierung wächst nicht linear im Speicher mit der Zeilenzahl.
- Parallele Worker verlieren keine fehlgeschlagenen Aufgaben und wiederholen sie.
- Nach einer langen Pause wird das Token erneuert und ein veralteter Cursor durch Wiederaufnahme per ID ersetzt.
Wie man testet
- Starten Sie eine vollständige Auslagerung eines kleinen Satzes und merken Sie sich die Anzahl der Datensätze.
- Starten Sie erneut auf demselben Satz und stellen Sie sicher, dass sich die Anzahl nicht geändert hat.
- Unterbrechen Sie die Auslagerung an verschiedenen Stellen: am Anfang, in der Mitte, näher am Ende.
- Starten Sie nach jeder Unterbrechung neu und prüfen Sie, dass das Ergebnis gleich ist.
Erfolgsindikatoren. Die Anzahl eindeutiger Datensätze ist zwischen den Läufen stabil. Der Speicherverbrauch wächst nicht unkontrolliert. Die Wiederaufnahme setzt immer vom gespeicherten Punkt fort. Es gibt keine beschädigten Dateien nach der Fortsetzung.
Typische Fehler und ihre Lösungen
Problem: Die Datei öffnet sich nach der Fortsetzung nicht. Ursache: Die Daten wurden im Modus ab angehängt, obwohl der Server mit Code 200 geantwortet und die ganze Datei gesendet hat. Lösung: Prüfen Sie den Status der Antwort, wechseln Sie bei 200 zum vollständigen Überschreiben der Datei von vorn.
Problem: Nach einem Abbruch beginnt die Auslagerung mit der ersten Seite. Ursache: Der Checkpoint wurde nur am Ende der Arbeit gespeichert oder gar nicht. Lösung: Speichern Sie den Checkpoint nach jeder verarbeiteten Seite, direkt nach dem Schreiben der Daten.
Problem: Der Checkpoint ist nicht lesbar, das JSON ist beschädigt. Ursache: Das Programm ist während des Schreibens direkt in die Zieldatei abgestürzt. Lösung: Schreiben Sie in eine temporäre Datei und ersetzen Sie atomar über os.replace.
Problem: Im Endergebnis sind Duplikate. Ursache: Es gibt keinen Deduplizierungsschlüssel oder er ist instabil. Lösung: Legen Sie einen PRIMARY KEY auf einem zuverlässigen Schlüssel fest und verwenden Sie INSERT OR IGNORE.
Problem: Der Prozess stürzt bei großen Volumina wegen Speichermangels ab. Ursache: Alle gesehenen Schlüssel werden in einem Set im Speicher gehalten. Lösung: Verlagern Sie die Eindeutigkeitsprüfung in die Datenbank oder verwenden Sie einen Bloom-Filter.
Problem: Nach einer langen Pause geben Anfragen einen Autorisierungsfehler zurück. Ursache: Token oder Sitzung sind während der Pause veraltet. Lösung: Erneuern Sie das Token und erstellen Sie bei jedem Start eine neue Sitzung.
Problem: Nach einer Pause fehlen Datensätze oder sind doppelt vorhanden. Ursache: Die Wiederaufnahme erfolgte anhand der Seitennummer, und die Daten in der Quelle haben sich geändert. Lösung: Nehmen Sie anhand der ID des letzten Datensatzes wieder auf, nicht anhand des Offsets.
Problem: Parallele Worker verlieren einen Teil der Daten. Ursache: Mehrere Worker schrieben in einen Checkpoint und überschrieben sich gegenseitig. Lösung: Schreiben Sie Ergebnisse in eine Datenbank mit Transaktionen, nicht in eine gemeinsame Zustandsdatei.
Zusätzliche Möglichkeiten und Optimierung
Stapelweises Schreiben
Schreiben Sie nicht Zeile für Zeile in die Datenbank. Sammeln Sie einen Stapel von mehreren hundert Datensätzen und fügen Sie sie auf einmal über executemany ein. Das beschleunigt das Schreiben bei großen Volumina um ein Vielfaches.
Periodisches Festschreiben
Rufen Sie commit nicht bei jedem Schreiben auf, sondern alle paar hundert Zeilen. Zu häufiges commit verlangsamt die Datenbank. Zu seltenes riskiert, bei einem Abbruch mehr Daten zu verlieren. Finden Sie die Balance für Ihre Last.
Fortschrittsbericht
Fügen Sie eine Schätzung der verbleibenden Zeit hinzu. Wenn Sie die Verarbeitungsgeschwindigkeit der Seiten und die Gesamtzahl der Datensätze kennen, können Sie abschätzen, wie lange es noch dauert. Das ist praktisch für lange Auslagerungen.
Getrennte Speicher für Roh- und verarbeitete Daten
Speichern Sie Rohantworten getrennt von den geparsten Datensätzen. Wenn Sie später die Parsing-Logik ändern, müssen Sie die Daten nicht erneut herunterladen. Es reicht, die Rohantworten durch den neuen Parser zu schicken.
Tipp: Konfigurieren Sie den Proxeon-Proxy so, dass die Verbindungen während der gesamten Auslagerung stabil sind. Eine stabile Verbindung reduziert die Anzahl der Abbrüche, sodass Ihr Downloader seltener in die Wiederaufnahme geht und schneller arbeitet.
FAQ: Häufige Fragen zum robusten Auslagern
Wie erkenne ich, ob der Server die Datei-Wiederaufnahme unterstützt? Senden Sie eine HEAD-Anfrage und sehen Sie sich den Header Accept-Ranges an. Der Wert bytes bedeutet Unterstützung. Ein fehlender Header oder none bedeutet, dass Wiederaufnahme nicht möglich ist.
Was tun, wenn die API keinen Cursor liefert, sondern nur Seiten? Speichern Sie die Seitennummer und, wenn möglich, die ID des letzten Datensatzes. Nehmen Sie bevorzugt anhand der ID wieder auf, weil sich Seitennummern bei Datenänderungen verschieben.
Wie oft sollte der Checkpoint gespeichert werden? Nach jeder erfolgreich verarbeiteten und geschriebenen Seite. So verlieren Sie bei einem Abbruch höchstens eine Seite Arbeit, nicht die ganze Auslagerung.
Kann man ohne Datenbank deduplizieren? Bei kleinen Volumina ja, mit einem gewöhnlichen Set im Speicher. Bei Millionen Zeilen ist das wegen des Speichers gefährlich. Besser eine Datenbank mit PRIMARY KEY oder einen Bloom-Filter verwenden.
Was wählt man als Deduplizierungsschlüssel? Die natürliche eindeutige ID der Quelle, falls vorhanden. Wenn nicht, eine Kombination stabiler Felder. Im Extremfall den Hash des gesamten Datensatzes mit sortierten Schlüsseln.
Warum ist die Wiederaufnahme per ID zuverlässiger als per Seitennummer? Weil sich die Daten in der Quelle ändern können. Neue Datensätze verschieben die Seiten, und anhand der Nummer überspringen oder verdoppeln Sie Daten. Die ID hängt davon nicht ab.
Wie viele parallele Worker sollte man einstellen? Beginnen Sie mit einer kleinen Anzahl und erhöhen Sie sie, während Sie Stabilität und Limits der Quelle beobachten. Übermäßige Parallelität schadet mehr, als sie hilft.
Wie speichert man die Verbindungszeichenfolge zum Proxy sicher? In einer Umgebungsvariable, nicht im Code. Lesen Sie sie über os.environ. So gerät das Passwort nicht ins Versionskontrollsystem.
Was tun, wenn der Cursor während einer langen Pause veraltet? Fangen Sie den Fehler des ungültigen Cursors ab und wechseln Sie zur Wiederaufnahme anhand der gespeicherten ID des letzten Datensatzes.
Muss man die Integrität der heruntergeladenen Datei prüfen? Ja. Vergleichen Sie die tatsächliche Dateigröße mit dem Header Content-Length. Wenn der Server eine Prüfsumme liefert, prüfen Sie auch diese.
Fazit: Was Sie jetzt können und wohin Sie weitergehen
Sie haben den Weg von einer zerbrechlichen Auslagerung, die beim ersten Abbruch zusammenbricht, zu einem robusten Downloader zurückgelegt. Jetzt haben Sie alle Werkzeuge, damit ein Abbruch keine Katastrophe mehr ist, sondern eine normale Arbeitssituation.
Was Sie beherrschen. Die Wiederaufnahme von Dateien über den Range-Header mit Prüfung von Accept-Ranges. Das Speichern des Fortschritts über atomare Checkpoints. Idempotenz auf der Ebene des Deduplizierungsschlüssels. Deduplizierung bei Millionen von Zeilen ohne Speicheraufblähung. Parallelität mit Wiederholung fehlgeschlagener Aufgaben. Wiederaufnahme nach langer Pause mit Token-Erneuerung und Ersetzung veralteter Cursor. Und das Wichtigste: ein fertiges Python-Gerüst, das all das vereint.
Was als Nächstes zu tun ist. Nehmen Sie Ihre echte Datenquelle und passen Sie die Funktion zum Abrufen einer Seite daran an. Beginnen Sie mit einem kleinen Volumen, debuggen Sie die Wiederaufnahme bei Unterbrechungen und skalieren Sie dann. Konfigurieren Sie einen stabilen Proxeon-Proxy, damit die Verbindungen während der gesamten Auslagerung vorhersehbar sind.
Wohin Sie sich weiterentwickeln können. Studieren Sie stapelweises Schreiben und Bloom-Filter tiefer. Fügen Sie Fortschrittsmonitoring und Zeitschätzung hinzu. Trennen Sie die Speicherung von Roh- und verarbeiteten Daten. Nach und nach