Introduction : pourquoi une longue extraction se coupe presque toujours, et c'est normal

Si vous avez déjà extrait via proxy plusieurs millions d'enregistrements ou un fichier de plusieurs dizaines de gigaoctets, vous connaissez cette sensation. Le script tournait depuis six heures, affichait 83 pour cent, puis il a planté avec une erreur de connexion. Et tout ce que vous avez, c'est un fichier incomplet et la certitude qu'il va falloir tout reprendre depuis le début.

La première chose à accepter : une longue extraction se coupe toujours. Pas « parfois », pas « quand le réseau est mauvais », mais toujours, dès qu'elle dure assez longtemps. Les causes sont multiples, et la plupart échappent à votre contrôle :

  • Le proxy change d'adresse IP externe. Chez les proxies mobiles Proxeon, c'est un comportement normal : rotation par minuteur ou à la demande. Au moment du changement d'IP, la connexion TCP ouverte est coupée, et le serveur source voit déjà un autre client.
  • Le serveur source ferme la connexion selon son propre timeout, redémarre, déploie une mise à jour ou renvoie simplement une erreur 5xx.
  • Votre propre processus redémarre : mise à jour système, disque plein, bug sur un enregistrement atypique, Ctrl+C accidentel.
  • Le token d'autorisation expire, la session se termine, le curseur de pagination devient obsolète.
  • Le portable passe en veille, le Wi-Fi bascule sur un autre point d'accès, le fournisseur d'accès change de route.

Lutter contre chacune de ces causes séparément n'a pas de sens. La bonne approche est différente : concevoir l'extraction pour qu'une coupure à n'importe quel moment vous coûte non pas six heures, mais une page ou un morceau de fichier. C'est précisément l'objet de ce guide.

Ce que vous obtiendrez au final

Après avoir suivi cette procédure, vous aurez :

  • Une fonction opérationnelle de reprise de téléchargement depuis le milieu grâce au protocole HTTP via l'en-tête Range, avec vérification que le serveur le supporte.
  • Un schéma clair de checkpoints pour les extractions paginées : quoi sauvegarder exactement et où stocker l'état pour qu'il ne disparaisse pas avec le processus.
  • Une écriture idempotente des résultats, où le rechargement d'une même page ne crée pas de doublons.
  • Une déduplication qui ne consomme pas toute la mémoire vive sur des millions de lignes.
  • Un loader parallèle avec file de tâches, répétition des tâches échouées et limitation de la concurrence.
  • Un squelette prêt à l'emploi d'un loader robuste en Python, que vous adapterez à votre source en une heure.

À qui s'adresse ce guide

Aux développeurs et analystes qui savent déjà faire des requêtes HTTP en Python et qui ont perdu au moins une fois les résultats d'une longue extraction. Niveau intermédiaire : on explique les bases, mais on n'apprend pas à programmer depuis zéro. Les lecteurs avancés trouveront des sections sur le stockage des hachages hors mémoire et sur l'écriture parallèle sécurisée dans SQLite.

Ce qu'il faut savoir au préalable

  • Python au niveau des fonctions, boucles, dictionnaires et gestion des exceptions.
  • Les bases de HTTP : méthode, en-tête, code de réponse, corps.
  • Une notion générale de la façon de définir un proxy dans la bibliothèque requests.

Une précision importante : les codes de réponse et les stratégies de retry avec pause exponentielle ne sont pas traités ici. Un article dédié traite de l'erreur 429 et des retries. Dans ce guide, l'accent est mis sur autre chose : l'état de l'extraction et sa reprise. Les retries répondent à la question « quand relancer la requête », nous répondons à la question « à partir de quel endroit reprendre le travail une fois les retries épuisés et le processus mort ».

Combien de temps cela prendra

Lecture et exécution des exemples sur une source de test : deux à trois heures. Adaptation du squelette à votre API réelle ou serveur de fichiers : encore une à deux heures selon le degré de non-standard de la pagination. Au total, une journée de travail avec de la marge.

Préparation et notions de base

Outils et accès

  1. Installez Python version 3.11 ou plus récente. En 2026, les branches 3.12 et 3.13 sont d'actualité, tous les exemples sont vérifiés dessus. Pour vérifier la version : ouvrez le terminal et tapez python --version. Si vous voyez 3.11 ou plus, tout va bien.
  2. Installez la bibliothèque requests : pip install requests. La version 2.32 et plus suffit. Le module sqlite3 fait partie de la bibliothèque standard Python, rien à installer séparément.
  3. Obtenez l'accès au proxy. Ouvrez votre espace client Proxeon, choisissez le canal voulu et copiez quatre valeurs : hôte, port, login et mot de passe. En général, ils sont regroupés en une seule chaîne du type http://USER:PASS@HOST:PORT. Cette chaîne vous servira partout par la suite.
  4. Placez la chaîne du proxy dans une variable d'environnement, pas dans le code. Sous Linux et macOS : export PROXY_URL=http://USER:PASS@HOST:PORT. Sous Windows PowerShell : $env:PROXY_URL='http://USER:PASS@HOST:PORT'. Ainsi, vous ne committrez pas le mot de passe par accident.
  5. Vérifiez que le proxy répond. Exécutez dans le terminal : curl -x $PROXY_URL -I https://api.example.com/, en remplaçant par l'adresse de votre source. Vous devez voir une ligne avec un code de réponse, par exemple HTTP/2 200. Si vous voyez une erreur d'autorisation proxy 407, revérifiez le login et le mot de passe.

Configuration requise

N'importe quelle machine avec 2 Go de mémoire vive libre et un disque pouvant contenir le résultat de l'extraction plus 20 pour cent de marge pour les index SQLite. Si vous prévoyez de charger des millions d'enregistrements, le disque compte plus que la mémoire : toute l'approche repose sur le fait que l'état vit sur le disque, pas dans les variables du processus.

Sauvegardes

Le fichier d'état que vous allez créer ci-dessous (dans les exemples, c'est export.sqlite) deviendra l'artefact le plus précieux de tout le travail. Prenez l'habitude de le copier avant toute expérimentation de code : cp export.sqlite export.sqlite.bak. Une fois, cela vous économisera une journée d'extraction.

