Einleitung: Warum lange Exporte fast immer abbrechen – und das ist normal

Wenn Sie schon einmal über einen Proxy mehrere Millionen Datensätze oder eine Datei mit Dutzenden Gigabyte heruntergeladen haben, kennen Sie dieses Gefühl. Das Skript lief sechs Stunden, zeigte 83 Prozent an und stürzte dann mit einem Verbindungsfehler ab. Und alles, was Sie haben, ist eine unvollständige Datei und die Erkenntnis, dass Sie von vorne beginnen müssen.

Das Erste, was Sie akzeptieren müssen: Ein langer Download bricht immer ab. Nicht „manchmal“, nicht „bei schlechter Verbindung“, sondern immer, wenn er lange genug dauert. Es gibt Dutzende Gründe, und die meisten liegen außerhalb Ihrer Kontrolle:

  • Der Proxy ändert die externe IP-Adresse. Bei mobilen Proxys von Proxeon ist das Standardverhalten: Rotation per Timer oder auf Anfrage. Im Moment des IP-Wechsels wird die offene TCP-Verbindung getrennt, und der Quellserver sieht bereits einen anderen Client.
  • Der Quellserver schließt die Verbindung nach seinem Timeout, startet neu, spielt ein Update ein oder antwortet einfach mit einem 5xx-Fehler.
  • Ihr eigener Prozess startet neu: Systemupdate, voller Speicher, ein Fehler im Code bei einem untypischen Datensatz, versehentliches Ctrl+C.
  • Das Autorisierungstoken läuft ab, die Sitzung endet, der Cursor für den seitenbasierten Export wird ungültig.
  • Der Laptop geht in den Ruhemodus, das WLAN wechselt zu einem anderen Zugangspunkt, der Provider ändert die Route.

Es ist sinnlos, jeden dieser Gründe einzeln bekämpfen zu wollen. Der richtige Ansatz ist ein anderer: Gestalten Sie den Export so, dass ein Abbruch zu jedem Zeitpunkt nicht sechs Stunden, sondern nur eine Seite oder ein Dateistück kostet. Genau darum geht es in dieser Anleitung.

Was Sie am Ende erhalten

Nach dem Durcharbeiten der Anleitung verfügen Sie über:

  • Eine funktionierende Funktion zum Fortsetzen eines Dateidownloads ab der Mitte über das HTTP-Header-Feld Range, mit Prüfung, ob der Server dies unterstützt.
  • Ein klares Checkpoint-Schema für seitenbasierte Exporte: Was genau gespeichert wird und wo der Zustand liegt, damit er nicht mit dem Prozess verloren geht.
  • Idempotentes Schreiben von Ergebnissen, bei dem ein erneutes Laden derselben Seite keine Duplikate erzeugt.
  • Deduplizierung, die bei Millionen von Zeilen nicht den gesamten Arbeitsspeicher frisst.
  • Einen parallelen Loader mit Aufgabenwarteschlange, Wiederholung fehlgeschlagener Aufgaben und Begrenzung der Gleichzeitigkeit.
  • Ein fertiges Gerüst für einen robusten Loader in Python, das Sie in einer Stunde an Ihre Quelle anpassen.

Für wen diese Anleitung ist

Für Entwickler und Analysten, die bereits HTTP-Anfragen aus Python stellen können und mindestens einmal die Ergebnisse eines langen Downloads verloren haben. Mittleres Niveau: Grundlagen werden erklärt, aber wir bringen Ihnen nicht das Programmieren von Null bei. Fortgeschrittene Leser finden Abschnitte über die Speicherung von Hashes außerhalb des Speichers und über sicheres paralleles Schreiben in SQLite.

Was Sie vorher wissen sollten

  • Python auf dem Niveau von Funktionen, Schleifen, Dictionaries und Ausnahmebehandlung.
  • HTTP-Grundlagen: Was ist eine Methode, ein Header, ein Antwortcode, ein Body.
  • Eine allgemeine Vorstellung davon, wie man einen Proxy in der requests-Bibliothek einstellt.

Eine Klarstellung vorab: Antwortcodes und Wiederholungsstrategien mit exponentiellem Backoff werden hier nicht behandelt. Dazu gibt es einen separaten Artikel über Fehler 429 und Retries. In dieser Anleitung liegt der Fokus auf etwas anderem: auf dem Zustand des Exports und seiner Wiederaufnahme. Retries beantworten die Frage „Wann soll die Anfrage wiederholt werden?“, während wir die Frage beantworten: „Ab welcher Stelle soll die Arbeit fortgesetzt werden, nachdem die Wiederholungen ausgeschöpft sind und der Prozess beendet wurde?“

Wie viel Zeit Sie benötigen

Lesen und Ausführen der Beispiele an einer Testquelle: zwei bis drei Stunden. Anpassung des Gerüsts an Ihre echte API oder Ihren Dateiserver: ein bis zwei Stunden mehr, je nachdem, wie ungewöhnlich die Paginierung ist. Insgesamt ein Arbeitstag mit Puffer.

Vorbereitung und grundlegende Konzepte

Werkzeuge und Zugänge

  1. Installieren Sie Python 3.11 oder neuer. 2026 sind die Branches 3.12 und 3.13 aktuell; alle Beispiele wurden damit getestet. Version prüfen: Öffnen Sie das Terminal und geben Sie python --version ein. Wenn Sie 3.11 oder höher sehen, ist alles in Ordnung.
  2. Installieren Sie die requests-Bibliothek: pip install requests. Version 2.32 oder neuer reicht. Das Modul sqlite3 ist Teil der Python-Standardbibliothek, es muss nichts separat installiert werden.
  3. Holen Sie sich Zugang zum Proxy. Öffnen Sie Ihr Proxeon-Kundenkonto, wählen Sie den gewünschten Kanal und kopieren Sie vier Werte: Host, Port, Benutzername und Passwort. Normalerweise sind sie in einer Zeile zusammengefasst, z. B. http://USER:PASS@HOST:PORT. Diese Zeile verwenden Sie überall weiter.
  4. Legen Sie die Proxy-Zeile in eine Umgebungsvariable, nicht in den Code. Unter Linux und macOS: export PROXY_URL=http://USER:PASS@HOST:PORT. Unter Windows PowerShell: $env:PROXY_URL='http://USER:PASS@HOST:PORT'. So committen Sie das Passwort nicht versehentlich ins Repository.
  5. Prüfen Sie, ob der Proxy antwortet. Führen Sie im Terminal aus: curl -x $PROXY_URL -I https://api.example.com/ und setzen Sie die Adresse Ihrer Quelle ein. Sie sollten eine Zeile mit einem Antwortcode sehen, z. B. HTTP/2 200. Bei einem Proxy-Authentifizierungsfehler 407 prüfen Sie Benutzername und Passwort.

Systemanforderungen

Jeder Rechner mit 2 GB freiem Arbeitsspeicher und einer Festplatte, auf der das Ergebnis des Exports plus 20 Prozent Reserve für SQLite-Indizes Platz finden. Wenn Sie Millionen von Datensätzen laden möchten, ist die Festplatte wichtiger als der Arbeitsspeicher: Der gesamte Ansatz beruht darauf, dass der Zustand auf der Festplatte lebt, nicht in Prozessvariablen.

Sicherungskopien

Die Zustandsdatei, die Sie unten erstellen (in den Beispielen export.sqlite), wird das wertvollste Artefakt der gesamten Arbeit. Gewöhnen Sie sich an, sie vor Experimenten mit dem Code zu kopieren: cp export.sqlite export.sqlite.bak. Einmal wird Ihnen das einen ganzen Tag Exportarbeit retten.

Achtung: Bearbeiten Sie die SQLite-Datei niemals manuell, während der Loader läuft. Selbst das Lesen aus einem anderen Programm im falschen Modus kann das Schreiben blockieren und den Prozess zum Absturz bringen. Wenn Sie den Zustand sehen möchten, stoppen Sie den Loader oder verwenden Sie den WAL-Modus, den wir im Abschnitt über Parallelität erklären.

Wichtige Begriffe einfach erklärt

  • Checkpoint – eine auf der Festplatte gespeicherte Markierung „Bis hierhin ist alles exportiert und geschrieben“. Nach einem Abbruch liest der Loader den Checkpoint und macht von dort weiter.
  • Cursor – eine undurchsichtige Zeichenkette, die die API zusammen mit einer Seite liefert und die übergeben werden muss, um die nächste Seite zu erhalten. Sie erstellen oder analysieren den Cursor nicht selbst.
  • Seitenbasierter Export per Offset – wenn Sie „Seite 37 mit 500 Datensätzen“ anfordern. Einfaches Schema, aber wenn neue Datensätze zur Quelle hinzugefügt werden, verschieben sich die Seiten und es entstehen Duplikate oder Lücken.
  • Seitenbasierter Export per Schlüssel (Keyset) – wenn Sie „alles mit einer ID größer als 184203, sortiert nach ID“ anfordern. Das robusteste Schema für die Wiederaufnahme, wenn die Quelle es unterstützt.
  • Idempotenz – die Eigenschaft einer Operation, bei der ihre wiederholte Ausführung dasselbe Ergebnis liefert wie eine einmalige. Sie haben eine Seite zweimal geschrieben, aber in der Datenbank liegt sie nur einmal.
  • Deduplizierungsschlüssel – der Wert, anhand dessen zwei Datensätze als derselbe gelten. Ideal ist die ID aus der Quelle; wenn es keine gibt, wird der Schlüssel als Hash stabiler Felder berechnet.
  • Range-Header – eine Möglichkeit, einen HTTP-Server zu bitten, nicht die ganze Datei, sondern einen Teil davon zu liefern, z. B. Bytes von 1048576 bis zum Ende.
  • At-least-once-Semantik – die Garantie, dass jeder Datensatz mindestens einmal abgerufen wird. Wiederholungen sind möglich, aber keine Auslassungen. Genau das erreichen Sie mit dieser Anleitung, und Duplikate entfernt die Deduplizierung.

