Cómo exportar grandes volúmenes de datos a través de un proxy y no empezar de cero tras una interrupción
Contenido del artículo
- Introducción: por qué una extracción larga casi siempre se interrumpe y eso es normal
- Preparación previa: herramientas, accesos y entorno
- Conceptos básicos: diccionario de la extracción robusta en palabras simples
- Paso 1: reanudación por http a través del encabezado range
- Paso 2: checkpoints para extracciones paginadas
- Paso 3: idempotencia, para que repetir no cree duplicados
- Paso 4: deduplicación de resultados sin inflar la memoria
- Paso 5: paralelismo sin pérdidas
- Paso 6: reanudación tras una pausa larga
- Paso 7: esqueleto listo de un cargador robusto en python
- Verificación del resultado: lista de control de una extracción robusta
- Errores típicos y sus soluciones
- Funciones adicionales y optimización
- Faq: preguntas frecuentes sobre extracción robusta
- Conclusión: lo que ya sabes hacer y hacia dónde avanzar
Introducción: por qué una extracción larga casi siempre se interrumpe y eso es normal
Si alguna vez lanzaste una extracción grande de datos, conoces esa sensación. El proceso llevaba varias horas, llegó al noventa por ciento y se cortó. La conexión se cayó, el servidor devolvió un error, la laptop entró en suspensión. Y todo hay que empezarlo otra vez. Esta guía está escrita para que eso no vuelva a pasar.
Qué vas a lograr al final. Vas a aprender a construir un cargador que sobrevive a las interrupciones. Reanuda los archivos desde el byte donde quedaron, recuerda en qué página se detuvo, no crea duplicados al repetir y puede retomarse incluso después de una pausa larga. Vas a tener un esqueleto listo en Python que puedes adaptar a tu caso.
Para quién es esta guía. Para ingenieros, analistas y desarrolladores que extraen datos de APIs, descargan archivos grandes o arman resultados paginados a través de un proxy. Nivel intermedio. Debes entender lo básico de HTTP y poder leer código en Python. No se requieren conocimientos profundos de programación de redes.
Qué necesitas saber de antemano. Python básico, noción de petición y respuesta HTTP, qué son los encabezados y los códigos de estado. Si trabajaste con la librería requests, con eso basta.
Cuánto tiempo tomará. Leer y entender los conceptos te llevará unos cuarenta minutos. Armar un cargador funcional con nuestro esqueleto, entre una y tres horas, según tu fuente de datos.
Aclaración importante sobre el tema. No vamos a analizar códigos de estado como el 429 ni estrategias de reintentos con espera (backoff). Sobre eso hay material aparte. Aquí el foco es uno solo: el estado del proceso y su reanudación. Cómo guardar el progreso, cómo no perder ni duplicar datos, cómo continuar desde donde te detuviste.
Consejo: Ten a mano un cuaderno o un archivo aparte donde anotes los parámetros de tu fuente: si soporta reanudación, si tiene navegación paginada, cuál es el formato de su cursor. Estas notas te servirán en cada paso.
Preparación previa: herramientas, accesos y entorno
Antes de escribir código, armemos el entorno de trabajo. Toma diez minutos, pero te ahorrará horas de depuración.
Qué necesitas instalar
- Instala Python 3.10 o superior. Verifica la versión con el comando
python --versionen la terminal. - Crea un entorno virtual con el comando
python -m venv venv, para que las dependencias del proyecto no se mezclen con las del sistema. - Activa el entorno. En Windows es
venv\Scripts\activate, en macOS y Linux essource venv/bin/activate. - Instala la librería para peticiones HTTP con el comando
pip install requests. - Para trabajar más rápido con la base de estado no necesitas instalar nada adicional: el módulo
sqlite3ya viene incluido en la librería estándar de Python.
Qué necesitas de los accesos
- Acceso a tu fuente de datos: URL, token o clave de API, si se requiere.
- Un proxy de Proxeon con dirección, puerto y datos de autorización. Sin un proxy estable, una extracción robusta pierde sentido, porque es justamente el proxy el que distribuye la carga y hace que las conexiones sean predecibles.
- Espacio en disco para el archivo de estado y para los datos extraídos.
Verificación del proxy de Proxeon
- Toma la cadena de conexión con el formato
http://usuario:contraseña@dirección:puerto. - Verifícala con una petición simple. En la terminal ejecuta
curl -x http://usuario:contraseña@dirección:puerto https://api.ipify.orgy confirma que devolvió la dirección IP del proxy y no la tuya.
Consejo: Guarda la cadena de conexión al proxy en una variable de entorno, no en el código. Así no enviarás la contraseña por accidente al sistema de control de versiones. En el código léela a través de os.environ.
⚠️ Atención: Trabaja siempre solo con las fuentes de datos a las que tengas acceso legal. Respeta los términos de uso del servicio y los límites que indique su propietario. El proxy de Proxeon está pensado para trabajo de ingeniería legal: distribución de carga, estabilidad de conexiones y extracción de datos correcta.
✅ Verificación: El entorno está listo si el comando python -c "import requests, sqlite3" corrió sin errores y la petición a través del proxy devolvió la dirección del servidor proxy.
Conceptos básicos: diccionario de la extracción robusta en palabras simples
Antes de escribir código, repasemos los términos clave. Sin ellos, los siguientes pasos sonarán a conjuro.
Reanudación
Reanudación es continuar la descarga de un archivo desde el byte en el que se interrumpió. En lugar de descargar el archivo de nuevo, le pides al servidor que entregue solo el fragmento faltante. Funciona a través del encabezado HTTP Range.
Checkpoint
Checkpoint es un punto de progreso guardado. Imagina el guardado en un videojuego. Si algo sale mal, vuelves a la última partida guardada, no al inicio del juego. En una extracción, el checkpoint guarda en qué página o registro te detuviste.
Cursor
Cursor es una marca que la API te da para que puedas pedir la siguiente porción de datos. A menudo es una cadena como eyJvZmZzZXQiOjEwMH0. La envías de vuelta y el servidor entiende desde dónde continuar.
Idempotencia
Idempotencia es la propiedad de una operación en la que repetirla no cambia el resultado. Si escribiste dos veces la misma fila con la misma clave, al final queda una sola fila, no dos. Es la protección contra duplicados al repetir.
Deduplicación
Deduplicación es descartar registros repetidos. Incluso trabajando con cuidado, un mismo objeto puede llegar dos veces. La deduplicación garantiza que en tu conjunto final quede solo uno.
Principio fundamental
Un cargador robusto se construye sobre una idea: el progreso debe guardarse constantemente, no solo al final. Cualquier paso puede ser el último antes de una interrupción. Entonces, después de cada bloque de trabajo exitoso, el estado debe quedar escrito en disco. Así, reanudar es simplemente leer el estado y continuar.
Consejo: Memoriza la regla de las tres preguntas para cualquier extracción. Primera: ¿dónde me detuve? Segunda: ¿cómo no duplicar lo ya obtenido? Tercera: ¿qué se vence mientras no estoy? Las respuestas a estas preguntas son la esencia de la robustez.
Paso 1: Reanudación por HTTP a través del encabezado Range
Objetivo de la etapa. Aprender a descargar un archivo grande de modo que, tras una interrupción, continúe desde el byte pendiente y no desde cero.
Cómo funciona
HTTP permite pedir no todo el archivo, sino una parte. Para eso se añade el encabezado Range a la petición. Por ejemplo, Range: bytes=1048576- significa: entrégame todo a partir del byte número 1048576. Pero primero hay que confirmar que el servidor puede hacerlo.
- Envía al archivo una petición con método HEAD o un GET normal y revisa los encabezados de respuesta.
- Busca el encabezado
Accept-Ranges. Si su valor esbytes, el servidor soporta reanudación. - Si el encabezado no está o su valor es
none, la reanudación no es posible. En ese caso habrá que descargar el archivo completo de una sola vez o buscar una fuente alternativa.
Verificación del soporte de reanudación
Aquí está el código que revisa si el servidor puede entregar partes del archivo.
import requests
def supports_resume(url, proxies):
resp = requests.head(url, proxies=proxies, timeout=30, allow_redirects=True)
accept = resp.headers.get("Accept-Ranges", "none")
total = resp.headers.get("Content-Length")
return accept.lower() == "bytes", totalReanudación del archivo desde la mitad
Ahora el código principal. Revisa cuántos bytes ya están descargados localmente y le pide al servidor solo el resto.
import os
import requests
def download_resumable(url, dest, proxies):
already = 0
if os.path.exists(dest):
already = os.path.getsize(dest)
headers = {}
if already > 0:
headers["Range"] = f"bytes={already}-"
mode = "ab" if already > 0 else "wb"
with requests.get(url, headers=headers, proxies=proxies,
stream=True, timeout=60) as r:
if already > 0 and r.status_code == 200:
mode = "wb"
already = 0
with open(dest, mode) as f:
for chunk in r.iter_content(chunk_size=65536):
if chunk:
f.write(chunk)
return os.path.getsize(dest)Repasemos los puntos importantes. Si el servidor devolvió el estado 206, entregó honestamente una parte del archivo y la escritura por añadido saldrá bien. Si el servidor devolvió 200 pese al encabezado Range, significa que ignoró la reanudación y entrega el archivo completo. En ese caso cambiamos al modo de reescritura total, para no pegar el fragmento viejo con el nuevo y dañar el archivo.
⚠️ Atención: Nunca añadas datos en modo ab si no estás seguro de que el servidor respondió con el código 206. De lo contrario obtendrás un archivo corrupto, donde el inicio es el resto del intento anterior y la continuación es el nuevo archivo completo. Ese archivo se abrirá con error y perderás tiempo buscando la causa.
Consejo: Descarga no directamente al archivo destino, sino a un archivo temporal con extensión .part. Cuando la descarga termine por completo, renómbralo al nombre final. Así nunca confundirás un archivo listo con uno a medio descargar.
Control de integridad
Después de la descarga completa, conviene verificar que el archivo no se dañó. Si el servidor entregaba el encabezado Content-Length, compáralo con el tamaño real del archivo en disco. Si los tamaños coinciden, el archivo llegó completo.
def verify_size(dest, expected):
if expected is None:
return True
return os.path.getsize(dest) == int(expected)✅ Verificación: Interrumpe la descarga a la mitad cerrando el programa. Vuelve a ejecutarlo. En los registros debes ver que la petición salió con el encabezado Range y que el archivo se completó, en lugar de empezar de nuevo. El tamaño final coincide con el esperado.
Paso 2: Checkpoints para extracciones paginadas
Objetivo de la etapa. Configurar el guardado del progreso para APIs que entregan datos por páginas, de modo que tras una interrupción continúe desde la página correcta.
Qué guardar exactamente
Un archivo se reanuda por bytes, y una extracción paginada sigue una lógica completamente distinta. Aquí no hay bytes, hay páginas y registros. Entonces, en el checkpoint hay que guardar otra cosa.
- Cursor si la API trabaja con cursores. Es la opción más confiable, porque el cursor sabe por sí mismo desde dónde continuar.
- Número de página o desplazamiento si la API trabaja con offset y limit. Guarda el número de la última página procesada con éxito.
- Identificador del último registro si se puede ordenar por ID o fecha ascendente. Entonces la siguiente petición pide registros con ID mayor al guardado.
- Contador de registros procesados para control y reporte.
Dónde guardar el estado
Tienes tres opciones principales, de la más simple a la más robusta.
- Archivo JSON. Lo más simple. Escribes un diccionario con el cursor y el contador en un archivo después de cada página. Sirve para extracciones únicas, no paralelas.
- Base SQLite. Más robusta. Ofrece transacciones, así que el estado no se corrompe si hay un corte justo al escribir. Va bien cuando hay muchos datos y se necesita deduplicación.
- Base de datos externa. Para extracciones grandes y distribuidas, cuando varios procesos comparten el mismo trabajo.
Guardado del checkpoint en JSON
import json
import os
def save_checkpoint(path, cursor, page, last_id, count):
tmp = path + ".tmp"
data = {
"cursor": cursor,
"page": page,
"last_id": last_id,
"count": count,
}
with open(tmp, "w") as f:
json.dump(data, f)
os.replace(tmp, path)
def load_checkpoint(path):
if not os.path.exists(path):
return {"cursor": None, "page": 0, "last_id": None, "count": 0}
with open(path) as f:
return json.load(f)Fíjate en el truco del archivo temporal. Escribimos en un archivo con sufijo .tmp y luego lo renombramos atómicamente con os.replace. Esto protege contra la situación en que el programa se cae justo mientras escribe el checkpoint. El checkpoint anterior queda intacto, en lugar de convertirse en medio JSON imposible de leer.
⚠️ Atención: Nunca escribas el checkpoint directamente sobre el mismo archivo viejo sin un archivo temporal. Una interrupción a mitad de la escritura te dejará un checkpoint corrupto y la reanudación será imposible. El reemplazo atómico resuelve este problema por completo.
Consejo: Guarda el checkpoint solo después de que los datos de la página estén realmente escritos en el almacenamiento. El orden es este: recibiste la página, escribiste los datos, luego actualizaste el checkpoint. Si cambias el orden, ante una interrupción saltarás una página y perderás datos.
Ciclo principal con checkpoints
def paginate(fetch_page, save_data, cp_path, proxies):
cp = load_checkpoint(cp_path)
cursor = cp["cursor"]
count = cp["count"]
while True:
items, next_cursor = fetch_page(cursor, proxies)
if not items:
break
save_data(items)
count += len(items)
last_id = items[-1].get("id")
save_checkpoint(cp_path, next_cursor, cp["page"] + 1,
last_id, count)
cursor = next_cursor
if next_cursor is None:
break
return count✅ Verificación: Lanza la extracción, deja que procese varias páginas, interrúmpela. Abre el archivo de checkpoint y confirma que contiene el cursor y el contador actuales. Vuelve a lanzarla: la extracción debe continuar desde el cursor guardado, no desde la primera página.
Paso 3: Idempotencia, para que repetir no cree duplicados
Objetivo de la etapa. Hacer que una nueva ejecución o la repetición de una petición concreta no genere registros iguales en tu almacenamiento.
Por qué surgen los duplicados
Imagina: recibiste una página de datos, la escribiste en el archivo, pero el programa se cayó antes de actualizar el checkpoint. En la siguiente ejecución pedirás la misma página otra vez. Los datos llegarán de nuevo y se escribirán por segunda vez. Así nacen los duplicados. Es una consecuencia inevitable de las interrupciones y hay que combatirla a nivel de arquitectura.
Clave de deduplicación
La herramienta principal de la idempotencia es la clave de deduplicación. Es un campo o una combinación de campos que identifican unívocamente un registro. Elegir bien la clave resuelve la mitad del problema.
- ID natural. Si el registro tiene un identificador único de la fuente, úsalo. Es la clave ideal.
- Combinación de campos. Si no hay un ID único, arma la clave con varios campos estables. Por ejemplo, email más fecha de registro.
- Hash del contenido. Si no hay campos estables en absoluto, calcula un hash de todo el registro. Es la última opción, porque cualquier cambio en un campo creará una clave nueva.
Escritura sin duplicados mediante UPSERT
Si guardas el resultado en SQLite u otra base, usa la inserción con ignorar conflicto. Así, volver a escribir con la misma clave simplemente no hace nada.
import sqlite3
def init_db(path):
conn = sqlite3.connect(path)
conn.execute(
"CREATE TABLE IF NOT EXISTS records ("
"dedup_key TEXT PRIMARY KEY, payload TEXT)"
)
conn.commit()
return conn
def save_records(conn, items):
rows = [(item["id"], json.dumps(item)) for item in items]
conn.executemany(
"INSERT OR IGNORE INTO records (dedup_key, payload) "
"VALUES (?, ?)", rows
)
conn.commit()El detalle clave aquí es la PRIMARY KEY en el campo dedup_key. La base rechazará por sí sola la inserción repetida con la misma clave, porque INSERT OR IGNORE se tragará el conflicto en silencio. No necesitas verificar manualmente si ese registro ya existe. La base lo hace por ti, y lo hace rápido.
Consejo: Elige la clave de deduplicación una sola vez al inicio del proyecto y fíjala en la documentación. Cambiar la clave a mitad de la extracción significa que los registros viejos y nuevos dejarán de corresponderse, y los duplicados sí aparecerán. La estabilidad de la clave importa más que su elegancia.
✅ Verificación: Lanza la extracción dos veces seguidas sobre el mismo rango de datos. Cuenta la cantidad de filas en la base con el comando SELECT COUNT(*) FROM records. El número debe ser igual después de la primera y de la segunda ejecución.
Paso 4: Deduplicación de resultados sin inflar la memoria
Objetivo de la etapa. Descartar registros repetidos en millones de filas sin cargar toda la memoria de la computadora con un montón de claves ya vistas.
Enfoque ingenuo y su problema
La deduplicación más simple: mantener en memoria un conjunto set con todas las claves vistas. Para cada registro nuevo, verificar si la clave está en el conjunto. Funciona perfecto con cientos de miles de filas. Pero con millones y decenas de millones, el conjunto crece y devora gigabytes de memoria RAM. El programa se ralentiza o se cae.
Primera solución: confiar en la base
La forma más simple y confiable con grandes volúmenes es no guardar lo visto en memoria en absoluto, y dejar la verificación a la base mediante la PRIMARY KEY, como hicimos en el paso anterior. La base guarda el índice en disco, no en la memoria de tu proceso. Aguanta decenas de millones de claves sin cargar tu RAM.
Segunda solución: hash del registro
Cuando no hay una clave natural, calcula un hash compacto del registro. El hash ocupa un volumen fijo y pequeño, sin importar el tamaño del registro.
import hashlib
import json
def record_hash(item):
raw = json.dumps(item, sort_keys=True, ensure_ascii=False)
return hashlib.sha256(raw.encode("utf-8")).hexdigest()El parámetro sort_keys=True aquí es crítico. Garantiza que registros con el mismo contenido den el mismo hash, aunque los campos vengan en distinto orden. Sin esa ordenación, dos objetos idénticos pueden obtener hashes distintos y pasar como registros diferentes.
Tercera solución: filtro de Bloom para ahorrar memoria
Si de todos modos necesitas una verificación rápida en memoria con volúmenes enormes, se usa el filtro de Bloom. Es una estructura que ocupa poco espacio y responde rápido si vimos la clave o definitivamente no. Tiene una particularidad: puede decir por error, de vez en cuando, que la clave ya existía aunque no fuera así. Por eso el filtro de Bloom se usa como descarte previo rápido, y la verificación final se deja a la base.
- Verificamos la clave con el filtro de Bloom.
- Si el filtro dice que definitivamente no la vimos, escribimos directo en la base.
- Si el filtro dice que posiblemente la vimos, hacemos la verificación exacta en la base.
⚠️ Atención: No intentes deduplicar decenas de millones de filas con un conjunto común en memoria. En una laptop típica eso agotará la memoria y hará colapsar el proceso a mitad de la extracción. Traslada la carga al disco mediante la base o usa un filtro de Bloom.
Consejo: Si extraes datos por porciones y dentro de una misma porción puede haber duplicados, deduplica la porción en memoria con un conjunto común antes de escribir en la base. La porción es pequeña, la memoria no sufre, y a la base llegan menos inserciones innecesarias.
✅ Verificación: Ejecuta la deduplicación sobre un conjunto de prueba grande con repeticiones intencionales. Verifica que el número final de registros únicos sea correcto y que el consumo de memoria del proceso se mantenga estable, sin crecer linealmente con el número de filas.
Paso 5: Paralelismo sin pérdidas
Objetivo de la etapa. Acelerar la extracción con peticiones en paralelo, sin perder ninguna tarea y reintentando correctamente las que fallen.
Cola de tareas
La base del paralelismo seguro es la cola de tareas. De antemano divides el trabajo en partes independientes. Por ejemplo, una lista de páginas o de rangos. Las pones en una cola. Varios trabajadores toman tareas de la cola, las ejecutan y guardan el resultado. Si un trabajador falla, su tarea puede volver a la cola y entregarse a otro.
Límite de concurrencia
No se puede lanzar un número infinito de peticiones en paralelo. Eso sobrecargará la fuente y tu proxy. El enfoque correcto es limitar la cantidad de trabajadores simultáneos a un valor razonable. Empieza con un número bajo y auméntalo observando la estabilidad.
from concurrent.futures import ThreadPoolExecutor, as_completed
def run_parallel(tasks, worker, proxies, max_workers=5):
results = []
failed = []
with ThreadPoolExecutor(max_workers=max_workers) as pool:
future_map = {
pool.submit(worker, t, proxies): t for t in tasks
}
for future in as_completed(future_map):
task = future_map[future]
try:
results.append(future.result())
except Exception:
failed.append(task)
return results, failedRepetición de tareas fallidas
La lista failed que se arma no son datos perdidos, sino la lista de lo que hay que repetir. Después de la primera pasada, vuelves a ejecutar las tareas fallidas una vez más. Normalmente eso basta para completar el resto.
def run_with_retry(tasks, worker, proxies, rounds=3):
remaining = tasks
for _ in range(rounds):
done, remaining = run_parallel(remaining, worker, proxies)
if not remaining:
break
return remainingEl papel del proxy de Proxeon en el paralelismo. En el trabajo paralelo, el proxy distribuye las conexiones, lo que hace la extracción más estable y predecible. Cada trabajador opera a través de su propia conexión, y la carga no se concentra en un solo punto.
⚠️ Atención: Al escribir en paralelo en un mismo archivo o en un mismo checkpoint surgen condiciones de carrera. Dos trabajadores pueden sobrescribir el estado del otro. Escribe los resultados solo en una base con transacciones o usa un archivo separado por trabajador, y arma el checkpoint consolidado en un hilo aparte.
Consejo: Haz las tareas pequeñas e independientes. Si una tarea abarca un rango demasiado grande, su interrupción descartará mucho trabajo. Las tareas pequeñas se repiten barato y casi sin que se note.
✅ Verificación: Lanza la extracción en paralelo y derriba a propósito parte de los trabajadores. Tras las rondas de repetición, la lista remaining debe quedar vacía y el conjunto final de datos, completo. Compara la cantidad de registros obtenidos con la esperada.
Paso 6: Reanudación tras una pausa larga
Objetivo de la etapa. Continuar correctamente la extracción si entre intentos pasó mucho tiempo, y entender qué pudo vencerse en ese período.
Qué se vence con el tiempo
Una interrupción de un minuto y una pausa de un día son situaciones distintas. Durante una pausa larga, parte de tu estado puede volverse inválido.
- Sesión. Muchos servicios mantienen la sesión por tiempo limitado. Después de una pausa larga el servidor la olvidará, y las peticiones empezarán a devolver error de autorización.
- Token de acceso. Los tokens de API a menudo tienen un tiempo de vida de minutos u horas. Un token vencido hay que renovarlo antes de continuar.
- Cursor. Algunos cursores viven poco. Si el cursor se venció, habrá que empezar desde el punto estable más cercano, por ejemplo por el identificador del último registro.
- Los propios datos. Durante la pausa pueden haber aparecido registros nuevos en la fuente o haber cambiado los viejos. Eso afecta los desplazamientos en la navegación paginada por offset.
Estrategia de reanudación segura
- Al arrancar, verifica la antigüedad del checkpoint. Si es viejo, prepárate para que parte del estado esté desactualizado.
- Renueva el token de acceso y crea una nueva sesión antes de la primera petición. No confíes en las viejas.
- Prefiere la reanudación por identificador del último registro, no por número de página. El ID es estable, mientras que el número de página se desplaza si los datos cambiaron.
- Haz una petición de prueba con el cursor guardado. Si devolvió error de cursor inválido, cambia a la reanudación por last_id.
def resume(cp, fetch_by_id, fetch_by_cursor, proxies):
if cp["cursor"]:
try:
return fetch_by_cursor(cp["cursor"], proxies)
except CursorExpired:
pass
return fetch_by_id(cp["last_id"], proxies)Por qué la reanudación por ID es más confiable. Imagina que te detuviste en la página 50 con orden por fecha. Mientras no estabas, se agregaron registros nuevos al inicio. Ahora la página 50 contiene datos completamente distintos, y saltarás parte de los registros. La reanudación por el identificador del último registro no sufre esto: simplemente pides todo lo que sea mayor que el ID guardado.
Consejo: Guarda siempre en el checkpoint tanto el cursor como el identificador del último registro a la vez. El cursor es más rápido, pero el ID es tu cuerda de seguridad por si el cursor se vence durante una pausa larga.
✅ Verificación: Detén la extracción, espera lo suficiente para que el token o el cursor se venzan, y lánzala de nuevo. El cargador debe renovar el token, detectar el cursor vencido y continuar por el identificador sin perder ni duplicar registros.
Paso 7: Esqueleto listo de un cargador robusto en Python
Objetivo de la etapa. Reunir todo lo aprendido en un armazón funcional único, que adaptarás a tu fuente de datos.
Abajo está el armazón que combina checkpoints, deduplicación mediante base, renovación de token y reanudación. Las funciones de obtención de página las pones tú, según la API concreta.
import os
import json
import sqlite3
import requests
class ResilientLoader:
def __init__(self, cp_path, db_path, proxies):
self.cp_path = cp_path
self.proxies = proxies
self.conn = sqlite3.connect(db_path)
self.conn.execute(
"CREATE TABLE IF NOT EXISTS records ("
"dedup_key TEXT PRIMARY KEY, payload TEXT)"
)
self.conn.commit()
def load_cp(self):
if not os.path.exists(self.cp_path):
return {"cursor": None, "last_id": None, "count": 0}
with open(self.cp_path) as f:
return json.load(f)
def save_cp(self, cp):
tmp = self.cp_path + ".tmp"
with open(tmp, "w") as f:
json.dump(cp, f)
os.replace(tmp, self.cp_path)
def save_records(self, items):
rows = [(str(i["id"]), json.dumps(i)) for i in items]
self.conn.executemany(
"INSERT OR IGNORE INTO records "
"(dedup_key, payload) VALUES (?, ?)", rows
)
self.conn.commit()
def run(self, fetch_page):
cp = self.load_cp()
while True:
items, next_cursor = fetch_page(
cp["cursor"], cp["last_id"], self.proxies
)
if not items:
break
self.save_records(items)
cp["count"] += len(items)
cp["last_id"] = items[-1]["id"]
cp["cursor"] = next_cursor
self.save_cp(cp)
if next_cursor is None:
break
return cp["count"]Ejemplo de función de obtención de página para tu fuente. Aquí implementas la lógica de la petición y el análisis de la respuesta.
def fetch_page(cursor, last_id, proxies):
params = {"limit": 100}
if cursor:
params["cursor"] = cursor
elif last_id:
params["after_id"] = last_id
r = requests.get(
"https://example-source/api/records",
params=params, proxies=proxies, timeout=60
)
r.raise_for_status()
data = r.json()
return data["items"], data.get("next_cursor")Lanzar todo el mecanismo es simple.
proxies = {
"http": os.environ["PROXEON_URL"],
"https": os.environ["PROXEON_URL"],
}
loader = ResilientLoader("state.json", "out.db", proxies)
total = loader.run(fetch_page)
print("Total de registros:", total)Consejo: Agrega al ciclo un registro de cada cien registros: hora, contador, cursor actual. Así verás el progreso y entenderás fácilmente si la extracción se atascó en un punto.
✅ Verificación: Lanza el armazón sobre una fuente real, interrúmpelo a la mitad, vuelve a lanzarlo. El número final de registros tras la reanudación coincidirá con el número total de registros de la fuente, y una nueva ejecución no aumentará el contador de filas únicas.
Verificación del resultado: lista de control de una extracción robusta
Repasa esta lista. Si se cumplen todos los puntos, tu cargador es realmente robusto.
- La reanudación de archivo continúa desde el byte pendiente, no desde cero.
- El cargador maneja correctamente el caso en que el servidor ignora el encabezado Range.
- El checkpoint se guarda después de cada página procesada, no solo al final.
- El checkpoint se escribe de forma atómica mediante archivo temporal y renombrado.
- Una nueva ejecución no crea duplicados en el conjunto final.
- La deduplicación no crece en memoria linealmente con el número de filas.
- Los trabajadores en paralelo no pierden tareas fallidas y las repiten.
- Después de una pausa larga el token se renueva y el cursor vencido se reemplaza por la reanudación por ID.
Cómo probar
- Lanza una extracción completa de un conjunto pequeño y anota el número de registros.
- Lánzala otra vez sobre el mismo conjunto y confirma que el número no cambió.
- Interrumpe la extracción en distintos puntos: al inicio, a la mitad, cerca del final.
- Después de cada interrupción, reinicia y verifica que el resultado final sea el mismo.
Indicadores de éxito. El número de registros únicos es estable entre ejecuciones. El consumo de memoria no crece sin control. La reanudación siempre continúa desde el punto guardado. No hay archivos corruptos tras la reanudación.
Errores típicos y sus soluciones
Problema: el archivo no abre tras la reanudación. Causa: los datos se añadían en modo ab, aunque el servidor respondió con el código 200 y entregó el archivo completo. Solución: verifica el estado de la respuesta, y ante un 200 cambia a reescritura completa del archivo desde cero.
Problema: tras una interrupción la extracción empieza desde la primera página. Causa: el checkpoint se guardaba solo al final del trabajo o no se guardaba. Solución: guarda el checkpoint después de cada página procesada, justo tras escribir los datos.
Problema: el checkpoint no se lee, el JSON está corrupto. Causa: el programa se cayó mientras escribía directamente en el archivo destino. Solución: escribe en un archivo temporal y reemplaza atómicamente con os.replace.
Problema: hay duplicados en el conjunto final. Causa: no hay clave de deduplicación o no es estable. Solución: define una PRIMARY KEY sobre una clave confiable y usa INSERT OR IGNORE.
Problema: el proceso se cae por falta de memoria con grandes volúmenes. Causa: todas las claves vistas se mantienen en un conjunto en memoria. Solución: traslada la verificación de unicidad a la base o aplica un filtro de Bloom.
Problema: tras una pausa larga las peticiones devuelven error de autorización. Causa: el token o la sesión se vencieron durante la pausa. Solución: renueva el token y crea una nueva sesión en cada arranque.
Problema: tras la pausa se omitieron o duplicaron parte de los registros. Causa: la reanudación se hacía por número de página, y los datos en la fuente cambiaron. Solución: reanuda por el identificador del último registro, no por offset.
Problema: los trabajadores en paralelo pierden parte de los datos. Causa: varios trabajadores escribían en un mismo checkpoint y se sobrescribían entre sí. Solución: escribe los resultados en una base con transacciones, no en un archivo de estado compartido.
Funciones adicionales y optimización
Escritura por lotes
No escribas en la base fila por fila. Arma un lote de varios cientos de registros e insértalos de golpe con executemany. Eso acelera la escritura varias veces en grandes volúmenes.
Confirmación periódica
Llama a commit no en cada escritura, sino cada varios cientos de filas. Un commit demasiado frecuente ralentiza la base. Uno demasiado espaciado arriesga perder más datos ante una interrupción. Encuentra el equilibrio según tu carga.
Reporte de progreso
Agrega una estimación del tiempo restante. Conociendo la velocidad de procesamiento de páginas y el número total de registros, calcularás cuánto falta. Es cómodo para extracciones largas.
Almacenamiento separado de datos crudos y procesados
Guarda las respuestas crudas aparte de los registros procesados. Si más adelante cambias la lógica de análisis, no tendrás que volver a descargar los datos. Basta con reprocesar las respuestas crudas con el nuevo analizador.
Consejo: Configura el proxy de Proxeon para que las conexiones sean estables durante toda la extracción. Una conexión estable reduce el número de interrupciones, así que tu cargador recurre menos a la reanudación y trabaja más rápido.
FAQ: preguntas frecuentes sobre extracción robusta
¿Cómo saber si el servidor soporta la reanudación de archivo? Envía una petición HEAD y revisa el encabezado Accept-Ranges. El valor bytes significa que la soporta. La ausencia del encabezado o el valor none significa que la reanudación no es posible.
¿Qué hacer si la API no entrega cursor, solo páginas? Guarda el número de página y, si es posible, el identificador del último registro. Reanuda preferentemente por ID, porque los números de página se desplazan al cambiar los datos.
¿Con qué frecuencia guardar el checkpoint? Después de cada página procesada y escrita con éxito. Así, ante una interrupción, pierdes como máximo una página de trabajo, no toda la extracción.
¿Se puede deduplicar sin base de datos? Con volúmenes pequeños sí, con un conjunto común en memoria. Con millones de filas es peligroso por la memoria. Mejor usa una base con PRIMARY KEY o un filtro de Bloom.
¿Qué elegir como clave de deduplicación? El ID único natural de la fuente, si existe. Si no, una combinación de campos estables. En el peor caso, un hash de todo el registro con ordenación de claves.
¿Por qué la reanudación por ID es más confiable que por número de página? Porque los datos en la fuente pueden cambiar. Los registros nuevos desplazan las páginas, y por número omitirás o duplicarás datos. El ID no depende de eso.
¿Cuántos trabajadores en paralelo poner? Empieza con un número bajo y auméntalo observando la estabilidad y los límites de la fuente. El paralelismo excesivo perjudica más de lo que ayuda.
¿Cómo guardar la cadena de conexión al proxy de forma segura? En una variable de entorno, no en el código. Léela a través de os.environ. Así la contraseña no llegará al sistema de control de versiones.
¿Qué hacer si el cursor se venció durante una pausa larga? Captura el error de cursor inválido y cambia a la reanudación por el identificador guardado del último registro.
¿Hay que verificar la integridad del archivo descargado? Sí. Compara el tamaño real del archivo con el encabezado Content-Length. Si el servidor entrega una suma de control, verifícala también.
Conclusión: lo que ya sabes hacer y hacia dónde avanzar
Recorriste el camino desde una extracción frágil, que se derrumba ante la primera interrupción, hasta un cargador robusto. Ahora tienes todas las herramientas para que una interrupción deje de ser una catástrofe y se vuelva una situación de trabajo habitual.
Lo que dominaste. La reanudación de archivos mediante el encabezado Range con verificación de Accept-Ranges. El guardado del progreso con checkpoints atómicos. La idempotencia a nivel de clave de deduplicación. La deduplicación en millones de filas sin inflar la memoria. El paralelismo con repetición de tareas fallidas. La reanudación tras una pausa larga con renovación de tokens y reemplazo de cursores vencidos. Y lo más importante, un armazón listo en Python que reúne todo esto.
Qué hacer después. Toma tu fuente de datos real y adapta la función de obtención de página. Empieza con un volumen pequeño, depura la reanudación ante interrupciones y luego escala. Configura un proxy de Proxeon estable para que las conexiones sean predecibles durante toda la extracción.
Hacia dónde desarrollarte. Estudia más a fondo la escritura por lotes y los filtros de Bloom. Agrega monitoreo del progreso y estimación de tiempo. Separa el almacenamiento de datos crudos y procesados. Poco a poco