Attention : ne modifiez jamais le fichier SQLite à la main pendant que le loader tourne. Même une lecture depuis un programme tiers en mode incorrect peut bloquer l'écriture et faire planter le processus. Si vous voulez consulter l'état, arrêtez le loader ou utilisez le mode WAL, que nous décrirons dans la section sur le parallélisme.

Termes clés en langage simple

  • Checkpoint — une marque sauvegardée sur disque indiquant « jusqu'ici tout est extrait et enregistré ». Après une coupure, le loader lit le checkpoint et reprend à partir de là.
  • Curseur — une chaîne opaque que l'API renvoie avec la page et qu'il faut transmettre pour obtenir la page suivante. Vous ne composez ni ne décomposez le curseur vous-même.
  • Extraction paginée par offset — quand vous demandez « page 37 avec 500 enregistrements ». Schéma simple, mais quand de nouveaux enregistrements sont ajoutés à la source, les pages se décalent et des doublons ou des manques apparaissent.
  • Extraction paginée par clé (keyset) — quand vous demandez « tout ce qui a un identifiant supérieur à 184203, trié par identifiant ». Le schéma le plus robuste pour la reprise, si la source le supporte.
  • Idempotence — propriété d'une opération dont la répétition donne le même résultat qu'une exécution unique. Vous avez enregistré la page deux fois, mais en base elle n'apparaît qu'une seule fois.
  • Clé de déduplication — valeur selon laquelle deux enregistrements sont considérés comme identiques. Idéalement, c'est l'identifiant de la source ; s'il n'existe pas, la clé est calculée comme un hachage de champs stables.
  • En-tête Range — façon de demander au serveur HTTP de renvoyer non pas tout le fichier, mais une partie, par exemple les octets de 1048576 jusqu'à la fin.
  • Sémantique at-least-once — garantie que chaque enregistrement sera récupéré au moins une fois. Des répétitions sont possibles, mais pas de manques. C'est exactement ce que vous obtiendrez après ce guide, et la déduplication éliminera les doublons.

Principe directeur

Les sept étapes ci-dessous se résument à une idée : chaque unité de travail doit être atomique et répétable. Une unité de travail, c'est soit un morceau de fichier, soit une page API, soit une tâche de la file. Atomique, cela signifie que le résultat et la marque de son achèvement sont sauvegardés ensemble. Répétable, cela signifie que si l'unité est exécutée deux fois, rien ne casse. Quand ces deux propriétés sont respectées, une coupure à n'importe quel point devient inoffensive.

Étape 1 : Reprise d'un fichier via HTTP avec l'en-tête Range

Objectif de l'étape : apprendre à télécharger un gros fichier via proxy de sorte qu'après une coupure, le téléchargement reprenne à l'octet où il s'est arrêté, et non depuis zéro.

Comment ça marche

Le protocole HTTP permet au client de demander une partie d'une ressource. Pour cela, on ajoute à la requête l'en-tête Range: bytes=DEBUT-, où DEBUT est le décalage en octets. Si le serveur supporte les requêtes partielles, il répond avec le code 206 Partial Content et l'en-tête Content-Range: bytes DEBUT-FIN/TOTAL. S'il ne supporte pas, il ignore Range et renvoie tout le fichier avec un code 200. Votre travail est de distinguer ces deux cas.

Pour savoir à l'avance si le serveur supporte la reprise, une requête HEAD aide : elle renvoie uniquement les en-têtes sans le corps. Regardez Accept-Ranges: bytes. La valeur none ou l'absence de l'en-tête signifie généralement qu'il n'y a pas de reprise, bien que certains serveurs traitent quand même correctement Range, donc on fait la vérification finale via le code de réponse.

Procédure pas à pas

  1. Faites une requête HEAD via le proxy et sauvegardez les en-têtes Accept-Ranges, Content-Length, ETag et Last-Modified. L'ETag servira à vérifier si le fichier a changé sur le serveur entre vos tentatives.
  2. Regardez combien d'octets se trouvent déjà dans le fichier local. Si le fichier n'existe pas, considérez que c'est zéro.
  3. Si la taille locale est déjà égale à Content-Length, le fichier est complet, rien à faire.
  4. Si la taille locale est supérieure à zéro et que le serveur déclare supporter les plages, ajoutez l'en-tête Range avec la taille actuelle. Ajoutez aussi If-Range avec l'ETag sauvegardé : le serveur ne renverra alors une réponse partielle que si le fichier n'a pas changé, sinon il renverra tout le fichier avec un code 200.
  5. Envoyez impérativement Accept-Encoding: identity. Sans cela, le serveur peut appliquer une compression à la volée, et les décalages en octets ne correspondront plus à votre fichier.
  6. Ouvrez le fichier local en mode ajout ab si vous obtenez 206, ou en mode réécriture wb si vous obtenez 200.
  7. Lisez le corps en flux par morceaux de 256 Ko et écrivez sur le disque. Ne chargez pas toute la réponse en mémoire.
  8. À la fin, comparez la taille finale avec Content-Length. Si elles ne correspondent pas, la connexion s'est coupée silencieusement, et il faut un passage de plus.

Code fonctionnel

import os
import requests