Das Hauptprinzip

Alle sieben Schritte unten laufen auf eine Idee hinaus: Jede Arbeitseinheit muss atomar und wiederholbar sein. Eine Arbeitseinheit ist entweder ein Dateistück, eine API-Seite oder eine Aufgabe aus der Warteschlange. Atomar bedeutet: Ergebnis und Markierung über seinen Abschluss werden zusammen gespeichert. Wiederholbar bedeutet: Wenn die Einheit zweimal ausgeführt wird, geht nichts kaputt. Wenn beide Eigenschaften erfüllt sind, wird ein Abbruch an jeder Stelle sicher.

Schritt 1: Datei-Download per HTTP über den Range-Header fortsetzen

Ziel dieser Phase: Eine große Datei über einen Proxy so herunterladen, dass der Download nach einem Abbruch ab dem Byte fortgesetzt wird, an dem er aufgehört hat, und nicht von vorne.

Wie es funktioniert

Das HTTP-Protokoll erlaubt es dem Client, einen Teil einer Ressource anzufordern. Dazu wird der Header Range: bytes=START- zur Anfrage hinzugefügt, wobei START der Offset in Bytes ist. Wenn der Server Teilanfragen unterstützt, antwortet er mit dem Code 206 Partial Content und dem Header Content-Range: bytes START-ENDE/GESAMT. Wenn nicht, ignoriert er Range und liefert die gesamte Datei mit dem Code 200. Ihre Aufgabe ist es, diese Fälle zu unterscheiden.

Um im Voraus zu erfahren, ob der Server die Fortsetzung unterstützt, hilft eine HEAD-Anfrage: Sie liefert nur die Header ohne Body. Achten Sie auf Accept-Ranges: bytes. Der Wert none oder ein fehlender Header bedeutet normalerweise, dass keine Fortsetzung möglich ist, obwohl einige Server Range trotzdem korrekt verarbeiten. Deshalb prüfen wir den endgültigen Fall anhand des Antwortcodes.

Schritt-für-Schritt-Anleitung

  1. Stellen Sie eine HEAD-Anfrage über den Proxy und speichern Sie die Header Accept-Ranges, Content-Length, ETag und Last-Modified. ETag benötigen Sie, um festzustellen, ob sich die Datei auf dem Server zwischen Ihren Versuchen geändert hat.
  2. Sehen Sie nach, wie viele Bytes bereits in der lokalen Datei liegen. Wenn keine Datei vorhanden ist, nehmen Sie null an.
  3. Wenn die lokale Größe bereits der Content-Length entspricht, ist die Datei vollständig, es ist nichts zu tun.
  4. Wenn die lokale Größe größer als null ist und der Server die Unterstützung von Bereichen angibt, fügen Sie den Range-Header mit der aktuellen Größe zur Anfrage hinzu. Fügen Sie außerdem If-Range mit dem gespeicherten ETag hinzu: Dann liefert der Server eine Teilantwort nur, wenn sich die Datei nicht geändert hat, andernfalls liefert er die gesamte Datei mit dem Code 200.
  5. Senden Sie unbedingt Accept-Encoding: identity. Ohne diesen Header kann der Server eine On-the-fly-Komprimierung anwenden, und die Byte-Offsets stimmen nicht mehr mit Ihrer Datei überein.
  6. Öffnen Sie die lokale Datei im Anhängemodus ab, wenn Sie 206 erhalten haben, oder im Überschreibmodus wb, wenn Sie 200 erhalten haben.
  7. Lesen Sie den Body als Stream in Blöcken von 256 KB und schreiben Sie auf die Festplatte. Laden Sie nicht die gesamte Antwort in den Speicher.
  8. Vergleichen Sie nach Abschluss die endgültige Größe mit Content-Length. Wenn sie nicht übereinstimmt, wurde die Verbindung stillschweigend unterbrochen, und ein weiterer Durchgang ist nötig.

Funktionierender Code

import os
import requests

PROXY_URL = os.environ['PROXY_URL'] # строка из личного кабинета Proxeon
PROXIES = {'http': PROXY_URL, 'https': PROXY_URL}


def probe(url):
r = requests.head(url, proxies=PROXIES, allow_redirects=True, timeout=30,
 headers={'Accept-Encoding': 'identity'})
r.raise_for_status()
return {
'ranges': r.headers.get('Accept-Ranges', 'none').lower(),
'length': int(r.headers.get('Content-Length', 0) or 0),
'etag': r.headers.get('ETag'),
}


def download_resumable(url, path):
meta = probe(url)
have = os.path.getsize(path) if os.path.exists(path) else 0
if meta['length'] and have >= meta['length']:
print('файл уже полный:', have, 'байт')
return True
headers = {'Accept-Encoding': 'identity'}
if have > 0 and meta['ranges'] == 'bytes':
headers['Range'] = 'bytes=%d-' % have
if meta['etag']:
headers['If-Range'] = meta['etag']
with requests.get(url, headers=headers, proxies=PROXIES, stream=True,
 timeout=(30, 120)) as r:
if r.status_code == 206:
expected = 'bytes %d-' % have
if not r.headers.get('Content-Range', '').startswith(expected):
raise IOError('сервер отдал не тот диапазон: ' + r.headers.get('Content-Range', ''))
mode = 'ab'
elif r.status_code == 200:
print('сервер отдаёт файл целиком, начинаем с нуля')
mode = 'wb'
have = 0
elif r.status_code == 416:
raise IOError('запрошенный диапазон вне файла, проверьте локальный размер')
else:
r.raise_for_status()
with open(path, mode) as f:
for chunk in r.iter_content(chunk_size=256 * 1024):
if chunk:
f.write(chunk)
have += len(chunk)
if meta['length'] and have != meta['length']:
print('обрыв: получено %d из %d' % (have, meta['length']))
return False
return True


def download_until_done(url, path, max_rounds=50):
for i in range(max_rounds):
try:
if download_resumable(url, path):
return
except (requests.ConnectionError, requests.Timeout, IOError) as e:
print('заход %d прерван: %s' % (i + 1, type(e).__name__))
# пауза перед следующим заходом: стратегия описана в статье про 429 и ретраи
raise RuntimeError('не удалось докачать файл за %d заходов' % max_rounds)


if __name__ == '__main__':
download_until_done('https://files.example.com/export-2026.csv.gz', 'export-2026.csv.gz')

Achten Sie auf die Funktion download_until_done: Sie enthält keine Wartelogik zwischen den Versuchen. Das ist Absicht. Fügen Sie dort Ihre Pausenstrategie aus dem Artikel über Retries ein; hier ist nur die Schleife „Größe geprüft, nachgeladen, erneut geprüft“ wichtig.

Tipp: Wenn die Datei als Archiv ausgeliefert wird, entpacken Sie sie nicht während des Nachladens. Holen Sie zuerst die vollständige Datei, prüfen Sie die Größe und vergleichen Sie sie, falls der Server eine Prüfsumme liefert. Erst dann entpacken Sie. Ein teilweise heruntergeladenes gzip sieht wie eine beschädigte Datei aus, und Sie verschwenden Zeit mit der Suche nach einem nicht existierenden Fehler.

Erwartetes Ergebnis

Prüfung: Starten Sie das Skript mit einer Datei von mindestens 200 MB und unterbrechen Sie es nach zehn Sekunden mit Ctrl+C. Schauen Sie sich die Größe der lokalen Datei an, z. B. 41.943.040 Bytes. Starten Sie das Skript erneut. In der Konsole darf nicht die Zeile „начинаем с нуля“ erscheinen, und die Dateigröße muss weiter wachsen, nicht zurückgesetzt werden. Am Ende muss die endgültige Größe genau mit der Content-Length aus der HEAD-Anfrage übereinstimmen.