PROXY_URL = os.environ['PROXY_URL'] # chaîne depuis l'espace client 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('fichier déjà complet :', have, 'octets')
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('serveur a renvoyé la mauvaise plage : ' + r.headers.get('Content-Range', ''))
mode = 'ab'
elif r.status_code == 200:
print('le serveur renvoie le fichier entier, on repart de zéro')
mode = 'wb'
have = 0
elif r.status_code == 416:
raise IOError('plage demandée hors du fichier, vérifiez la taille locale')
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('coupure : reçu %d sur %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('tentative %d interrompue : %s' % (i + 1, type(e).__name__))
# pause avant la tentative suivante : stratégie décrite dans l'article sur 429 et retries
raise RuntimeError('impossible de terminer le téléchargement en %d tentatives' % max_rounds)


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

Notez la fonction download_until_done : elle ne contient pas de logique d'attente entre les tentatives. C'est intentionnel. Insérez-y votre stratégie de pauses issue de l'article sur les retries ; ici, seule la boucle « vérifier la taille, reprendre, vérifier à nouveau » compte.

Conseil : si le fichier est distribué sous forme d'archive, ne le décompressez pas à la volée pendant la reprise. Obtenez d'abord le fichier complet, vérifiez la taille et, si le serveur fournit une somme de contrôle, comparez-la. Ensuite seulement, décompressez. Un gzip partiellement téléchargé ressemble à un fichier corrompu, et vous perdrez du temps à chercher une erreur inexistante.

Résultat attendu

Vérification : lancez le script sur un fichier d'au moins 200 Mo, interrompez-le au bout de dix secondes avec Ctrl+C. Regardez la taille du fichier local, par exemple 41 943 040 octets. Relancez le script. La console ne doit pas afficher « on repart de zéro », et la taille du fichier doit continuer à croître, pas se réinitialiser. À la fin, la taille finale doit correspondre exactement au Content-Length de la requête HEAD.

Problèmes possibles

  • Le serveur renvoie toujours 200 au lieu de 206. Cela signifie que la reprise n'est pas supportée. La seule issue pour une telle source est de télécharger le fichier entier en une seule passe avec un timeout important, ou de chercher chez la source un format d'extraction par morceaux, par exemple un découpage par dates.
  • Pas d'en-tête Content-Length. Le serveur envoie le fichier en mode chunked sans annoncer la taille. On ne peut ni vérifier la complétude par la taille, ni reprendre : Range exige des décalages connus. Voyez avec la source ou utilisez une somme de contrôle, si elle est publiée.
  • Réponse 416 dès la première tentative. Le fichier local est plus grand que le fichier sur le serveur. Le fichier sur le serveur a changé et est devenu plus court. Supprimez le fichier local et repartez de zéro.
  • La taille correspond, mais le fichier est corrompu. Il y a probablement eu une coupure au milieu avec un code 200 et une réécriture depuis zéro, puis un ajout. Recréez le fichier. Pour éviter que cela se reproduise, conservez l'ETag dans un fichier séparé à côté et comparez-le avant chaque tentative.

Étape 2 : Checkpoints pour les extractions paginées

Objectif de l'étape : sauvegarder l'état de l'extraction pour qu'après n'importe quel plantage, le processus reprenne à la dernière page enregistrée avec succès.

Que sauvegarder

Le checkpoint minimal dépend du type de pagination de la source. Examinons trois cas.

  1. Pagination par curseur. L'API renvoie avec les données un champ du type next_cursor. Sauvegardez-le précisément. C'est le cas le plus simple : le curseur contient déjà tout ce dont le serveur a besoin pour continuer.
  2. Pagination par numéro de page ou offset. Sauvegardez le numéro de la dernière page entièrement enregistrée et la taille de page. Rappelez-vous que si des enregistrements sont ajoutés à la source pendant l'extraction, les offsets se décalent, donc la déduplication de l'étape 4 est obligatoire.
  3. Pagination par clé. Sauvegardez l'identifiant du dernier enregistrement enregistré. À la reprise, vous demandez tout ce qui est supérieur à cet identifiant. Le schéma ne craint ni les insertions, ni les longues pauses.

Quel que soit le type de pagination, il vaut la peine d'ajouter au checkpoint des champs auxiliaires : l'identifiant du dernier enregistrement (même pour un schéma à curseur, c'est un ancrage de secours si le curseur expire), des compteurs de pages et de lignes pour suivre la progression, l'heure de démarrage de l'extraction et l'heure de la dernière mise à jour.

Où stocker l'état

Il existe deux options opérationnelles, et toutes deux valent mieux que des variables en mémoire.

Option A : fichier JSON avec remplacement atomique

Convient si les résultats sont écrits dans des fichiers séparés, et non dans une base. Le piège principal : si vous écrivez l'état directement dans le fichier cible et que le processus plante au milieu de l'écriture, vous obtiendrez un JSON tronqué, illisible. La solution : écrire dans un fichier temporaire à côté et le renommer par-dessus le principal. L'opération de renommage dans un même système de fichiers est atomique.

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)

Option B : table dans SQLite à côté des données

Option préférable si vous stockez les enregistrements dans une base. Le checkpoint est mis à jour dans la même transaction que l'insertion des lignes de la page. Soit les données et la marque sont enregistrées, soit rien. Le décalage entre « données présentes, marque absente » n'existe pas par principe.

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: # une transaction par page
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))

Ordre des opérations

Retenez la règle : d'abord les données, ensuite le checkpoint, et idéalement dans une seule transaction. Si la transaction n'est pas possible (par exemple, les données sont écrites dans des fichiers), l'ordre est exactement celui-ci : on écrit le fichier de la page, on synchronise sur disque, puis on met à jour l'état. En cas de plantage entre ces deux actions, vous obtiendrez un rechargement d'une page, ce qui est sans danger grâce à l'étape 3. L'ordre inverse entraînerait un saut de page, et donc une perte de données.

Conseil : stockez dans le checkpoint non pas le curseur actuel, mais le curseur de la page suivante renvoyé par le serveur. Ainsi, à la reprise, vous demandez directement ce que vous n'avez pas encore, sans requête supplémentaire pour la page déjà obtenue.

Résultat attendu

Vérification : lancez l'extraction, attendez dix pages et terminez le processus de force. Ouvrez la base avec sqlite3 export.sqlite et exécutez SELECT pages, last_id FROM checkpoint;. Vous devez voir le nombre 10 et l'identifiant. Ensuite, exécutez SELECT count(*) FROM records; et vérifiez que le nombre d'enregistrements est égal à dix fois la taille de page. Relancez le loader : le premier message dans la console doit être quelque chose comme « démarrage : pages 10 ».

Problèmes possibles

  • Erreur « database is locked ». Un autre processus tient la connexion. Fermez toutes les fenêtres sqlite3 et autres outils ayant ouvert le fichier. Pour le travail multithread, activez le mode WAL, voir étape 5.
  • Le curseur est sauvegardé, mais pas les données. Vous avez mis à jour le checkpoint en dehors d'une transaction avec les données. Revenez au code ci-dessus et assurez-vous que les deux opérations sont dans un même bloc with con:.
  • Le JSON d'état est vide ou corrompu. Vous avez écrit directement dans le fichier sans fichier temporaire et remplacement. Utilisez intégralement la fonction save_state.

Étape 3 : Idempotence de l'écriture des résultats

Objectif de l'étape : faire en sorte que le retraitement de n'importe quelle page ne produise pas de doublons et ne casse pas les données.

Pourquoi la répétition est inévitable

Après l'étape 2, vous avez déjà vu le scénario où une page est enregistrée deux fois : le processus a planté après l'insertion des données, mais avant la mise à jour du checkpoint. En outre, les répétitions viennent de la pagination par offset lorsque la source change, des workers parallèles qui ont reçu la même tâche après un redémarrage, et simplement d'un redémarrage manuel « au cas où ». Lutter contre les répétitions côté requête est inutile. Il faut rendre l'écriture elle-même inoffensive en cas de répétition.

Choix de la clé de déduplication

  1. La source a un identifiant. Utilisez-le. C'est le champ id, uuid, order_number ou similaire, que la source garantit unique. Si vous extrayez depuis plusieurs sources vers une même table, faites une clé composite : nom de la source plus identifiant.
  2. Pas d'identifiant, mais un ensemble de champs qui ensemble définissent l'enregistrement. Par exemple, pour une ligne de tarif, c'est la référence plus l'entrepôt plus la date. Composez la clé à partir de ces champs après normalisation : mettez les chaînes dans la même casse, supprimez les espaces en début et fin, convertissez les dates dans un format unique.
  3. Rien de stable. Alors la clé devient le hachage de tout l'enregistrement après canonisation. Ce cas est détaillé à l'étape 4. Notez que si la source modifie un enregistrement (met à jour un prix), le hachage change, et vous obtiendrez les deux versions. Parfois c'est exactement ce qu'on veut, parfois non.

Attention : n'utilisez pas comme clé le numéro d'ordre de la ligne dans la réponse ni le numéro de page. Ces valeurs changent à la moindre modification de la source, et la déduplication se transforme en générateur de doublons.

Insertion idempotente en base

Dans SQLite et la plupart des bases relationnelles, il existe une construction qui soit ignore le conflit sur la clé primaire, soit met à jour la ligne existante. La première variante INSERT OR IGNORE, vous l'avez déjà vue à l'étape 2. Elle convient quand les enregistrements sont immuables. La seconde variante est nécessaire si la source peut mettre à jour les enregistrements et que vous voulez la version fraîche :

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])

Écriture idempotente dans des fichiers

Si le résultat doit être dans des fichiers plutôt qu'en base, appliquez le même principe : une page égale un fichier au nom déterministe. Le nom dépend des paramètres de la page, pas de l'heure ni d'un compteur. Avant le téléchargement, vérifiez si le fichier final existe ; si oui, vous sautez la page. Écrivez sous un nom temporaire et renommez à la fin, comme dans la fonction 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 # la page existe déjà, pas besoin de réécrire
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

Les fichiers avec l'extension .part restés après un plantage peuvent être supprimés sans crainte au démarrage : ils sont par définition incomplets.

Conseil : pour une extraction par offset, ne vous fiez pas uniquement à « le fichier existe, donc la page est prête ». Vérifiez en plus que le nombre de lignes dans le fichier est égal à la taille de page (sauf la dernière). Un fichier vide ou trop court portant un nom existant vaut mieux re-téléchargé.

Résultat attendu

Vérification : appelez trois fois de suite la fonction d'écriture d'une même page. Ensuite, exécutez SELECT count(*) FROM records;. Le nombre doit être égal à la taille d'une page, et non au triple. Pour la variante fichier, le répertoire doit contenir exactement un fichier de page et aucun fichier .part.

Étape 4 : Déduplication des résultats sans explosion mémoire

Objectif de l'étape : filtrer les enregistrements répétés à la volée sur des millions de lignes, sans garder toutes les clés en mémoire vive.

Hachage de l'enregistrement

Quand un enregistrement n'a pas d'identifiant, la clé devient le hachage de son contenu. Pour que des enregistrements identiques donnent un hachage identique, le contenu doit être canonisé : trier les clés du dictionnaire, supprimer les espaces superflus, fixer les séparateurs. Sinon, le même enregistrement arrivé avec un autre ordre de champs obtiendra un hachage différent.

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()

La fonction renvoie 16 octets. Cela suffit : la probabilité de collision accidentelle pour des centaines de millions d'enregistrements est négligeable. Le paramètre fields permet de calculer le hachage uniquement sur des champs stables, en excluant par exemple l'heure de dernière mise à jour, qui change à chaque requête.

Pourquoi un ensemble en mémoire ne tient pas sur des millions

La première idée qui vient à l'esprit : créer seen = set() et y déposer les clés. Calculons. Un objet bytes de 16 octets occupe en Python environ 49 octets plus les données elles-mêmes, soit environ 65 octets. Un emplacement dans l'ensemble, compte tenu du coefficient de remplissage, ajoute environ 30 octets. On obtient près de 95 octets par clé. Sur 10 millions d'enregistrements, cela fait environ 950 Mo, sur 50 millions près de 5 Go. Et surtout : après un redémarrage du processus, l'ensemble est vide, et toute la déduplication repart de zéro.

Trois façons de ne pas exploser la mémoire

  1. Stocker les clés dans la base elle-même. La voie la plus simple et la plus fiable. Si la clé est la clé primaire de la table records, la déduplication est déjà faite par la construction INSERT OR IGNORE de l'étape 3. L'index vit sur le disque, survit aux redémarrages, et SQLite met en cache les pages chaudes de l'index. Pour 10 millions de clés de 16 octets, l'index occupera environ 400-500 Mo sur disque, mais pas en mémoire.
  2. Table séparée des clés vues, sans rowid. Nécessaire si les données elles-mêmes ne sont pas écrites dans SQLite, mais par exemple dans des fichiers. SQLite sert alors uniquement d'ensemble compact sur disque.
  3. Ensemble compressé en mémoire comme pré-filtre. Variante avancée : tronquer le hachage à 8 octets et le stocker comme entier dans un tableau trié ou utiliser un filtre de Bloom. La mémoire est réduite plusieurs fois, mais une probabilité de faux positif apparaît. C'est pourquoi un tel pré-filtre ne sert qu'à écarter rapidement les enregistrements manifestement nouveaux, et la vérification finale se fait toujours via la base.