Mögliche Probleme

  • Der Server liefert immer 200 statt 206. Das bedeutet, dass die Fortsetzung nicht unterstützt wird. Der einzige Ausweg für eine solche Quelle ist, die Datei in einem Durchgang mit einem großen Timeout herunterzuladen oder nach einem alternativen Format für den Export in Teilen zu suchen, z. B. einer Aufteilung nach Datum.
  • Kein Content-Length-Header. Der Server liefert die Datei im Chunked-Modus ohne Angabe der Größe. Die Vollständigkeit kann nicht anhand der Größe geprüft werden, und ein Nachladen ist ebenfalls nicht möglich: Range erfordert bekannte Offsets. Verhandeln Sie mit der Quelle oder verwenden Sie eine Prüfsumme, falls eine veröffentlicht wird.
  • Antwort 416 beim ersten Versuch. Die lokale Datei ist größer als die Datei auf dem Server. Die Datei auf dem Server hat sich geändert und ist kürzer geworden. Löschen Sie die lokale Datei und beginnen Sie von vorne.
  • Die Größe stimmt überein, aber die Datei ist defekt. Wahrscheinlich gab es irgendwo in der Mitte einen Abbruch mit Code 200 und einem Schreiben von Null, gefolgt von einem Anhängen. Erstellen Sie die Datei neu. Um dies zu vermeiden, speichern Sie ETag in einer separaten Datei daneben und vergleichen Sie es vor jedem Versuch.

Schritt 2: Checkpoints für seitenbasierte Exporte

Ziel dieser Phase: Den Exportzustand so speichern, dass der Prozess nach jedem Absturz mit der letzten erfolgreich geschriebenen Seite fortfährt.

Was gespeichert werden sollte

Der minimale Checkpoint hängt vom Paginierungstyp der Quelle ab. Betrachten wir drei Fälle.

  1. Paginierung per Cursor. Die API liefert zusammen mit den Daten ein Feld wie next_cursor. Speichern Sie genau dieses. Das ist der einfachste Fall: Der Cursor enthält bereits alles, was der Server zur Fortsetzung benötigt.
  2. Paginierung per Seitennummer oder Offset. Speichern Sie die Nummer der letzten vollständig geschriebenen Seite und die Seitengröße. Denken Sie daran, dass sich die Offsets verschieben, wenn während des Exports Datensätze zur Quelle hinzugefügt werden. Deshalb ist die Deduplizierung aus Schritt 4 unerlässlich.
  3. Paginierung per Schlüssel. Speichern Sie die ID des letzten geschriebenen Datensatzes. Bei der Wiederaufnahme fordern Sie alles an, was größer als diese ID ist. Das Schema fürchtet weder Einfügungen noch lange Pausen.

Unabhängig vom Paginierungstyp sollten Sie dem Checkpoint Hilfsfelder hinzufügen: die ID des letzten Datensatzes (auch für das Cursor-Schema, als Backup-Anker, falls der Cursor ungültig wird), Zähler für Seiten und Zeilen zur Fortschrittskontrolle, die Startzeit des Exports und die Zeit der letzten Aktualisierung.

Wo der Zustand gespeichert wird

Es gibt zwei funktionierende Varianten, und beide sind besser als Variablen im Speicher.

Variante A: JSON-Datei mit atomarem Ersetzen

Geeignet, wenn die Ergebnisse in separate Dateien geschrieben werden und nicht in eine Datenbank. Die Hauptfalle: Wenn Sie den Zustand direkt in die Zieldatei schreiben und der Prozess mitten im Schreiben abstürzt, erhalten Sie eine abgeschnittene JSON-Datei, die nicht gelesen werden kann. Die Lösung: Schreiben Sie in eine temporäre Datei daneben und benennen Sie sie über die Hauptdatei um. Der Umbenennungsvorgang ist innerhalb eines Dateisystems atomar.

import json
import os
import tempfile


def save_state(path, state):
directory = os.path.dirname(os.path.abspath(path))
fd, tmp = tempfile.mkstemp(dir=directory, prefix='.state-')
with os.fdopen(fd, 'w', encoding='utf-8') as f:
json.dump(state, f, ensure_ascii=False)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)


def load_state(path, default):
if not os.path.exists(path):
return dict(default)
with open(path, encoding='utf-8') as f:
return json.load(f)

Variante B: Tabelle in SQLite neben den Daten

Die bevorzugte Variante, wenn Sie Datensätze in eine Datenbank legen. Der Checkpoint wird in derselben Transaktion aktualisiert wie das Einfügen der Zeilen einer Seite. Entweder werden sowohl die Daten als auch die Markierung geschrieben oder nichts. Eine Lücke zwischen „Daten sind da, Markierung fehlt“ gibt es grundsätzlich nicht.

import sqlite3

con = sqlite3.connect('export.sqlite')
con.executescript('''
CREATE TABLE IF NOT EXISTS records(
id TEXT PRIMARY KEY,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS checkpoint(
job TEXT PRIMARY KEY,
cursor TEXT,
last_id TEXT,
pages INTEGER NOT NULL DEFAULT 0,
updated_at TEXT
);
''')


def commit_page(job, rows, next_cursor, pages, ts):
with con: # одна транзакция на страницу
con.executemany(
'INSERT OR IGNORE INTO records(id, payload, fetched_at) VALUES (?, ?, ?)',
[(str(r['id']), json.dumps(r, ensure_ascii=False), ts) for r in rows])
con.execute(
'INSERT INTO checkpoint(job, cursor, last_id, pages, updated_at) VALUES (?, ?, ?, ?, ?) '
'ON CONFLICT(job) DO UPDATE SET cursor=excluded.cursor, last_id=excluded.last_id, '
'pages=excluded.pages, updated_at=excluded.updated_at',
(job, next_cursor, str(rows[-1]['id']), pages, ts))

Reihenfolge der Operationen

Merken Sie sich die Regel: Zuerst die Daten, dann der Checkpoint, und möglichst in einer Transaktion. Wenn keine Transaktion möglich ist (z. B. wenn Daten in Dateien geschrieben werden), ist die Reihenfolge genau so: Datei der Seite schreiben, auf die Festplatte synchronisieren, dann den Zustand aktualisieren. Bei einem Absturz zwischen diesen beiden Aktionen erhalten Sie ein erneutes Laden einer Seite, was dank Schritt 3 sicher ist. Die umgekehrte Reihenfolge führt zu einer ausgelassenen Seite, und das ist ein Datenverlust.

Tipp: Speichern Sie im Checkpoint nicht den aktuellen Cursor, sondern den Cursor der nächsten Seite, den der Server zurückgegeben hat. Dann fordern Sie bei der Wiederaufnahme sofort das an, was Sie noch nicht haben, ohne eine zusätzliche Anfrage für die bereits erhaltene Seite.

Erwartetes Ergebnis

Prüfung: Starten Sie den Export, warten Sie zehn Seiten ab und beenden Sie den Prozess gewaltsam. Öffnen Sie die Datenbank mit sqlite3 export.sqlite und führen Sie SELECT pages, last_id FROM checkpoint; aus. Sie sollten die Zahl 10 und die ID sehen. Führen Sie dann SELECT count(*) FROM records; aus und stellen Sie sicher, dass die Anzahl der Datensätze zehn Seitengrößen entspricht. Starten Sie den Loader erneut: Die erste Meldung in der Konsole sollte etwa „старт: страниц 10“ lauten.

Mögliche Probleme

  • Fehler „database is locked“. Ein anderer Prozess hält die Verbindung. Schließen Sie alle sqlite3-Fenster und andere Tools, die die Datei geöffnet haben. Für Multithread-Arbeit aktivieren Sie den WAL-Modus, siehe Schritt 5.
  • Der Cursor wurde gespeichert, die Daten aber nicht. Sie haben den Checkpoint außerhalb einer Transaktion mit den Daten aktualisiert. Kehren Sie zum obigen Code zurück und stellen Sie sicher, dass beide Operationen innerhalb eines with con:-Blocks liegen.
  • Die JSON-Zustandsdatei ist leer oder beschädigt. Sie haben direkt in die Datei geschrieben, ohne temporäre Datei und Ersetzen. Verwenden Sie die Funktion save_state vollständig.

Schritt 3: Idempotentes Schreiben von Ergebnissen

Ziel dieser Phase: Dafür sorgen, dass die erneute Verarbeitung jeder Seite keine Duplikate erzeugt und die Daten nicht beschädigt.

Warum Wiederholungen unvermeidlich sind

Nach Schritt 2 haben Sie bereits das Szenario gesehen, bei dem eine Seite zweimal geschrieben wird: Der Prozess stürzte nach dem Einfügen der Daten, aber vor der Aktualisierung des Checkpoints ab. Darüber hinaus entstehen Wiederholungen durch Offset-Paginierung bei Änderungen der Quelle, durch parallele Worker, die nach einem Neustart dieselbe Aufgabe erhalten haben, und einfach durch manuelles Neustarten „nur zur Sicherheit“. Es ist sinnlos, Wiederholungen auf der Anforderungsseite zu bekämpfen. Richtig ist, das Schreiben selbst so zu gestalten, dass eine Wiederholung harmlos ist.