Implémentation de l'ensemble sur disque

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]
# déduplication à l'intérieur de la page
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])

Appelez filter_new avant d'enregistrer la page, et remember dans la même transaction que l'écriture des données et du checkpoint. Ainsi, après un plantage, l'ensemble des clés vues, les données et la marque de progression sont toujours cohérents entre eux.

Conseil : une requête avec IN sur 500 valeurs s'exécute via l'index en millisecondes. Ne vérifiez pas les clés une par une en boucle : c'est des dizaines de fois plus lent à cause des frais généraux de chaque requête.

Résultat attendu

Vérification : formez une page de test de 500 enregistrements, dont 100 sont répétés deux fois dans la page et 100 autres sont déjà dans la table seen. La fonction filter_new doit renvoyer exactement 300 enregistrements. Après un redémarrage du processus, les mêmes 500 enregistrements doivent donner zéro nouveau.

Problèmes possibles

  • Les doublons passent quand même. Vérifiez la canonisation : il y a probablement un champ avec l'heure de la requête ou un ordre aléatoire des éléments d'une liste. Excluez-le via le paramètre fields ou triez les listes imbriquées avant le hachage.
  • Insertion lente après plusieurs millions de lignes. L'index ne tient plus dans le cache. Augmentez le cache SQLite avec PRAGMA cache_size=-200000 (cela fait 200 Mo) et assurez-vous que les insertions se font par lots dans une transaction par page, et non ligne par ligne.

Étape 5 : Parallélisme sans perte

Objectif de l'étape : accélérer l'extraction avec plusieurs workers simultanés via proxy, de sorte que la chute de l'un d'eux ne perde aucune tâche et ne casse pas la base.

Quand paralléliser est possible, quand non

La pagination par curseur est séquentielle par nature : le curseur suivant n'est connu qu'après réception de la page précédente. Impossible de la paralléliser directement. Mais presque toujours, on peut découper l'extraction en shards indépendants : par jours, par catégories, par régions, par premiers caractères de l'identifiant. Chaque shard est extrait séquentiellement avec son propre checkpoint, et les shards avancent en parallèle. La pagination par offset et par clé avec des bornes connues se parallélise directement : des tâches du type « pages 1 à 100 » ou « identifiants de 0 à 100000 ».

File de tâches sur disque

Une file en mémoire meurt avec le processus. C'est pourquoi les tâches vivent dans une table avec des statuts :

  • pending — en attente d'exécution ;
  • running — prise par un worker ;
  • done — exécutée et enregistrée ;
  • failed — tentatives épuisées, nécessite l'attention humaine.

Au démarrage, le loader commence par remettre toutes les tâches running en pending : si elles sont dans ce statut, c'est que le processus précédent est mort en pleine tâche. Ensuite, les workers prennent les pending.

Un point d'écriture unique

SQLite autorise plusieurs lecteurs simultanés, mais un seul écrivain. Le patron le plus simple et le plus sûr : les workers téléchargent et renvoient les données, et tout l'écrit dans la base est fait par le flux principal. Pas de verrou dans le code, pas de « database is locked ». De plus, activez le mode WAL pour que la lecture de l'état depuis un autre processus ne gêne pas l'écriture.

Limitation de la concurrence

Limitez le nombre de workers par deux choses. Premièrement, les capacités du proxy : si dans votre espace client Proxeon vous avez plusieurs canaux, il est raisonnable de garder un à deux workers par canal, pour que la rotation d'IP sur un canal ne coupe pas les connexions de tous les threads d'un coup. Deuxièmement, la politesse envers la source : même sans limites formelles, dix flux parallèles sur une petite API créeront une charge, à cause de laquelle vous recevrez des refus. Commencez avec trois à quatre workers et augmentez en observant le taux d'erreur.

Code du loader parallèle

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: # données et statut de la tâche dans une seule transaction
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('file vide, tâches en erreur :', failed)

La fonction fetch_fn exécute la requête via proxy et renvoie une liste d'enregistrements. Elle travaille dans le thread et ne touche pas à la base. La fonction write_fn est appelée dans le flux principal à l'intérieur d'une transaction et fait l'insertion idempotente de l'étape 3. La tâche en échec est automatiquement remise en pending et sera reprise au prochain cycle claim ; après épuisement des tentatives, elle reçoit le statut failed, et vous vous en occuperez manuellement.

Attention : ne transmettez pas l'objet de connexion sqlite3 aux workers. La connexion est liée au thread dans lequel elle a été créée, et une tentative de l'utiliser depuis un autre thread entraînera une erreur ou, pire, une corruption silencieuse des données. Soit chaque thread a sa propre connexion, soit, comme dans l'exemple ci-dessus, aucune.

Conseil : utilisez pour chaque worker un objet requests.Session distinct avec sa propre adresse proxy. Si vous avez plusieurs canaux dans Proxeon, répartissez-les sur les workers en rond : worker 0 prend le canal 0, worker 1 le canal 1, etc. Ainsi, la chute de connexion sur un canal n'affectera qu'un seul thread.

Résultat attendu

Vérification : mettez 100 tâches en file, lancez quatre workers et tuez le processus au bout d'une demi-minute. Exécutez SELECT status, count(*) FROM tasks GROUP BY status;. Vous verrez quelques done, quelques running et le reste en pending. Relancez : les running doivent disparaître au démarrage, et à la fin toutes les tâches doivent être en done, sauf celles qui ont vraiment échoué et se trouvent en failed avec le texte de l'erreur dans last_error.

Étape 6 : Reprise après une longue pause

Objectif de l'étape : reprendre correctement une extraction interrompue pendant plusieurs heures ou jours, sans tomber sur un état périmé.

Ce qui se périme