Wahl des Deduplizierungsschlüssels

  1. Es gibt eine ID in der Quelle. Verwenden Sie diese. Das ist das Feld id, uuid, order_number oder ein ähnliches, das die Quelle als eindeutig garantiert. Wenn Sie aus mehreren Quellen in eine Tabelle exportieren, verwenden Sie einen zusammengesetzten Schlüssel: Name der Quelle plus ID.
  2. Keine ID, aber eine Gruppe von Feldern, die zusammen den Datensatz bestimmen. Zum Beispiel für eine Preislistenzeile: Artikelnummer plus Lager plus Datum. Bilden Sie den Schlüssel aus diesen Feldern, nachdem Sie sie normalisiert haben: Bringen Sie Zeichenketten auf eine einheitliche Groß-/Kleinschreibung, entfernen Sie Leerzeichen an den Rändern, konvertieren Sie Daten in ein einheitliches Format.
  3. Nichts Stabiles vorhanden. Dann wird der Schlüssel zum Hash des gesamten Datensatzes nach der Kanonisierung. Dieser Fall wird in Schritt 4 ausführlich behandelt. Beachten Sie: Wenn die Quelle einen Datensatz ändert (z. B. den Preis aktualisiert), ändert sich der Hash, und Sie erhalten beide Versionen. Manchmal ist das genau das, was Sie wollen, manchmal nicht.

Achtung: Verwenden Sie niemals die laufende Nummer einer Zeile in der Antwort oder die Seitennummer als Schlüssel. Diese Werte ändern sich bei jeder Änderung der Quelle, und die Deduplizierung wird zum Duplikat-Generator.

Idempotentes Einfügen in die Datenbank

In SQLite und den meisten relationalen Datenbanken gibt es ein Konstrukt, das entweder den Konflikt mit dem Primärschlüssel ignoriert oder die vorhandene Zeile aktualisiert. Die erste Variante INSERT OR IGNORE haben Sie bereits in Schritt 2 gesehen. Sie eignet sich, wenn die Datensätze unveränderlich sind. Die zweite Variante ist nötig, wenn die Quelle Datensätze aktualisieren kann und Sie die aktuelle Version benötigen:

def upsert_rows(con, rows, ts):
con.executemany(
'INSERT INTO records(id, payload, fetched_at) VALUES (?, ?, ?) '
'ON CONFLICT(id) DO UPDATE SET payload=excluded.payload, fetched_at=excluded.fetched_at',
[(str(r['id']), json.dumps(r, ensure_ascii=False), ts) for r in rows])

Idempotentes Schreiben in Dateien

Wenn das Ergebnis in Dateien statt in einer Datenbank liegen soll, wenden Sie dasselbe Prinzip an: Eine Seite entspricht einer Datei mit deterministischem Namen. Der Name hängt von den Parametern der Seite ab, nicht von der Zeit oder einem Zähler. Vor dem Download prüfen Sie, ob die endgültige Datei existiert; wenn ja, überspringen Sie die Seite. Schreiben Sie in einen temporären Namen und benennen Sie ihn nach Abschluss um, wie in der Funktion save_state.

def page_path(base_dir, job, cursor_or_page):
safe = str(cursor_or_page).replace('/', '_').replace(':', '_')[:120]
return os.path.join(base_dir, job, 'page-%s.jsonl' % safe)


def write_page_idempotent(path, rows):
if os.path.exists(path):
return False # страница уже есть, повторная запись не нужна
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + '.part'
with open(tmp, 'w', encoding='utf-8') as f:
for r in rows:
f.write(json.dumps(r, ensure_ascii=False))
f.write(chr(10))
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
return True

Dateien mit der Endung .part, die nach einem Absturz übrig bleiben, können beim Start bedenkenlos gelöscht werden: Sie sind per Definition unvollständig.

Tipp: Verlassen Sie sich bei einem Export per Offset nicht nur darauf, dass „die Datei existiert, also ist die Seite fertig“. Prüfen Sie zusätzlich, dass die Anzahl der Zeilen in der Datei der Seitengröße entspricht (außer bei der letzten). Eine leere oder zu kurze Datei bei vorhandenem Namen sollte besser neu geladen werden.

Erwartetes Ergebnis

Prüfung: Rufen Sie die Schreibfunktion für dieselbe Seite dreimal hintereinander auf. Führen Sie dann SELECT count(*) FROM records; aus. Die Zahl muss der Größe einer Seite entsprechen, nicht der dreifachen. Bei der Dateivariante darf im Verzeichnis genau eine Seitendatei und keine .part-Datei liegen.

Schritt 4: Deduplizierung von Ergebnissen ohne Speicheraufblähung

Ziel dieser Phase: Doppelte Datensätze im Fluss von Millionen von Zeilen herausfiltern, ohne alle Schlüssel im Arbeitsspeicher zu halten.

Hash des Datensatzes

Wenn ein Datensatz keine ID hat, wird der Hash seines Inhalts zum Schlüssel. Damit gleiche Datensätze denselben Hash ergeben, muss der Inhalt kanonisiert werden: Schlüssel des Dictionaries sortieren, überflüssige Leerzeichen entfernen, Trennzeichen festlegen. Andernfalls erhält derselbe Datensatz mit anderer Feldreihenfolge einen anderen Hash.

import hashlib
import json


def record_key(rec, fields=None):
src = rec if fields is None else {k: rec.get(k) for k in fields}
canon = json.dumps(src, sort_keys=True, ensure_ascii=False, separators=(',', ':'))
return hashlib.blake2b(canon.encode('utf-8'), digest_size=16).digest()

Die Funktion gibt 16 Bytes zurück. Das reicht: Die Wahrscheinlichkeit einer zufälligen Übereinstimmung für Hunderte Millionen Datensätze ist vernachlässigbar gering. Der Parameter fields ermöglicht es, den Hash nur über stabile Felder zu berechnen und beispielsweise die Zeit der letzten Aktualisierung auszuschließen, die sich bei jeder Anfrage ändert.

Warum eine Menge im Speicher bei Millionen nicht funktioniert

Das Erste, was einem in den Sinn kommt: ein seen = set() anlegen und Schlüssel hineinlegen. Rechnen wir nach. Ein bytes-Objekt der Länge 16 belegt in Python etwa 49 Bytes plus die Daten selbst, insgesamt etwa 65 Bytes. Ein Slot in der Menge fügt mit Füllfaktor noch etwa 30 Bytes hinzu. Wir erhalten etwa 95 Bytes pro Schlüssel. Bei 10 Millionen Datensätzen sind das etwa 950 MB, bei 50 Millionen fast 5 GB. Und das Wichtigste: Nach einem Prozessneustart ist die Menge leer, und die gesamte Deduplizierung beginnt von vorne.

Drei Möglichkeiten, den Speicher nicht aufzublähen

  1. Schlüssel in der Datenbank selbst speichern. Der einfachste und zuverlässigste Weg. Wenn der Schlüssel der Primärschlüssel der Tabelle records ist, ist die Deduplizierung bereits durch das Konstrukt INSERT OR IGNORE aus Schritt 3 erledigt. Der Index lebt auf der Festplatte, überlebt Neustarts, und SQLite cached heiße Indexseiten selbst. Für 10 Millionen 16-Byte-Schlüssel belegt der Index etwa 400–500 MB auf der Festplatte, aber nicht im Speicher.
  2. Separate Tabelle gesehener Schlüssel ohne rowid. Nötig, wenn Sie die Daten selbst nicht in SQLite schreiben, sondern z. B. in Dateien. Dann wird SQLite nur als kompakte Menge auf der Festplatte verwendet.
  3. Komprimierte Menge im Speicher als Vorfilter. Fortgeschrittene Variante: Hash auf 8 Bytes kürzen und als Ganzzahl in einem sortierten Array speichern oder einen Bloom-Filter verwenden. Der Speicher wird um ein Mehrfaches reduziert, aber es entsteht die Möglichkeit eines False Positives. Deshalb wird ein solcher Vorfilter nur verwendet, um offensichtlich neue Datensätze schnell auszusortieren, während die endgültige Prüfung trotzdem in der Datenbank erfolgt.

Implementierung einer Menge auf der Festplatte

class DiskSeen:
def __init__(self, con):
self.con = con
con.execute('CREATE TABLE IF NOT EXISTS seen(key BLOB PRIMARY KEY) WITHOUT ROWID')

def filter_new(self, rows):
keyed = [(record_key(r), r) for r in rows]
keys = [k for k, _ in keyed]
placeholders = ','.join('?' * len(keys))
known = {row[0] for row in self.con.execute(
'SELECT key FROM seen WHERE key IN (%s)' % placeholders, keys)}
fresh = [(k, r) for k, r in keyed if k not in known]
# дедупликация внутри самой страницы
unique = {}
for k, r in fresh:
unique.setdefault(k, r)
return unique

def remember(self, keys):
self.con.executemany('INSERT OR IGNORE INTO seen(key) VALUES (?)', [(k,) for k in keys])

Rufen Sie filter_new vor dem Schreiben der Seite auf und remember innerhalb derselben Transaktion wie das Schreiben der Daten und des Checkpoints. Dann sind nach einem Absturz die Menge der gesehenen Schlüssel, die Daten und die Fortschrittsmarkierung immer konsistent.

Tipp: Eine Abfrage mit IN auf 500 Werte wird über den Index in Millisekunden ausgeführt. Prüfen Sie Schlüssel nicht einzeln in einer Schleife: Das ist aufgrund des Aufwands pro Abfrage um ein Vielfaches langsamer.

Erwartetes Ergebnis

Prüfung: Erstellen Sie eine Testseite mit 500 Datensätzen, von denen 100 innerhalb der Seite doppelt vorkommen und weitere 100 bereits in der Tabelle seen liegen. Die Funktion filter_new muss genau 300 Datensätze zurückgeben. Nach einem Prozessneustart müssen dieselben 500 Datensätze null neue ergeben.

Mögliche Probleme

  • Duplikate kommen trotzdem durch. Prüfen Sie die Kanonisierung: Wahrscheinlich gibt es im Datensatz ein Feld mit der Anfragezeit oder einer zufälligen Reihenfolge von Listenelementen. Schließen Sie es über den Parameter fields aus oder sortieren Sie verschachtelte Listen vor dem Hashen.
  • Langsames Einfügen nach mehreren Millionen Zeilen. Der Index passt nicht mehr in den Cache. Erhöhen Sie den SQLite-Cache mit PRAGMA cache_size=-200000 (das sind 200 MB) und stellen Sie sicher, dass Einfügungen in Stapeln in einer Transaktion pro Seite erfolgen, nicht Zeile für Zeile.

Schritt 5: Parallelität ohne Verluste

Ziel dieser Phase: Den Export durch mehrere gleichzeitige Worker über den Proxy beschleunigen, ohne dass der Absturz eines beliebigen Workers Aufgaben verliert oder die Datenbank beschädigt.

Wann Parallelisierung möglich ist und wann nicht

Die Cursor-Paginierung ist von Natur aus sequenziell: Der nächste Cursor ist erst nach Erhalt der vorherigen Seite bekannt. Sie kann nicht direkt parallelisiert werden. Aber fast immer kann der Export in unabhängige Shards aufgeteilt werden: nach Tagen, nach Kategorien, nach Regionen, nach den ersten Zeichen der ID. Jeder Shard wird sequenziell mit eigenem Checkpoint exportiert, und die Shards laufen parallel. Die Paginierung nach Offset und nach Schlüssel mit bekannten Grenzen lässt sich direkt parallelisieren: Aufgaben wie „Seiten 1 bis 100“ oder „IDs von 0 bis 100000“.

Aufgabenwarteschlange auf der Festplatte

Eine Warteschlange im Speicher stirbt mit dem Prozess. Deshalb leben Aufgaben in einer Tabelle mit Status:

  • pending – wartet auf Ausführung;
  • running – von einem Worker übernommen;
  • done – ausgeführt und geschrieben;
  • failed – Versuche erschöpft, erfordert menschliche Aufmerksamkeit.

Beim Start setzt der Loader zuerst alle Aufgaben von running zurück auf pending: Wenn sie in diesem Status hängen, ist der vorherige Prozess mitten in der Arbeit gestorben. Dann nehmen die Worker pending-Aufgaben ab.

Eine Schreibstelle

SQLite erlaubt viele gleichzeitige Leser, aber nur einen Schreiber. Das einfachste und sicherste Muster: Worker laden nur herunter und geben Daten zurück, und der Hauptthread erledigt das gesamte Schreiben in die Datenbank. Keine Sperren im Code, kein „database is locked“. Aktivieren Sie zusätzlich den WAL-Modus, damit das Lesen des Zustands aus einem anderen Prozess das Schreiben nicht behindert.

Begrenzung der Gleichzeitigkeit

Die Anzahl der Worker sollte durch zwei Dinge begrenzt werden. Erstens durch die Möglichkeiten des Proxys: Wenn Sie in Ihrem Proxeon-Kundenkonto mehrere Kanäle haben, ist es sinnvoll, ein bis zwei Worker pro Kanal zu halten, damit die IP-Rotation auf einem Kanal nicht die Verbindungen aller Threads gleichzeitig unterbricht. Zweitens durch die Höflichkeit gegenüber der Quelle: Selbst ohne formale Limits erzeugen zehn parallele Threads auf einer kleinen API eine Last, die zu Ablehnungen führt. Beginnen Sie mit drei bis vier Workern und erhöhen Sie sie, während Sie den Fehleranteil beobachten.

Code für einen parallelen Loader

import json
import sqlite3
from concurrent.futures import ThreadPoolExecutor, as_completed


def init_tasks(con):
con.execute('PRAGMA journal_mode=WAL')
con.executescript('''
CREATE TABLE IF NOT EXISTS tasks(
task_id TEXT PRIMARY KEY,
params TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
''')
with con:
con.execute('UPDATE tasks SET status=? WHERE status=?', ('pending', 'running'))


def enqueue(con, tasks):
with con:
con.executemany('INSERT OR IGNORE INTO tasks(task_id, params) VALUES (?, ?)',
[(t['task_id'], json.dumps(t['params'])) for t in tasks])


def claim(con, limit):
rows = con.execute('SELECT task_id, params FROM tasks WHERE status=? LIMIT ?',
 ('pending', limit)).fetchall()
with con:
con.executemany('UPDATE tasks SET status=? WHERE task_id=?',
[('running', r[0]) for r in rows])
return [(r[0], json.loads(r[1])) for r in rows]


def run_parallel(con, fetch_fn, write_fn, workers=4, max_attempts=5):
init_tasks(con)
with ThreadPoolExecutor(max_workers=workers) as pool:
while True:
batch = claim(con, workers * 2)
if not batch:
break
futures = {pool.submit(fetch_fn, params): task_id for task_id, params in batch}
for fut in as_completed(futures):
task_id = futures[fut]
try:
rows = fut.result()
except Exception as e:
with con:
con.execute(
'UPDATE tasks SET attempts=attempts+1, last_error=?, '
'status=CASE WHEN attempts+1 >= ? THEN ? ELSE ? END WHERE task_id=?',
(str(e)[:500], max_attempts, 'failed', 'pending', task_id))
continue
with con: # данные и статус задания в одной транзакции
write_fn(con, rows)
con.execute('UPDATE tasks SET status=? WHERE task_id=?', ('done', task_id))
failed = con.execute('SELECT count(*) FROM tasks WHERE status=?', ('failed',)).fetchone()[0]
print('очередь пуста, заданий с ошибкой:', failed)

Die Funktion fetch_fn führt die Anfrage über den Proxy aus und gibt eine Liste von Datensätzen zurück. Sie arbeitet in einem Thread und berührt die Datenbank nicht. Die Funktion write_fn wird im Hauptthread innerhalb einer Transaktion aufgerufen und führt das idempotente Einfügen aus Schritt 3 aus. Eine fehlgeschlagene Aufgabe wird automatisch wieder auf pending gesetzt und im nächsten claim-Zyklus erneut übernommen; nach Erschöpfung der Versuche erhält sie den Status failed, und Sie kümmern sich manuell darum.

Achtung: Übergeben Sie das sqlite3-Verbindungsobjekt nicht an Worker. Die Verbindung ist an den Thread gebunden, in dem sie erstellt wurde, und der Versuch, sie aus einem anderen Thread zu verwenden, führt zu einem Fehler oder, schlimmer, zu einer stillen Datenbeschädigung. Entweder erhält jeder Thread seine eigene Verbindung oder, wie im obigen Beispiel, keine.

Tipp: Verwenden Sie für jeden Worker ein separates requests.Session-Objekt mit eigener Proxy-Adresse. Wenn Sie in Proxeon mehrere Kanäle haben, verteilen Sie sie reihum auf die Worker: Worker 0 nimmt Kanal 0, Worker 1 Kanal 1 und so weiter. So betrifft ein Verbindungsabbruch auf einem Kanal nur einen Thread.

Erwartetes Ergebnis

Prüfung: Stellen Sie 100 Aufgaben in die Warteschlange, starten Sie vier Worker und beenden Sie den Prozess nach einer halben Minute. Führen Sie SELECT status, count(*) FROM tasks GROUP BY status; aus. Sie sehen einige done, einige running und den Rest pending. Starten Sie erneut: running sollten beim Start verschwinden, und am Ende sollten alle Aufgaben done sein, außer denen, die ehrlich fehlgeschlagen sind und in failed mit dem Fehlertext in last_error liegen.

Schritt 6: Wiederaufnahme nach einer langen Pause

Ziel dieser Phase: Einen Export, der für mehrere Stunden oder Tage unterbrochen wurde, korrekt fortsetzen, ohne auf veralteten Zustand zu stoßen.

Was veraltet