Une reprise après dix secondes et après une semaine, ce sont deux problèmes différents. Après une longue pause, une partie de l'état sauvegardé cesse d'être valide.

  1. Sessions et cookies. Les sessions côté serveur vivent en général de quelques heures à une journée. Les cookies sauvegardés après cela mèneront à des réponses 401 ou à une redirection vers le formulaire de connexion. Solution : au démarrage, effectuer une authentification complète à nouveau, plutôt que de restaurer les cookies depuis un fichier.
  2. Tokens d'accès. Les tokens OAuth vivent une heure, parfois moins. Si vous avez un refresh token, mettez à jour l'access token avant le démarrage et selon un calendrier pendant le travail, sans attendre le refus.
  3. Curseurs de pagination. Beaucoup d'API limitent la durée de vie d'un curseur à des minutes ou des heures. Un curseur périmé renverra une erreur 400 avec un message sur un curseur invalide. C'est précisément pour cela qu'à l'étape 2 nous avons sauvegardé un ancrage de secours : l'identifiant du dernier enregistrement. Si la source supporte un filtre par identifiant ou par date de modification, construisez une nouvelle requête depuis cet ancrage. Si elle ne le supporte pas, il faudra recommencer le shard depuis le début, et la déduplication de l'étape 4 écartera ce qui a déjà été obtenu.
  4. Contenu du fichier sur le serveur. Pour la reprise de l'étape 1, il est crucial que le fichier n'ait pas changé. Comparez l'ETag actuel avec celui sauvegardé avant chaque tentative ; en cas de non-correspondance, recommencez le fichier.
  5. Paramètres du proxy. En une semaine, dans l'espace client Proxeon, le port, le mot de passe ont pu changer, ou la durée du canal expirer. Vérifiez le proxy par une requête de test avant de commencer à traiter la file.
  6. L'ensemble des données lui-même. Si l'extraction dure une semaine, et que la source a ajouté et supprimé des enregistrements pendant ce temps, votre résultat sera un mélange d'états à différents moments. Pour beaucoup de tâches, c'est acceptable. Sinon, stockez l'heure de démarrage dans le checkpoint et, après la fin, faites une passe incrémentale séparée sur les enregistrements modifiés après cette heure.

Vérification pré-vol

Rassemblez toutes les vérifications dans une seule fonction exécutée au démarrage avant tout travail réel. Soit elle remet l'état en ordre, soit elle arrête le loader avec un message clair.

def preflight(session, state, probe_url):
# 1. le proxy est vivant et autorisé
r = session.head(probe_url, timeout=20)
if r.status_code == 407:
raise SystemExit('le proxy a rejeté le login ou le mot de passe, vérifiez les données dans l espace client Proxeon')
# 2. le token d'accès est frais
refresh_access_token(session)
# 3. le curseur est encore valide
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('curseur périmé, bascule sur l ancrage last_id =', state['last_id'])
state['cursor'] = None
state['resume_after_id'] = state['last_id']
# 4. rappel de l'âge de l'extraction
print('extraction démarrée', state['started_at'], 'pages enregistrées', state['pages'])
return state

La fonction refresh_access_token dépend de votre source : c'est généralement une requête POST avec le refresh token, après quoi vous mettez à jour l'en-tête Authorization dans la session. Le champ resume_after_id est ensuite utilisé dans la fonction de récupération de page comme filtre « identifiant supérieur à celui indiqué ».

Conseil : gardez le refresh token et le mot de passe proxy non pas dans le checkpoint, mais dans des variables d'environnement ou dans un fichier de secrets séparé avec des droits d'accès restreints. Vous copierez le checkpoint, l'enverrez à des collègues et le joindrez à des rapports d'erreurs ; les secrets n'y ont pas leur place.

Résultat attendu

Vérification : corrompez manuellement le curseur dans la table checkpoint avec UPDATE checkpoint SET cursor='broken'; et lancez le loader. La console doit afficher une ligne sur la bascule vers l'ancrage last_id, et l'extraction doit continuer sans planter. Le nombre d'enregistrements après la fin doit coïncider avec une exécution de contrôle sans corruption du curseur.

Étape 7 : Squelette prêt à l'emploi d'un loader robuste en Python

Objectif de l'étape : rassembler tout ce qui précède dans un seul fichier, que l'on peut lancer, interrompre, relancer et obtenir un résultat complet sans doublons.

Structure du squelette

  • Configuration par variables d'environnement : adresse du proxy Proxeon, adresse de l'API, token, nom du job.
  • Classe Store : SQLite avec les tables records et checkpoint, une transaction par page.
  • Fonction de clé d'enregistrement pour l'idempotence.
  • Fonction de récupération de page via proxy.
  • Boucle principale avec reprise par checkpoint et recréation de session après coupure.

Code complet

# 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 depuis l espace client 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émarrage : pages %d, lignes %d, uniques en base %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('tentatives épuisées, état sauvegardé, relancez plus tard')
raise
print('coupure (%s), tentative %d sur %d' % (type(e).__name__, attempts, MAX_ATTEMPTS))
time.sleep(min(60, 2 ** attempts)) # choix des pauses décrit dans l'article sur 429 et retries
session = make_session() # nouvelle session : la connexion via proxy est recréée
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('pages %d, lignes %d' % (state['pages'], state['rows']))
if next_cursor is None:
break
print('terminé : pages %d, lignes reçues %d, uniques en base %d'
 % (state['pages'], state['rows'], store.unique_count()))


if __name__ == '__main__':
run()

Comment adapter à votre source

  1. Remplacez le chemin /orders et les noms des champs items, next_cursor, id par ceux que renvoie votre API. Cela fait trois endroits dans les fonctions fetch_page et record_key.
  2. Si la source a une pagination par offset, remplacez le paramètre cursor par page et calculez la valeur suivante comme state['pages'] + 1. Dans le checkpoint, sauvegardez le numéro de page au lieu du curseur.
  3. Si la pagination est par clé, transmettez un paramètre du type after_id issu de state['last_id'] et supprimez la gestion du curseur.
  4. Si vous avez besoin de parallélisme, sortez fetch_page dans fetch_fn de l'étape 5, et utilisez commit_page comme write_fn. Découpez l'extraction en shards et remplissez la file de tâches.
  5. Ajoutez la fonction preflight de l'étape 6 avant la boucle principale.

Vérification du résultat : checklist