Eine Wiederaufnahme nach zehn Sekunden und nach einer Woche sind unterschiedliche Aufgaben. Nach einer langen Pause ist ein Teil des gespeicherten Zustands nicht mehr gültig.

  1. Sitzungen und Cookies. Serverseitige Sitzungen leben normalerweise von einigen Stunden bis zu einem Tag. Gespeicherte Cookies führen danach zu Antworten 401 oder einer Weiterleitung zum Anmeldeformular. Lösung: Beim Start eine vollständige erneute Autorisierung durchführen, nicht die Cookies aus der Datei wiederherstellen.
  2. Zugriffstoken. OAuth-Token leben etwa eine Stunde, manchmal weniger. Wenn Sie ein Refresh-Token haben, aktualisieren Sie das Access-Token vor dem Start und planmäßig während der Arbeit, ohne auf eine Ablehnung zu warten.
  3. Paginierungscursor. Viele APIs begrenzen die Lebensdauer eines Cursors auf Minuten oder Stunden. Ein veralteter Cursor führt zu einem Fehler 400 mit der Meldung über einen ungültigen Cursor. Deshalb haben wir in Schritt 2 einen Backup-Anker gespeichert: die ID des letzten Datensatzes. Wenn die Quelle einen Filter nach ID oder Änderungsdatum unterstützt, bauen Sie eine neue Anfrage von diesem Anker aus. Wenn nicht, muss der Shard von vorne beginnen, und die Deduplizierung aus Schritt 4 filtert das bereits Erhaltene heraus.
  4. Dateiinhalt auf dem Server. Für das Nachladen aus Schritt 1 ist entscheidend, dass sich die Datei nicht geändert hat. Vergleichen Sie das aktuelle ETag mit dem gespeicherten vor jedem Versuch; bei Nichtübereinstimmung beginnen Sie die Datei von vorne.
  5. Proxy-Einstellungen. Innerhalb einer Woche können sich in Ihrem Proxeon-Kundenkonto Port, Passwort oder die Laufzeit des Kanals geändert haben. Prüfen Sie den Proxy mit einer Testanfrage, bevor Sie die Warteschlange abarbeiten.
  6. Der Datensatz selbst. Wenn der Export eine Woche dauert und die Quelle in dieser Zeit Datensätze hinzugefügt und gelöscht hat, ist Ihr Ergebnis eine Mischung aus Zuständen zu verschiedenen Zeitpunkten. Für viele Aufgaben ist das akzeptabel. Wenn nicht, speichern Sie die Startzeit im Checkpoint und führen Sie nach Abschluss einen separaten inkrementellen Durchlauf über die seit dieser Zeit geänderten Datensätze durch.

Vorflugprüfung

Fassen Sie alle Prüfungen in einer Funktion zusammen, die beim Start vor jeder echten Arbeit ausgeführt wird. Sie bringt entweder den Zustand in Ordnung oder stoppt den Loader mit einer verständlichen Meldung.

def preflight(session, state, probe_url):
# 1. прокси жив и авторизован
r = session.head(probe_url, timeout=20)
if r.status_code == 407:
raise SystemExit('прокси отверг логин или пароль, проверьте данные в кабинете Proxeon')
# 2. токен доступа свежий
refresh_access_token(session)
# 3. курсор ещё действителен
if state['cursor']:
test = session.get(API_BASE + '/orders', params={'cursor': state['cursor'], 'limit': 1}, timeout=30)
if test.status_code == 400 and 'cursor' in test.text.lower():
print('курсор протух, переключаемся на якорь по last_id =', state['last_id'])
state['cursor'] = None
state['resume_after_id'] = state['last_id']
# 4. напоминание о возрасте выгрузки
print('выгрузка стартовала', state['started_at'], 'страниц записано', state['pages'])
return state

Die Funktion refresh_access_token hängt von Ihrer Quelle ab: In der Regel ist es eine POST-Anfrage mit dem Refresh-Token, nach der Sie den Authorization-Header in der Sitzung aktualisieren. Das Feld resume_after_id wird dann in der Funktion zum Abrufen der Seite als Filter „ID größer als angegeben“ verwendet.

Tipp: Speichern Sie das Refresh-Token und das Proxy-Passwort nicht im Checkpoint, sondern in Umgebungsvariablen oder einer separaten Datei mit eingeschränkten Zugriffsrechten. Den Checkpoint werden Sie kopieren, an Kollegen weiterleiten und Fehlerberichten beifügen; Geheimnisse haben dort nichts verloren.

Erwartetes Ergebnis

Prüfung: Beschädigen Sie den Cursor in der Tabelle checkpoint manuell mit UPDATE checkpoint SET cursor='broken'; und starten Sie den Loader. In der Konsole muss eine Zeile über das Umschalten auf den Anker per last_id erscheinen, und der Export muss ohne Absturz fortgesetzt werden. Die Anzahl der Datensätze nach Abschluss muss mit einem Kontrolllauf ohne Cursor-Beschädigung übereinstimmen.

Schritt 7: Fertiges Gerüst für einen robusten Loader in Python

Ziel dieser Phase: Alles aus den vorherigen Schritten in einer Datei zusammenführen, die gestartet, unterbrochen und erneut gestartet werden kann und ein vollständiges Ergebnis ohne Duplikate liefert.

Struktur des Gerüsts

  • Konfiguration aus Umgebungsvariablen: Proxeon-Proxy-Adresse, API-Adresse, Token, Jobname.
  • Klasse Store: SQLite mit Tabellen für Datensätze und Checkpoint, eine Transaktion pro Seite.
  • Funktion für den Datensatzschlüssel zur Idempotenz.
  • Funktion zum Abrufen einer Seite über den Proxy.
  • Hauptschleife mit Wiederaufnahme per Checkpoint und Neuerstellung der Sitzung nach einem Abbruch.

Vollständiger Code

# resumable_loader.py
import hashlib
import json
import os
import sqlite3
import time
from datetime import datetime, timezone

import requests

PROXY_URL = os.environ['PROXY_URL']# http://USER:PASS@HOST:PORT из кабинета Proxeon
API_BASE = os.environ.get('API_BASE', 'https://api.example.com')
API_TOKEN = os.environ.get('API_TOKEN', '')
DB_PATH = os.environ.get('DB_PATH', 'export.sqlite')
JOB = os.environ.get('JOB', 'orders-2026')
PAGE_SIZE = 500
MAX_ATTEMPTS = 8
FIELDS = ('cursor', 'last_id', 'pages', 'rows', 'started_at')


def now():
return datetime.now(timezone.utc).isoformat()


def record_key(rec):
canon = json.dumps({'id': rec['id']}, sort_keys=True, separators=(',', ':'))
return hashlib.blake2b(canon.encode('utf-8'), digest_size=16).digest()


class Store:
def __init__(self, path):
self.con = sqlite3.connect(path)
self.con.execute('PRAGMA journal_mode=WAL')
self.con.executescript('''
CREATE TABLE IF NOT EXISTS records(
key BLOB PRIMARY KEY,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
) WITHOUT ROWID;
CREATE TABLE IF NOT EXISTS checkpoint(
job TEXT PRIMARY KEY,
cursor TEXT,
last_id TEXT,
pages INTEGER NOT NULL,
rows INTEGER NOT NULL,
started_at TEXT,
updated_at TEXT
);
''')

def load(self, job):
row = self.con.execute(
'SELECT cursor, last_id, pages, rows, started_at FROM checkpoint WHERE job=?',
(job,)).fetchone()
if row is None:
return {'cursor': None, 'last_id': None, 'pages': 0, 'rows': 0, 'started_at': now()}
return dict(zip(FIELDS, row))

def commit_page(self, job, rows, state):
ts = now()
with self.con:
self.con.executemany(
'INSERT OR IGNORE INTO records(key, payload, fetched_at) VALUES (?, ?, ?)',
[(record_key(r), json.dumps(r, ensure_ascii=False), ts) for r in rows])
self.con.execute(
'INSERT INTO checkpoint(job, cursor, last_id, pages, rows, started_at, updated_at) '
'VALUES (?, ?, ?, ?, ?, ?, ?) '
'ON CONFLICT(job) DO UPDATE SET cursor=excluded.cursor, last_id=excluded.last_id, '
'pages=excluded.pages, rows=excluded.rows, updated_at=excluded.updated_at',
(job, state['cursor'], state['last_id'], state['pages'], state['rows'],
 state['started_at'], ts))

def unique_count(self):
return self.con.execute('SELECT count(*) FROM records').fetchone()[0]


def make_session():
s = requests.Session()
s.proxies = {'http': PROXY_URL, 'https': PROXY_URL}
s.headers['User-Agent'] = 'resumable-loader/1.0'
if API_TOKEN:
s.headers['Authorization'] = 'Bearer ' + API_TOKEN
return s


def fetch_page(session, cursor):
params = {'limit': PAGE_SIZE}
if cursor:
params['cursor'] = cursor
r = session.get(API_BASE + '/orders', params=params, timeout=(15, 90))
r.raise_for_status()
body = r.json()
return body['items'], body.get('next_cursor')


def run():
store = Store(DB_PATH)
state = store.load(JOB)
print('старт: страниц %d, строк %d, уникальных в базе %d'
 % (state['pages'], state['rows'], store.unique_count()))