Avant de lancer le loader sur un volume réel de plusieurs heures, passez-le à travers cette liste. Chaque point prend quelques minutes, et ensemble ils garantissent qu'une extraction nocturne ne perdra pas de données.

  1. Test d'interruption. Lancez le loader, au bout de 30 secondes appuyez sur Ctrl+C. Relancez. La première ligne de sortie doit afficher un nombre de pages non nul, et non « pages 0 ».
  2. Test de doublons. Réduisez PAGE_SIZE à 10, interrompez le loader cinq fois de suite à des moments aléatoires. À la fin, comparez le nombre d'enregistrements uniques avec le nombre de lignes reçues : les uniques doivent être inférieurs ou égaux, et avec une pagination par curseur propre sans insertions dans la source, presque égaux.
  3. Test de cohérence. Après n'importe quelle interruption, exécutez deux requêtes : SELECT rows FROM checkpoint; et SELECT count(*) FROM records;. La différence entre elles ne doit pas dépasser une taille de page. Si elle dépasse, le checkpoint et les données ne sont pas écrits dans la même transaction.
  4. Test du proxy. Désactivez temporairement la variable PROXY_URL ou indiquez un mot de passe erroné. Le loader doit planter à la première requête avec une erreur claire, pas se figer ni commencer à aller directement à la source.
  5. Test du curseur périmé. Corrompez le curseur dans la base, comme décrit à l'étape 6, et vérifiez que la bascule vers l'ancrage se déclenche.
  6. Test du disque. Vérifiez la taille du fichier export.sqlite après mille pages et multipliez par le nombre de pages attendu. Assurez-vous que l'espace disque suffit avec 20 pour cent de marge.

Vérification : le signe d'une exécution réussie est une situation où, après trois interruptions intentionnelles et trois redémarrages, le nombre final d'enregistrements uniques coïncide avec celui obtenu en une seule exécution continue sur la même source, et où la ligne « repart de zéro » n'est jamais apparue dans la console.

Possibilités supplémentaires et optimisation

  • Progression et estimation du temps. Si le nombre total d'enregistrements est connu, affichez le pourcentage et l'estimation du temps restant toutes les vingt pages. C'est utile pour vous et pour distinguer un blocage d'un fonctionnement lent.
  • Compression du payload. Pour des dizaines de millions d'enregistrements, le JSON en texte prend beaucoup de place. Compressez le champ payload avec zlib.compress avant l'écriture et stockez-le en BLOB. L'économie est généralement de trois à huit fois.
  • Migration vers une base serveur. Le schéma avec checkpoint dans la même transaction que les données se transpose presque sans changement sur PostgreSQL. La construction ON CONFLICT y est supportée, et la contrainte « un seul écrivain » disparaît.
  • Extractions incrémentales. Sauvegardez l'heure de démarrage de chaque tâche et, après une extraction complète, lancez une tâche séparée avec le filtre « modifié après ». Ainsi, vous maintenez une copie à jour sans tout re-télécharger.
  • Processus séparé par shard. Au lieu de threads, vous pouvez lancer plusieurs instances du script avec des valeurs JOB différentes et des fichiers DB_PATH différents, puis fusionner les résultats à la fin. C'est plus simple à déboguer et cela élimine complètement la question de l'écriture concurrente.
  • Métriques de coupures. Journalisez chaque coupure avec le type d'exception et l'heure. Au bout d'une journée, vous verrez que les coupures se regroupent autour des intervalles de rotation d'IP sur le canal Proxeon, et vous pourrez ajuster l'intervalle de rotation à la durée de vos requêtes.

Erreurs typiques et solutions

Voici les situations auxquelles presque tout le monde fait face aux premiers lancements. Format : problème, cause, solution.

  1. Problème : après un redémarrage, l'extraction repart de zéro à chaque fois. Cause : le checkpoint est écrit en mémoire ou dans un fichier qui ne survit pas au plantage, ou le loader ne le lit pas au démarrage. Solution : assurez-vous que la première action dans la fonction run est store.load, et que l'état est mis à jour après chaque page dans une transaction.
  2. Problème : la base contient une fois et demie plus d'enregistrements que la source. Cause : la clé de déduplication n'est pas stable : l'heure de la requête, le numéro de page ou un champ à l'ordre aléatoire y figurent. Solution : calculez la clé uniquement sur l'identifiant de la source ou sur une liste explicite de champs stables via le paramètre fields.
  3. Problème : la base contient moins d'enregistrements que la source, alors que l'extraction s'est terminée sans erreur. Cause : le checkpoint a été mis à jour avant l'écriture des données, et après un plantage la page a été sautée. Ou bien la pagination par offset, lors de la suppression d'enregistrements dans la source, a décalé les pages en arrière. Solution : remplacez l'ordre par « données, puis checkpoint » dans une seule transaction ; pour les sources avec suppressions, passez à la pagination par clé.
  4. Problème : la reprise du fichier donne une archive corrompue avec une taille identique. Cause : le serveur a répondu une fois 200 au lieu de 206, le fichier a été partiellement réécrit, puis complété. Solution : conservez l'ETag à côté du fichier, supprimez le fichier en cas de changement ; vérifiez que l'en-tête Content-Range correspond au décalage demandé.
  5. Problème : erreur « database is locked » en travail parallèle. Cause : plusieurs threads écrivent dans SQLite simultanément, ou la connexion est transmise entre threads. Solution : un point d'écriture unique dans le flux principal, les workers ne font que télécharger ; mode WAL ; la connexion est créée dans le thread où elle est utilisée.
  6. Problème : au bout d'une heure, toutes les requêtes renvoient 401. Cause : le token d'accès a expiré. Solution : mettez à jour le token selon un calendrier avant expiration, et en cas de 401, déclenchez la mise à jour et relancez la requête une fois, sans la considérer comme une coupure.
  7. Problème : la mémoire vive monte à plusieurs gigaoctets. Cause : l'ensemble des clés vues ou la liste de tous les enregistrements est gardé en mémoire du processus. Solution : déduplication via la clé primaire en base ou la table seen ; écriture des données page par page, sans accumulation.
  8. Problème : les coupures se produisent strictement toutes les quelques minutes. Cause : elles coïncident avec l'intervalle de rotation d'IP sur le canal proxy. Solution : c'est une situation normale, le loader doit la surmonter. Si les requêtes sont longues, ajustez l'intervalle de rotation dans l'espace client Proxeon pour qu'il soit nettement supérieur au temps typique d'une requête, ou utilisez la rotation à la demande entre les pages.

FAQ : questions fréquentes sur l'extraction robuste

Est-il obligatoire d'utiliser SQLite si le résultat doit être en CSV ?