session = make_session()
attempts = 0
while True:
try:
items, next_cursor = fetch_page(session, state['cursor'])
attempts = 0
except (requests.ConnectionError, requests.Timeout, requests.HTTPError) as e:
attempts += 1
if attempts > MAX_ATTEMPTS:
print('попытки исчерпаны, состояние сохранено, запустите снова позже')
raise
print('обрыв (%s), попытка %d из %d' % (type(e).__name__, attempts, MAX_ATTEMPTS))
time.sleep(min(60, 2 ** attempts)) # выбор пауз описан в статье про 429 и ретраи
session = make_session() # новая сессия: соединение через прокси пересоздаётся
continue
if not items:
break
state['pages'] += 1
state['rows'] += len(items)
state['last_id'] = str(items[-1]['id'])
state['cursor'] = next_cursor
store.commit_page(JOB, items, state)
if state['pages'] % 20 == 0:
print('страниц %d, строк %d' % (state['pages'], state['rows']))
if next_cursor is None:
break
print('готово: страниц %d, строк получено %d, уникальных в базе %d'
 % (state['pages'], state['rows'], store.unique_count()))


if __name__ == '__main__':
run()

Wie Sie es an Ihre Quelle anpassen

  1. Ersetzen Sie den Pfad /orders und die Feldnamen items, next_cursor, id durch die Ihrer API. Das sind drei Stellen in den Funktionen fetch_page und record_key.
  2. Wenn die Quelle eine Offset-Paginierung hat, ersetzen Sie den Parameter cursor durch page und berechnen Sie den nächsten Wert als state['pages'] + 1. Speichern Sie im Checkpoint statt des Cursors die Seitennummer.
  3. Bei einer Schlüssel-Paginierung übergeben Sie einen Parameter wie after_id aus state['last_id'] und entfernen Sie die Arbeit mit dem Cursor.
  4. Wenn Sie Parallelität benötigen, verschieben Sie fetch_page in fetch_fn aus Schritt 5 und verwenden Sie commit_page als write_fn. Teilen Sie den Export in Shards auf und füllen Sie die Aufgabenwarteschlange.
  5. Fügen Sie die Funktion preflight aus Schritt 6 vor der Hauptschleife hinzu.

Ergebnisprüfung: Checkliste

Bevor Sie den Loader mit einem echten mehrstündigen Volumen starten, gehen Sie diese Liste durch. Jeder Punkt dauert ein paar Minuten, und zusammen garantieren sie, dass der nächtliche Export keine Daten verliert.

  1. Unterbrechungstest. Starten Sie den Loader und drücken Sie nach 30 Sekunden Ctrl+C. Starten Sie ihn erneut. Die erste Ausgabezeile muss eine von null verschiedene Seitenzahl zeigen, nicht „страниц 0“.
  2. Duplikattest. Reduzieren Sie PAGE_SIZE auf 10 und unterbrechen Sie den Loader fünfmal hintereinander zu zufälligen Zeitpunkten. Vergleichen Sie nach Abschluss die Anzahl der eindeutigen Datensätze mit der Anzahl der erhaltenen Zeilen: Die eindeutigen sollten kleiner oder gleich sein, und bei sauberer Cursor-Paginierung ohne Einfügungen in die Quelle nahezu gleich.
  3. Konsistenztest. Führen Sie nach jeder Unterbrechung zwei Abfragen aus: SELECT rows FROM checkpoint; und SELECT count(*) FROM records;. Die Differenz zwischen ihnen darf eine Seitengröße nicht überschreiten. Wenn sie größer ist, werden Checkpoint und Daten nicht in derselben Transaktion geschrieben.
  4. Proxytest. Deaktivieren Sie vorübergehend die Variable PROXY_URL oder geben Sie ein falsches Passwort an. Der Loader muss bei der ersten Anfrage mit einem verständlichen Fehler abstürzen und nicht hängen bleiben oder direkt auf die Quelle zugreifen.
  5. Test mit veraltetem Cursor. Beschädigen Sie den Cursor in der Datenbank wie in Schritt 6 beschrieben und stellen Sie sicher, dass das Umschalten auf den Anker funktioniert.
  6. Festplattentest. Prüfen Sie die Größe der Datei export.sqlite nach tausend Seiten und multiplizieren Sie sie mit der erwarteten Seitenzahl. Stellen Sie sicher, dass auf der Festplatte genügend Platz mit 20 Prozent Reserve vorhanden ist.

Prüfung: Als erfolgreiche Ausführung gilt eine Situation, in der nach drei absichtlichen Unterbrechungen und drei Neustarts die endgültige Anzahl eindeutiger Datensätze mit der Anzahl übereinstimmt, die bei einem einzigen kontinuierlichen Durchlauf an derselben Quelle erhalten wurde, und in der Konsole nie die Zeile über einen Start von null erschienen ist.

Zusätzliche Möglichkeiten und Optimierung

  • Fortschritt und Zeitschätzung. Wenn die Gesamtzahl der Datensätze bekannt ist, geben Sie alle zwanzig Seiten den Prozentsatz und eine Schätzung der verbleibenden Zeit aus. Das ist sowohl für Sie nützlich als auch, um ein Hängen von langsamer Arbeit zu unterscheiden.
  • Kompression der Nutzdaten. Bei Dutzenden Millionen Datensätzen belegt JSON im Textformat viel Platz. Komprimieren Sie das Feld payload mit zlib.compress vor dem Schreiben und speichern Sie es als BLOB. Die Einsparung liegt normalerweise zwischen dem Drei- und Achtfachen.
  • Umzug auf eine Serverdatenbank. Das Schema mit dem Checkpoint in einer Transaktion mit den Daten lässt sich fast unverändert auf PostgreSQL übertragen. Das Konstrukt ON CONFLICT wird dort unterstützt, und die Beschränkung „ein Schreiber“ entfällt.
  • Inkrementelle Exporte. Speichern Sie die Startzeit jedes Jobs und starten Sie nach dem vollständigen Export einen separaten Job mit dem Filter „geändert nach“. So halten Sie eine aktuelle Kopie ohne vollständiges Neuladen.
  • Separater Prozess pro Shard. Statt Threads können Sie mehrere Instanzen des Skripts mit unterschiedlichen Werten für JOB und unterschiedlichen DB_PATH-Dateien starten und die Ergebnisse am Ende zusammenführen. Das ist einfacher zu debuggen und beseitigt die Frage der konkurrierenden Schreibzugriffe vollständig.
  • Metriken zu Abbrüchen. Protokollieren Sie jeden Abbruch mit Ausnahmetyp und Zeit. Nach einem Tag werden Sie sehen, dass sich Abbrüche um die Intervalle der IP-Rotation auf dem Proxeon-Kanal gruppieren, und können das Rotationsintervall an die Dauer Ihrer Anfragen anpassen.

Typische Fehler und Lösungen

Unten sind Situationen zusammengestellt, mit denen fast jeder bei den ersten Starts konfrontiert wird. Format: Problem, Ursache, Lösung.

  1. Problem: Nach einem Neustart beginnt der Export jedes Mal von vorne. Ursache: Der Checkpoint wird im Speicher oder in einer Datei geschrieben, die einen Absturz nicht überlebt, oder der Loader liest ihn beim Start nicht. Lösung: Stellen Sie sicher, dass die erste Aktion in der Funktion run store.load ist und der Zustand nach jeder Seite innerhalb einer Transaktion aktualisiert wird.
  2. Problem: In der Datenbank sind anderthalbmal mehr Datensätze als in der Quelle. Ursache: Der Deduplizierungsschlüssel ist instabil: Er enthält die Anfragezeit, die Seitennummer oder ein Feld mit zufälliger Reihenfolge. Lösung: Berechnen Sie den Schlüssel nur anhand der Quell-ID oder anhand einer expliziten Liste stabiler Felder über den Parameter fields.
  3. Problem: In der Datenbank sind weniger Datensätze als in der Quelle, obwohl der Export ohne Fehler abgeschlossen wurde. Ursache: Der Checkpoint wurde vor dem Schreiben der Daten aktualisiert, und nach einem Absturz wurde die Seite übersprungen. Oder die Offset-Paginierung hat bei Löschungen in der Quelle die Seiten nach hinten verschoben. Lösung: Ändern Sie die Reihenfolge auf „Daten, dann Checkpoint“ in einer Transaktion; für Quellen mit Löschungen wechseln Sie zur Schlüssel-Paginierung.
  4. Problem: Das Nachladen der Datei ergibt bei übereinstimmender Größe ein defektes Archiv. Ursache: Der Server antwortete einmal mit 200 statt 206, die Datei wurde teilweise überschrieben und dann angehängt. Lösung: Speichern Sie ETag neben der Datei, löschen Sie die Datei bei Änderung; prüfen Sie den Content-Range-Header auf Übereinstimmung mit dem angeforderten Offset.
  5. Problem: Fehler „database is locked“ bei paralleler Arbeit. Ursache: Mehrere Threads schreiben gleichzeitig in SQLite, oder die Verbindung wurde zwischen Threads weitergegeben. Lösung: Eine Schreibstelle im Hauptthread, Worker laden nur herunter; WAL-Modus; die Verbindung wird in dem Thread erstellt, in dem sie verwendet wird.
  6. Problem: Nach einer Stunde Arbeit beginnen alle Anfragen, 401 zurückzugeben. Ursache: Das Zugriffstoken ist abgelaufen. Lösung: Aktualisieren Sie das Token planmäßig vor Ablauf, und bei Erhalt eines 401 rufen Sie die Aktualisierung auf und wiederholen Sie die Anfrage einmal, ohne dies als Abbruch zu zählen.
  7. Problem: Der Arbeitsspeicher wächst auf mehrere Gigabyte. Ursache: Die Menge der gesehenen Schlüssel oder die Liste aller Datensätze wird im Prozessspeicher gehalten. Lösung: Deduplizierung über den Primärschlüssel in der Datenbank oder die Tabelle seen; Schreiben der Daten seitenweise, ohne Ansammlung.
  8. Problem: Abbrüche treten streng alle paar Minuten auf. Ursache: Sie fallen mit dem Intervall der IP-Rotation auf dem Proxy-Kanal zusammen. Lösung: Das ist eine normale Situation, der Loader muss sie überstehen. Wenn die Anfragen lang sind, wählen Sie das Rotationsintervall in Ihrem Proxeon-Kundenkonto so, dass es deutlich größer als die typische Dauer einer Anfrage ist, oder verwenden Sie die Rotation auf Anfrage zwischen den Seiten.

FAQ: Häufige Fragen zum robusten Export

Muss ich SQLite verwenden, wenn das Ergebnis in CSV vorliegen soll?

Nein, aber es ist praktisch. SQLite dient hier als zuverlässiger Speicher für den Zustand und die Menge der gesehenen Schlüssel. Die endgültige CSV exportieren Sie mit einem Befehl aus der Tabelle records nach Abschluss. Wenn Sie ganz ohne Datenbank auskommen möchten, verwenden Sie die Dateivariante aus Schritt 3 mit einer Datei pro Seite und einem JSON-Checkpoint mit atomarem Ersetzen.

Wie oft sollte der Checkpoint gespeichert werden: nach jeder Seite oder seltener?

Nach jeder Seite. Eine SQLite-Transaktion mit einigen hundert Zeilen dauert Millisekunden, das ist vernachlässigbar im Vergleich zur Netzwerkanfrage über den Proxy. Die Einsparung bei seltenen Checkpoints ist das Risiko nicht wert, Dutzende Seiten zu verlieren.

Was tun, wenn die API weder Cursor noch IDs liefert, nur Seitennummern?

Arbeiten Sie mit der Seitennummer, speichern Sie sie im Checkpoint und fügen Sie unbedingt eine Deduplizierung nach dem Hash des Datensatzinhalts hinzu. Akzeptieren Sie, dass bei aktiven Änderungen in der Quelle ein Teil der Datensätze aufgrund der Seitenverschiebung ausgelassen werden kann. Für kritische Daten machen Sie einen zweiten Durchlauf in umgekehrter Seitenreihenfolge: Das beim ersten Durchlauf Ausgelassene wird mit hoher Wahrscheinlichkeit im zweiten erfasst.

Kann ich die Datei mit mehreren Threads in verschiedenen Bereichen nachladen?

Ja, wenn der Server Range unterstützt. Teilen Sie die Datei in Stücke von 50–100 MB auf, jedes Stück ist eine Aufgabe aus der Warteschlange von Schritt 5 mit eigener temporärer Datei, und nach Abschluss aller Stücke fügen Sie sie in der richtigen Reihenfolge zusammen. Prüfen Sie jedes Stück anhand der Größe und die gesamte Datei anhand der Prüfsumme, falls vorhanden.

Wie viele Worker sollte ich bei der Arbeit über einen Proxy einsetzen?

Beginnen Sie mit drei bis vier pro Proxeon-Kanal und beobachten Sie den Fehleranteil in der Tabelle tasks. Wenn die Fehler unter einem Prozent liegen, fügen Sie zwei weitere hinzu. Wenn die Fehler wachsen, reduzieren Sie. Mehr als zehn Threads pro Kanal bringen selten einen Gewinn: Sie stoßen entweder an die Bandbreite des Kanals oder an die Geduld der Quelle.

Muss ich die Cookies der Sitzung zwischen den Starts speichern?

Normalerweise nicht. Eine erneute Autorisierung beim Start dauert Sekunden und ist zuverlässiger als die Wiederherstellung von Cookies mit unbekannter Lebensdauer. Ausnahme: Die Quelle begrenzt die Anzahl der Anmeldungen pro Tag. Dann speichern Sie die Cookies, aber bei der ersten 401 oder Weiterleitung zur Anmeldung verwerfen Sie sie und melden Sie sich erneut an.

Wie erkenne ich, dass der Export vollständig abgeschlossen und nicht stillschweigend abgebrochen wurde?

Für Dateien: Die Größe entspricht Content-Length und die Prüfsumme stimmt überein. Für APIs: Es wurde eine Seite ohne next_cursor oder eine leere Seite empfangen, und die Anzahl der Zeilen entspricht der Gesamtzahl, falls die Quelle sie meldet. Schreiben Sie ein explizites Abschlussflag in den Checkpoint, damit ein erneuter Start keinen neuen Durchlauf beginnt.

Was tun mit Aufgaben im Status failed?

Sehen Sie sich das Feld last_error an. Wenn es Netzwerkfehler sind, setzen Sie die Aufgaben einfach mit UPDATE zurück auf pending und starten Sie den Loader erneut. Wenn es Parsing-Fehler sind, gibt es in der Quelle Datensätze mit ungewöhnlicher Form: Korrigieren Sie den Code und starten Sie neu. Löschen Sie failed niemals stillschweigend, das ist der einzige Beweis dafür, was im Export fehlt.

Kann ich diesen Ansatz statt mit requests mit einer asynchronen Bibliothek verwenden?

Ja, die Prinzipien sind dieselben: atomare Arbeitseinheit, Checkpoint zusammen mit den Daten, idempotentes Schreiben, Warteschlange auf der Festplatte. Es ändert sich nur der Transport. Ein einziger Hinweis: Lassen Sie das Schreiben in SQLite synchron und sequenziell, und halten Sie die Parallelität auf der Ebene der Netzwerkanfragen.

Fazit

Sie haben den Weg von einer unvollständigen Datei und nervösem Neustart zu einem Loader zurückgelegt, dem ein Abbruch gleichgültig ist. Fassen wir zusammen, was genau gemacht wurde.

  • Wir haben uns mit dem Nachladen per HTTP befasst: Prüfung von Accept-Ranges per HEAD, Range-Header mit der aktuellen Dateigröße, Unterscheidung der Codes 206 und 200, Schutz vor Dateiaustausch über ETag und If-Range.
  • Wir haben Checkpoints für seitenbasierte Exporte erstellt: Cursor, Seitennummer oder ID des letzten Datensatzes plus Hilfszähler, alles in einer Transaktion mit den Daten.
  • Wir haben das Schreiben idempotent gemacht über Primärschlüssel und das Konstrukt INSERT OR IGNORE oder ON CONFLICT DO UPDATE, und für Dateien über deterministische Namen und atomares Ersetzen.
  • Wir haben die Deduplizierung auf der Festplatte organisiert, damit Millionen von Schlüsseln nicht im Arbeitsspeicher leben und Neustarts überleben.
  • Wir haben Parallelität mit einer Aufgabenwarteschlange in SQLite hinzugefügt, automatischer Rückgabe fehlgeschlagener Aufgaben und einer einzigen Schreibstelle.
  • Wir haben die Wiederaufnahme nach einer langen Pause vorgesehen: Aktualisierung von Token, erneute Autorisierung, Umschalten von einem veralteten Cursor auf einen Anker per ID, Prüfung des Proxeon-Proxys vor dem Start.
  • Wir haben alles in einem funktionierenden Gerüst zusammengeführt, das durch Ersetzen von drei bis vier Zeilen an eine konkrete Quelle angepasst wird.

Wie es weitergeht

Nehmen Sie das Gerüst aus Schritt 7 und starten Sie es mit einem kleinen echten Volumen, sagen wir zehntausend Datensätzen. Arbeiten Sie die Checkliste aus dem Abschnitt zur Prüfung durch. Erst danach starten Sie den vollständigen Export über Nacht. Am Morgen sehen Sie entweder die Zeile „готово“ mit übereinstimmenden Zählern oder die Zeile über erschöpfte Versuche mit gespeichertem Zustand, und dann starten Sie das Skript einfach erneut.

Wohin Sie sich weiterentwickeln können

Die nächste Stufe sind inkrementelle Exporte nach Änderungszeit statt vollständiger Durchläufe, die Verlagerung des Zustands in eine Serverdatenbank für mehrere Maschinen sowie eine sinnvolle Wiederholungsstrategie unter Berücksichtigung der Antwortcodes, der ein separater Artikel über 429 und Retries gewidmet ist. Die Kombination aus gut durchdachten Retries aus jenem Artikel und dem robusten Zustand aus diesem ergibt einen Loader, den Sie eine Woche lang arbeiten lassen können, ohne ein Terminal zu öffnen.

Und zum Schluss: Ein Abbruch eines langen Exports ist kein Notfall, sondern eine Arbeitssituation, die Sie jetzt zu behandeln wissen. Viel Erfolg beim Exportieren.