Non, mais c'est pratique. SQLite joue ici le rôle de stockage fiable de l'état et de l'ensemble des clés vues. Le CSV final, vous l'exportez en une commande depuis la table records une fois terminé. Si vous voulez vraiment sans base, utilisez la variante fichier de l'étape 3 avec un fichier par page et un checkpoint JSON avec remplacement atomique.

À quelle fréquence sauvegarder le checkpoint : après chaque page ou moins souvent ?

Après chaque page. Une transaction SQLite avec quelques centaines de lignes s'exécute en quelques millisecondes, c'est négligeable par rapport à une requête réseau via proxy. L'économie sur des checkpoints rares ne vaut pas le risque de perdre des dizaines de pages.

Que faire si l'API ne renvoie ni curseur, ni identifiants, seulement des numéros de page ?

Travaillez par numéro de page, stockez-le dans le checkpoint et incluez impérativement une déduplication par hachage du contenu de l'enregistrement. Acceptez que, en cas de modifications actives dans la source, une partie des enregistrements puisse être manquée à cause du décalage des pages. Pour les données critiques, faites une seconde passe dans l'ordre inverse des pages : ce qui a été manqué à la première passe aura de fortes chances de tomber dans la seconde.

Peut-on reprendre le téléchargement d'un fichier avec plusieurs threads sur différentes plages ?

Oui, si le serveur supporte Range. Découpez le fichier en morceaux de 50-100 Mo, chaque morceau est une tâche de la file de l'étape 5 avec son propre fichier temporaire, et à la fin de tous les morceaux, assemblez-les dans le bon ordre. Vérifiez chaque morceau par sa taille, et le fichier entier par la somme de contrôle, si elle existe.

Combien de workers mettre en travaillant via proxy ?

Commencez avec trois à quatre par canal Proxeon et observez le taux d'erreur dans la table tasks. Si les erreurs sont inférieures à un pour cent, ajoutez-en deux. Si les erreurs augmentent, réduisez. Plus de dix flux sur un canal donnent rarement un gain : vous butez soit sur la bande passante du canal, soit sur la patience de la source.

Faut-il sauvegarder les cookies de session entre les lancements ?

En général non. Une nouvelle authentification au démarrage prend quelques secondes et est plus fiable que la restauration de cookies à la durée de vie inconnue. Exception : la source limite le nombre de connexions par jour. Dans ce cas, sauvegardez les cookies, mais au premier 401 ou redirection vers la connexion, jetez-les et authentifiez-vous à nouveau.

Comment savoir que l'extraction s'est terminée complètement, et ne s'est pas arrêtée silencieusement ?

Pour les fichiers : la taille est égale à Content-Length et la somme de contrôle correspond. Pour une API : une page sans next_cursor ou une page vide a été reçue, et le nombre de lignes correspond au total, si la source le communique. Écrivez dans le checkpoint un drapeau de fin explicite, pour qu'un relancement ne commence pas un nouveau parcours.

Que faire des tâches en statut failed ?

Regardez le champ last_error. Si ce sont des erreurs réseau, remettez simplement les tâches en pending avec la commande UPDATE et relancez le loader. Si ce sont des erreurs de parsing, cela signifie qu'il y a dans la source des enregistrements de forme non standard : corrigez le code et relancez. Ne supprimez jamais les failed en silence, c'est la seule preuve de ce qui manque dans l'extraction.

Peut-on utiliser cette approche non pas avec requests, mais avec une bibliothèque asynchrone ?

Oui, les principes sont les mêmes : unité de travail atomique, checkpoint avec les données, écriture idempotente, file sur disque. Seul le transport change. Un seul détail : gardez l'écriture dans SQLite synchrone et séquentielle, et la concurrence au niveau des requêtes réseau.

Conclusion

Vous êtes passé d'un fichier incomplet et d'un redémarrage stressant à un loader pour qui une coupure est indifférente. Fixons ce qui a été fait.

  • Nous avons traité la reprise HTTP : vérification d'Accept-Ranges via HEAD, en-tête Range avec la taille actuelle du fichier, distinction des codes 206 et 200, protection contre la substitution du fichier via ETag et If-Range.
  • Nous avons construit des checkpoints pour les extractions paginées : curseur, numéro de page ou identifiant du dernier enregistrement plus des compteurs auxiliaires, le tout dans une transaction avec les données.
  • Nous avons rendu l'écriture idempotente via la clé primaire et la construction INSERT OR IGNORE ou ON CONFLICT DO UPDATE, et pour les fichiers via des noms déterministes et un remplacement atomique.
  • Nous avons organisé la déduplication sur disque, pour que des millions de clés ne vivent pas en mémoire vive et survivent aux redémarrages.
  • Nous avons ajouté le parallélisme avec une file de tâches dans SQLite, le retour automatique des tâches échouées et un point d'écriture unique.
  • Nous avons prévu la reprise après une longue pause : mise à jour des tokens, nouvelle authentification, bascule d'un curseur périmé vers un ancrage par identifiant, vérification du proxy Proxeon avant le démarrage.
  • Nous avons tout rassemblé dans un squelette opérationnel, qui s'adapte à une source concrète en remplaçant trois ou quatre lignes.

Que faire ensuite

Prenez le squelette de l'étape 7 et lancez-le sur un petit volume réel, disons dix mille enregistrements. Passez la checklist de la section de vérification. Ensuite seulement, lancez l'extraction complète pour la nuit. Le matin, soit vous verrez la ligne « terminé » avec des compteurs cohérents, soit la ligne sur les tentatives épuisées avec l'état sauvegardé, et il suffira alors de relancer le script.

Où progresser

Le niveau suivant, ce sont les extractions incrémentales par heure de modification au lieu de parcours complets, le transfert de l'état vers une base serveur pour plusieurs machines, ainsi qu'une stratégie de retry réfléchie selon les codes de réponse, à laquelle un article dédié sur 429 et les retries est consacré. La combinaison de retries judicieux de cet article et d'un état robuste de celui-ci donne un loader que l'on peut laisser tourner une semaine sans ouvrir le terminal.

Et pour finir. La coupure d'une longue extraction n'est pas une panne, mais une situation de travail que vous savez désormais gérer. Bonnes extractions à vous.