Cómo exportar grandes volúmenes de datos a través de un proxy sin reiniciar desde cero tras una interrupción: guía paso a paso
Contenido del artículo
- Introducción: por qué una exportación larga casi siempre se corta, y eso es normal
- Preparación previa y conceptos básicos
- Paso 1: reanudación de descarga de archivos por http mediante el encabezado range
- Paso 2: checkpoints para exportaciones paginadas
- Paso 3: idempotencia de la escritura de resultados
- 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 resistente en python
- Errores típicos y soluciones
- Faq: preguntas frecuentes sobre exportación resistente
- Conclusión
Introducción: por qué una exportación larga casi siempre se corta, y eso es normal
Si alguna vez has exportado a través de un proxy varios millones de registros o un archivo de decenas de gigabytes, conoces esa sensación. El script llevaba seis horas funcionando, mostraba el 83 por ciento, y de repente falló con un error de conexión. Y todo lo que tienes es un archivo incompleto y la certeza de que tendrás que empezar de nuevo.
Lo primero que hay que aceptar: una exportación larga siempre se corta. No «a veces», no «con mala red», sino siempre, si dura lo suficiente. Las causas son decenas, y la mayoría están fuera de tu control:
- El proxy cambia la dirección IP externa. En los proxies móviles de Proxeon esto es el comportamiento normal: rotación por temporizador o bajo demanda. En el momento del cambio de IP, la conexión TCP abierta se rompe, y el servidor de origen ve ya a otro cliente.
- El servidor de origen cierra la conexión por su propio timeout, se reinicia, despliega una actualización o simplemente responde con un error 5xx.
- Tu propio proceso se reinicia: actualización del sistema, disco lleno, un error en el código con un registro atípico, un Ctrl+C accidental.
- Caduca el token de autorización, expira la sesión, se invalida el cursor de paginación.
- El portátil entra en suspensión, el Wi-Fi cambia de punto de acceso, el proveedor cambia la ruta.
Luchar contra cada una de estas causas por separado no tiene sentido. El enfoque correcto es otro: diseñar la exportación de modo que una interrupción en cualquier momento te cueste no seis horas, sino una página o un fragmento de archivo. A eso está dedicada esta guía.
Qué obtendrás al final
Después de seguir las instrucciones tendrás:
- Una función de reanudación de descarga de archivos desde la mitad mediante el protocolo HTTP con el encabezado Range y verificación de que el servidor lo soporta.
- Un esquema claro de checkpoints para exportaciones paginadas: qué guardar exactamente y dónde almacenar el estado para que no se pierda junto con el proceso.
- Escritura idempotente de resultados, en la que volver a cargar la misma página no crea duplicados.
- Deduplicación que no consume toda la memoria RAM con millones de filas.
- Un cargador paralelo con cola de tareas, reintento de tareas fallidas y límite de concurrencia.
- Un esqueleto listo de cargador resistente en Python que adaptarás a tu fuente en una hora.
Para quién es esta guía
Para desarrolladores y analistas que ya saben hacer peticiones HTTP desde Python y que al menos una vez han perdido los resultados de una exportación larga. Nivel medio: explicamos lo básico, pero no enseñamos a programar desde cero. Los lectores avanzados encontrarán secciones sobre almacenar hashes fuera de memoria y sobre escritura paralela segura en SQLite.
Qué hay que saber de antemano
- Python a nivel de funciones, bucles, diccionarios y manejo de excepciones.
- Fundamentos de HTTP: qué es un método, un encabezado, un código de respuesta, un cuerpo.
- Una idea general de cómo configurar un proxy en la librería requests.
Aclararemos por separado: los códigos de respuesta y las estrategias de reintento con pausa exponencial no se tratan aquí. A eso está dedicado un artículo aparte sobre el error 429 y los reintentos. En esta guía el foco está en otra cosa: en el estado de la exportación y su reanudación. Los reintentos responden a la pregunta «cuándo repetir la petición», y nosotros respondemos a la pregunta «desde qué punto continuar el trabajo después de que los reintentos se agoten y el proceso muera».
Cuánto tiempo tomará
Leer y ejecutar los ejemplos en una fuente de prueba: dos o tres horas. Adaptar el esqueleto a tu API real o servidor de archivos: una o dos horas más, según lo no estándar que sea la paginación. En total, un día de trabajo con margen.
Preparación previa y conceptos básicos
Herramientas y accesos
- Instala Python 3.11 o superior. En 2026 las ramas actuales son 3.12 y 3.13, todos los ejemplos están probados en ellas. Para comprobar la versión: abre la terminal y escribe
python --version. Si ves 3.11 o superior, todo está bien. - Instala la librería requests:
pip install requests. La versión 2.32 o superior es suficiente. El módulo sqlite3 viene en la librería estándar de Python, no hay que instalar nada aparte. - Obtén acceso al proxy. Abre el panel de Proxeon, elige el canal necesario y copia cuatro valores: host, puerto, usuario y contraseña. Normalmente vienen en una sola línea con el formato
http://USER:PASS@HOST:PORT. Usarás esa línea en todo lo que sigue. - Pon la línea del proxy en una variable de entorno, no en el código. En Linux y macOS:
export PROXY_URL=http://USER:PASS@HOST:PORT. En Windows PowerShell:$env:PROXY_URL='http://USER:PASS@HOST:PORT'. Así no subirás la contraseña al repositorio por accidente. - Comprueba que el proxy responde. Ejecuta en la terminal:
curl -x $PROXY_URL -I https://api.example.com/, sustituyendo la dirección de tu fuente. Deberías ver una línea con el código de respuesta, por ejemploHTTP/2 200. Si ves un error de autorización del proxy 407, revisa el usuario y la contraseña.
Requisitos del sistema
Cualquier máquina con 2 GB de memoria RAM libre y un disco donde quepa el resultado de la exportación más un 20 por ciento de margen para los índices de SQLite. Si planeas cargar millones de registros, el disco importa más que la memoria: todo el enfoque se basa en que el estado vive en el disco, no en variables del proceso.
Copias de seguridad
El archivo de estado que crearás abajo (en los ejemplos es export.sqlite) se convertirá en el artefacto más valioso de todo el trabajo. Adquiere el hábito de copiarlo antes de cualquier experimento con el código: cp export.sqlite export.sqlite.bak. Una vez te salvará un día de exportación.
Atención: nunca edites el archivo SQLite a mano mientras el cargador está funcionando. Incluso la lectura desde otro programa en modo incorrecto puede bloquear la escritura y tumbar el proceso. Si necesitas ver el estado, detén el cargador o usa el modo WAL, del que hablaremos en la sección sobre paralelismo.
Términos clave en lenguaje sencillo
- Checkpoint — una marca guardada en disco que indica «hasta aquí todo está exportado y escrito». Tras un corte, el cargador lee el checkpoint y continúa desde ahí.
- Cursor — una cadena opaca que la API devuelve junto con la página y que hay que pasar para obtener la siguiente página. Tú no construyes ni interpretas el cursor.
- Exportación paginada por offset — cuando pides «la página 37 con 500 registros». Esquema simple, pero al añadirse nuevos registros en la fuente las páginas se desplazan y aparecen duplicados o huecos.
- Exportación paginada por clave (keyset) — cuando pides «todo con identificador mayor que 184203, ordenado por identificador». El esquema más resistente para reanudar, si la fuente lo soporta.
- Idempotencia — propiedad de una operación en la que repetirla da el mismo resultado que hacerla una vez. Escribiste una página dos veces, pero en la base está una sola vez.
- Clave de deduplicación — el valor por el que dos registros se consideran el mismo. Lo ideal es que sea el identificador de la fuente; si no lo hay, la clave se calcula como un hash de campos estables.
- Encabezado Range — la forma de pedirle al servidor HTTP que entregue no todo el archivo, sino una parte, por ejemplo los bytes desde 1048576 hasta el final.
- Semántica at-least-once — garantía de que cada registro se obtendrá al menos una vez. Puede haber repeticiones, pero no omisiones. Es la que obtendrás después de esta guía, y la deduplicación eliminará los duplicados.
Principio fundamental
Los siete pasos siguientes se reducen a una sola idea: cada unidad de trabajo debe ser atómica y repetible. Una unidad de trabajo es un fragmento de archivo, una página de API o una tarea de la cola. Atómica significa que el resultado y la marca de su finalización se guardan juntos. Repetible significa que si la unidad se ejecuta dos veces, nada se rompe. Cuando se cumplen ambas propiedades, un corte en cualquier punto se vuelve seguro.
Paso 1: Reanudación de descarga de archivos por HTTP mediante el encabezado Range
Objetivo de la etapa: aprender a descargar un archivo grande a través de un proxy de modo que tras un corte la descarga continúe desde el byte donde se detuvo, y no desde cero.
Cómo funciona
El protocolo HTTP permite al cliente solicitar una parte del recurso. Para ello se añade al encabezado de la petición Range: bytes=INICIO-, donde INICIO es el offset en bytes. Si el servidor soporta peticiones parciales, responde con el código 206 Partial Content y el encabezado Content-Range: bytes INICIO-FIN/TOTAL. Si no lo soporta, ignora el Range y entrega todo el archivo con el código 200. Tu tarea es distinguir estos dos casos.
Para saber de antemano si el servidor soporta reanudación, ayuda una petición HEAD: devuelve solo los encabezados sin cuerpo. Fíjate en Accept-Ranges: bytes. El valor none o la ausencia del encabezado normalmente indican que no hay reanudación, aunque algunos servidores procesan Range correctamente de todos modos, por eso la comprobación final la hacemos por el código de respuesta.
Instrucciones paso a paso
- Haz una petición HEAD a través del proxy y guarda los encabezados Accept-Ranges, Content-Length, ETag y Last-Modified. ETag hará falta para entender si el archivo en el servidor cambió entre tus intentos.
- Mira cuántos bytes hay ya en el archivo local. Si el archivo no existe, considera que son cero.
- Si el tamaño local ya es igual a Content-Length, el archivo está completo y no hay que hacer nada.
- Si el tamaño local es mayor que cero y el servidor declara soporte de rangos, añade a la petición el encabezado Range con el tamaño actual. Añade también
If-Rangecon el ETag guardado: entonces el servidor entregará una respuesta parcial solo si el archivo no cambió, y si no, devolverá el archivo completo con el código 200. - Envía obligatoriamente
Accept-Encoding: identity. Sin él, el servidor puede aplicar compresión al vuelo, y los offsets de bytes dejarán de coincidir con tu archivo. - Abre el archivo local en modo de adición
absi recibiste 206, o en modo de sobrescriturawbsi recibiste 200. - Lee el cuerpo en streaming en fragmentos de 256 KB y escribe en disco. No cargues toda la respuesta en memoria.
- Al terminar, compara el tamaño final con Content-Length. Si no coincide, significa que la conexión se cortó silenciosamente y hace falta otra pasada.
Código funcional
import os
import requests
PROXY_URL = os.environ['PROXY_URL'] # строка из личного кабинета Proxeon
PROXIES = {'http': PROXY_URL, 'https': PROXY_URL}
def probe(url):
r = requests.head(url, proxies=PROXIES, allow_redirects=True, timeout=30,
headers={'Accept-Encoding': 'identity'})
r.raise_for_status()
return {
'ranges': r.headers.get('Accept-Ranges', 'none').lower(),
'length': int(r.headers.get('Content-Length', 0) or 0),
'etag': r.headers.get('ETag'),
}
def download_resumable(url, path):
meta = probe(url)
have = os.path.getsize(path) if os.path.exists(path) else 0
if meta['length'] and have >= meta['length']:
print('файл уже полный:', have, 'байт')
return True
headers = {'Accept-Encoding': 'identity'}
if have > 0 and meta['ranges'] == 'bytes':
headers['Range'] = 'bytes=%d-' % have
if meta['etag']:
headers['If-Range'] = meta['etag']
with requests.get(url, headers=headers, proxies=PROXIES, stream=True,
timeout=(30, 120)) as r:
if r.status_code == 206:
expected = 'bytes %d-' % have
if not r.headers.get('Content-Range', '').startswith(expected):
raise IOError('сервер отдал не тот диапазон: ' + r.headers.get('Content-Range', ''))
mode = 'ab'
elif r.status_code == 200:
print('сервер отдаёт файл целиком, начинаем с нуля')
mode = 'wb'
have = 0
elif r.status_code == 416:
raise IOError('запрошенный диапазон вне файла, проверьте локальный размер')
else:
r.raise_for_status()
with open(path, mode) as f:
for chunk in r.iter_content(chunk_size=256 * 1024):
if chunk:
f.write(chunk)
have += len(chunk)
if meta['length'] and have != meta['length']:
print('обрыв: получено %d из %d' % (have, meta['length']))
return False
return True
def download_until_done(url, path, max_rounds=50):
for i in range(max_rounds):
try:
if download_resumable(url, path):
return
except (requests.ConnectionError, requests.Timeout, IOError) as e:
print('заход %d прерван: %s' % (i + 1, type(e).__name__))
# пауза перед следующим заходом: стратегия описана в статье про 429 и ретраи
raise RuntimeError('не удалось докачать файл за %d заходов' % max_rounds)
if __name__ == '__main__':
download_until_done('https://files.example.com/export-2026.csv.gz', 'export-2026.csv.gz')Fíjate en la función download_until_done: no contiene lógica de espera entre intentos. Esto es intencional. Inserta ahí tu estrategia de pausas del artículo sobre reintentos, aquí solo importa el ciclo «comprobó tamaño, reanudó descarga, comprobó de nuevo».
Consejo: si el archivo se distribuye como comprimido, no lo descomprimas al vuelo durante la reanudación. Primero obtén el archivo completo, verifica el tamaño y, si el servidor proporciona una suma de control, compruébala. Solo entonces descomprime. Un gzip descargado parcialmente parece dañado y perderás tiempo buscando un error inexistente.
Resultado esperado
Verificación: ejecuta el script con un archivo de al menos 200 MB, interrúmpelo a los diez segundos con Ctrl+C. Mira el tamaño del archivo local, por ejemplo 41 943 040 bytes. Vuelve a ejecutar el script. En la consola no debe aparecer la línea «начинаем с нуля», y el tamaño del archivo debe seguir creciendo, no reiniciarse. Al finalizar, el tamaño total debe coincidir exactamente con el Content-Length de la petición HEAD.
Posibles problemas
- El servidor siempre devuelve 200 en lugar de 206. Significa que no soporta reanudación. La única salida para esa fuente es descargar el archivo completo en una sola pasada con un timeout grande o buscar en la fuente un formato alternativo de exportación por partes, por ejemplo dividido por fechas.
- No hay encabezado Content-Length. El servidor entrega el archivo en modo chunked sin declarar el tamaño. No se puede verificar la completitud por tamaño, ni reanudar: Range requiere offsets conocidos. Llega a un acuerdo con la fuente o usa la suma de control, si se publica.
- Respuesta 416 en la primera pasada. El archivo local es más grande que el archivo en el servidor. El archivo en el servidor cambió y se hizo más corto. Elimina el archivo local y empieza de nuevo.
- El tamaño coincide, pero el archivo está corrupto. Lo más probable es que en algún punto intermedio hubo un corte con código 200 y escritura desde cero, y luego adición. Vuelve a crear el archivo. Para que no se repita, guarda el ETag en un archivo aparte junto al archivo y compáralo antes de cada pasada.
Paso 2: Checkpoints para exportaciones paginadas
Objetivo de la etapa: guardar el estado de la exportación para que tras cualquier caída el proceso continúe desde la última página escrita con éxito.
Qué guardar
El checkpoint mínimo depende del tipo de paginación de la fuente. Veamos tres casos.
- Paginación por cursor. La API devuelve junto con los datos un campo como
next_cursor. Guarda precisamente ese. Es el caso más simple: el cursor ya contiene todo lo que el servidor necesita para continuar. - Paginación por número de página u offset. Guarda el número de la última página completamente escrita y el tamaño de la página. Recuerda que al añadirse registros en la fuente durante la exportación los offsets se desplazan, por eso la deduplicación del paso 4 es obligatoria.
- Paginación por clave. Guarda el identificador del último registro escrito. Al reanudar, pides todo lo que sea mayor que ese identificador. El esquema no teme ni a inserciones ni a pausas largas.
Independientemente del tipo de paginación, conviene añadir al checkpoint campos auxiliares: el identificador del último registro (aunque sea esquema por cursor, es un ancla de respaldo por si el cursor caduca), contadores de páginas y filas para controlar el progreso, la hora de inicio de la exportación y la hora de la última actualización.
Dónde guardar el estado
Hay dos variantes funcionales, y ambas son mejores que las variables en memoria.
Variante A: archivo JSON con reemplazo atómico
Es adecuada si los resultados se escriben en archivos separados, no en una base. La trampa principal: si escribes el estado directamente en el archivo destino y el proceso cae en medio de la escritura, obtendrás un JSON truncado que no se podrá leer. La solución es escribir en un archivo temporal junto al principal y renombrarlo sobre el principal. La operación de renombrado dentro del mismo sistema de archivos es atómica.
import json
import os
import tempfile
def save_state(path, state):
directory = os.path.dirname(os.path.abspath(path))
fd, tmp = tempfile.mkstemp(dir=directory, prefix='.state-')
with os.fdopen(fd, 'w', encoding='utf-8') as f:
json.dump(state, f, ensure_ascii=False)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
def load_state(path, default):
if not os.path.exists(path):
return dict(default)
with open(path, encoding='utf-8') as f:
return json.load(f)Variante B: tabla en SQLite junto a los datos
Es la variante preferida si guardas los registros en una base. El checkpoint se actualiza en la misma transacción que la inserción de filas de la página. O se escriben los datos y la marca, o no se escribe nada. No existe la posibilidad de un desajuste entre «datos hay, marca no».
import sqlite3
con = sqlite3.connect('export.sqlite')
con.executescript('''
CREATE TABLE IF NOT EXISTS records(
id TEXT PRIMARY KEY,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS checkpoint(
job TEXT PRIMARY KEY,
cursor TEXT,
last_id TEXT,
pages INTEGER NOT NULL DEFAULT 0,
updated_at TEXT
);
''')
def commit_page(job, rows, next_cursor, pages, ts):
with con: # одна транзакция на страницу
con.executemany(
'INSERT OR IGNORE INTO records(id, payload, fetched_at) VALUES (?, ?, ?)',
[(str(r['id']), json.dumps(r, ensure_ascii=False), ts) for r in rows])
con.execute(
'INSERT INTO checkpoint(job, cursor, last_id, pages, updated_at) VALUES (?, ?, ?, ?, ?) '
'ON CONFLICT(job) DO UPDATE SET cursor=excluded.cursor, last_id=excluded.last_id, '
'pages=excluded.pages, updated_at=excluded.updated_at',
(job, next_cursor, str(rows[-1]['id']), pages, ts))Orden de las operaciones
Recuerda la regla: primero los datos, después el checkpoint, y preferiblemente en una sola transacción. Si la transacción no es posible (por ejemplo, los datos se escriben en archivos), el orden es exactamente ese: escribir el archivo de la página, sincronizar a disco, luego actualizar el estado. Si el proceso cae entre estas dos acciones, obtendrás una recarga de la misma página, lo cual es seguro gracias al paso 3. El orden inverso provocará la omisión de una página, y eso ya es pérdida de datos.
Consejo: guarda en el checkpoint no el cursor actual, sino el cursor de la siguiente página que devolvió el servidor. Así al reanudar pides directamente lo que aún no tienes, sin una petición extra de la página ya obtenida.
Resultado esperado
Verificación: inicia la exportación, espera diez páginas y termina el proceso a la fuerza. Abre la base con el comando sqlite3 export.sqlite y ejecuta SELECT pages, last_id FROM checkpoint;. Debes ver el número 10 y el identificador. Luego ejecuta SELECT count(*) FROM records; y asegúrate de que el número de registros sea igual a diez tamaños de página. Vuelve a iniciar el cargador: el primer mensaje en la consola debe ser algo como «старт: страниц 10».
Posibles problemas
- Error «database is locked». Otro proceso mantiene la conexión. Cierra todas las ventanas de sqlite3 y otras herramientas que hayan abierto el archivo. Para trabajo multihilo activa el modo WAL, consulta el paso 5.
- El cursor se guardó, pero los datos no. Actualizaste el checkpoint fuera de la transacción con los datos. Vuelve al código de arriba y asegúrate de que ambas operaciones estén dentro del mismo bloque
with con:. - El JSON de estado quedó vacío o corrupto. Escribiste directamente en el archivo sin archivo temporal ni reemplazo. Usa la función save_state tal cual.
Paso 3: Idempotencia de la escritura de resultados
Objetivo de la etapa: lograr que el reprocesamiento de cualquier página no genere duplicados ni dañe los datos.
Por qué la repetición es inevitable
Después del paso 2 ya has visto el escenario en el que una página se escribe dos veces: el proceso cayó después de insertar los datos, pero antes de actualizar el checkpoint. Además, las repeticiones vienen de la paginación por offset al cambiar la fuente, de los trabajadores paralelos que recibieron la misma tarea tras un reinicio, y simplemente de un reinicio manual «por si acaso». Luchar contra las repeticiones desde la petición es inútil. Lo correcto es hacer que la propia escritura sea tal que la repetición sea inofensiva.
Elección de la clave de deduplicación
- Hay identificador en la fuente. Úsalo. Es el campo
id,uuid,order_numbero similar que la fuente garantiza único. Si exportas de varias fuentes a una sola tabla, haz una clave compuesta: nombre de la fuente más identificador. - No hay identificador, pero hay un conjunto de campos que juntos definen el registro. Por ejemplo, para una fila de lista de precios es el artículo más el almacén más la fecha. Compón la clave con esos campos, normalizándolos: pasa las cadenas a un mismo caso, elimina espacios en los bordes, convierte las fechas a un formato único.
- No hay nada estable. Entonces la clave es el hash de todo el registro después de la canonización. Este caso se analiza en detalle en el paso 4. Ten en cuenta que si la fuente modifica el registro (actualiza el precio), el hash cambiará y obtendrás ambas versiones. A veces eso es exactamente lo que se necesita, a veces no.
Atención: no uses como clave el número de orden de la fila en la respuesta ni el número de página. Estos valores cambian con cualquier modificación de la fuente, y la deduplicación se convertirá en un generador de duplicados.
Inserción idempotente en la base
En SQLite y en la mayoría de las bases relacionales existe una construcción que o bien ignora el conflicto por clave primaria, o bien actualiza la fila existente. La primera variante INSERT OR IGNORE ya la viste en el paso 2. Es adecuada cuando los registros son inmutables. La segunda variante es necesaria si la fuente puede actualizar registros y necesitas la versión más reciente:
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])Escritura idempotente en archivos
Si el resultado debe quedar en archivos y no en la base, aplica el mismo principio: una página equivale a un archivo con nombre determinista. El nombre depende de los parámetros de la página, no de la hora ni de un contador. Antes de la descarga compruebas si el archivo final existe; si existe, omites la página. Escribes en un nombre temporal y renombras al terminar, como en la función save_state.
def page_path(base_dir, job, cursor_or_page):
safe = str(cursor_or_page).replace('/', '_').replace(':', '_')[:120]
return os.path.join(base_dir, job, 'page-%s.jsonl' % safe)
def write_page_idempotent(path, rows):
if os.path.exists(path):
return False # страница уже есть, повторная запись не нужна
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + '.part'
with open(tmp, 'w', encoding='utf-8') as f:
for r in rows:
f.write(json.dumps(r, ensure_ascii=False))
f.write(chr(10))
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
return TrueLos archivos con extensión .part que quedaron tras una caída se pueden eliminar sin miedo al inicio: por definición están incompletos.
Consejo: para exportaciones por offset no te fíes solo de «el archivo existe, entonces la página está lista». Comprueba además que el número de líneas del archivo sea igual al tamaño de la página (excepto la última). Un archivo vacío o corto con un nombre existente es mejor volver a descargarlo.
Resultado esperado
Verificación: llama a la función de escritura de la misma página tres veces seguidas. Luego ejecuta SELECT count(*) FROM records;. El número debe ser igual al tamaño de una página, no al triple. Para la variante de archivos, en el directorio debe haber exactamente un archivo de página y ningún archivo .part.
Paso 4: Deduplicación de resultados sin inflar la memoria
Objetivo de la etapa: descartar registros repetidos en un flujo de millones de filas sin mantener todas las claves en la memoria RAM.
Hash del registro
Cuando un registro no tiene identificador, la clave pasa a ser el hash de su contenido. Para que registros idénticos den el mismo hash, el contenido debe canonizarse: ordenar las claves del diccionario, eliminar espacios sobrantes, fijar los separadores. De lo contrario, el mismo registro que llegue con otro orden de campos obtendrá otro hash.
import hashlib
import json
def record_key(rec, fields=None):
src = rec if fields is None else {k: rec.get(k) for k in fields}
canon = json.dumps(src, sort_keys=True, ensure_ascii=False, separators=(',', ':'))
return hashlib.blake2b(canon.encode('utf-8'), digest_size=16).digest()La función devuelve 16 bytes. Es suficiente: la probabilidad de coincidencia aleatoria para cientos de millones de registros es despreciable. El parámetro fields permite calcular el hash solo con campos estables, excluyendo, por ejemplo, la hora de la última actualización, que cambia en cada petición.
Por qué un conjunto en memoria no funciona con millones
Lo primero que se te ocurre: crear seen = set() y meter ahí las claves. Hagamos cuentas. Un objeto bytes de 16 bytes de longitud ocupa en Python unos 49 bytes más los propios datos, en total aproximadamente 65 bytes. Una ranura del conjunto, considerando el factor de ocupación, añade unos 30 bytes más. Obtenemos del orden de 95 bytes por clave. En 10 millones de registros eso es alrededor de 950 MB, en 50 millones casi 5 GB. Y lo principal: después de reiniciar el proceso el conjunto está vacío, y toda la deduplicación empieza desde cero.
Tres formas de no inflar la memoria
- Guardar las claves en la propia base. La vía más simple y fiable. Si la clave es la clave primaria de la tabla records, la deduplicación ya está hecha con la construcción INSERT OR IGNORE del paso 3. El índice vive en el disco, sobrevive a los reinicios, y SQLite cachea por sí mismo las páginas calientes del índice. Para 10 millones de claves de 16 bytes, el índice ocupará unos 400-500 MB en disco, pero no en memoria.
- Tabla separada de claves vistas sin rowid. Es necesaria si los propios datos los escribes no en SQLite, sino, por ejemplo, en archivos. Entonces SQLite se usa solo como un conjunto compacto en disco.
- Conjunto comprimido en memoria como prefiltro. Variante avanzada: truncar el hash a 8 bytes y guardarlo como entero en un arreglo ordenado o usar un filtro de Bloom. La memoria se reduce varias veces, pero aparece la probabilidad de un falso positivo. Por eso ese prefiltro se aplica solo para descartar rápidamente registros evidentemente nuevos, y la verificación final siempre se hace contra la base.
Implementación del conjunto en disco
class DiskSeen:
def __init__(self, con):
self.con = con
con.execute('CREATE TABLE IF NOT EXISTS seen(key BLOB PRIMARY KEY) WITHOUT ROWID')
def filter_new(self, rows):
keyed = [(record_key(r), r) for r in rows]
keys = [k for k, _ in keyed]
placeholders = ','.join('?' * len(keys))
known = {row[0] for row in self.con.execute(
'SELECT key FROM seen WHERE key IN (%s)' % placeholders, keys)}
fresh = [(k, r) for k, r in keyed if k not in known]
# дедупликация внутри самой страницы
unique = {}
for k, r in fresh:
unique.setdefault(k, r)
return unique
def remember(self, keys):
self.con.executemany('INSERT OR IGNORE INTO seen(key) VALUES (?)', [(k,) for k in keys])Llama a filter_new antes de escribir la página, y a remember dentro de la misma transacción que la escritura de datos y el checkpoint. Entonces, tras una caída, el conjunto de vistos, los datos y la marca de progreso siempre estarán coordinados entre sí.
Consejo: una consulta con IN sobre 500 valores se ejecuta por índice en milisegundos. No compruebes las claves de una en una en un bucle: es decenas de veces más lento por el coste de cada consulta.
Resultado esperado
Verificación: forma una página de prueba de 500 registros, donde 100 se repiten dos veces dentro de la página y otros 100 ya están en la tabla seen. La función filter_new debe devolver exactamente 300 registros. Después de reiniciar el proceso, esos mismos 500 registros deben dar cero nuevos.
Posibles problemas
- Los duplicados siguen pasando. Revisa la canonización: lo más probable es que en los registros haya un campo con la hora de la petición o un orden aleatorio de elementos en una lista. Exclúyelo con el parámetro fields u ordena las listas anidadas antes del hash.
- Inserción lenta tras varios millones de filas. El índice dejó de caber en caché. Aumenta la caché de SQLite con el comando
PRAGMA cache_size=-200000(esto es 200 MB) y asegúrate de que las inserciones van en lotes dentro de una transacción por página, y no de una en una.
Paso 5: Paralelismo sin pérdidas
Objetivo de la etapa: acelerar la exportación con varios trabajadores simultáneos a través del proxy de modo que la caída de cualquiera de ellos no pierda tareas ni rompa la base.
Cuándo se puede paralelizar y cuándo no
La paginación por cursor es por naturaleza secuencial: el siguiente cursor se conoce solo después de obtener la página anterior. No se puede paralelizar directamente. Pero casi siempre se puede dividir la exportación en shards independientes: por días, por categorías, por regiones, por los primeros caracteres del identificador. Cada shard se exporta secuencialmente con su propio checkpoint, y los shards van en paralelo. La paginación por offset y por clave con límites conocidos se paraleliza directamente: tareas del tipo «páginas de 1 a 100» o «identificadores de 0 a 100000».
Cola de tareas en disco
Una cola en memoria muere junto con el proceso. Por eso las tareas viven en una tabla con estados:
pending— espera ejecución;running— tomada por un trabajador;done— ejecutada y escrita;failed— intentos agotados, requiere atención humana.
Al iniciar, el cargador primero pasa todas las tareas de running de vuelta a pending: si están colgadas en ese estado, significa que el proceso anterior murió en medio del trabajo. Luego los trabajadores toman las pending.
Un único punto de escritura
SQLite admite muchos lectores simultáneos, pero solo un escritor. El patrón más simple y seguro: los trabajadores solo descargan y devuelven datos, y toda la escritura en la base la hace el hilo principal. Sin bloqueos en el código, sin «database is locked». Además activa el modo WAL para que la lectura del estado desde otro proceso no interfiera con la escritura.
Límite de concurrencia
Limita el número de trabajadores por dos cosas. Primero, las capacidades del proxy: si en el panel de Proxeon tienes varios canales, es razonable mantener uno o dos trabajadores por canal para que la rotación de IP en un canal no rompa las conexiones de todos los hilos a la vez. Segundo, la cortesía hacia la fuente: incluso sin límites formales, diez hilos paralelos sobre una API pequeña crearán una carga por la que recibirás rechazos. Empieza con tres o cuatro trabajadores y sube observando la proporción de errores.
Código del cargador paralelo
import json
import sqlite3
from concurrent.futures import ThreadPoolExecutor, as_completed
def init_tasks(con):
con.execute('PRAGMA journal_mode=WAL')
con.executescript('''
CREATE TABLE IF NOT EXISTS tasks(
task_id TEXT PRIMARY KEY,
params TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
''')
with con:
con.execute('UPDATE tasks SET status=? WHERE status=?', ('pending', 'running'))
def enqueue(con, tasks):
with con:
con.executemany('INSERT OR IGNORE INTO tasks(task_id, params) VALUES (?, ?)',
[(t['task_id'], json.dumps(t['params'])) for t in tasks])
def claim(con, limit):
rows = con.execute('SELECT task_id, params FROM tasks WHERE status=? LIMIT ?',
('pending', limit)).fetchall()
with con:
con.executemany('UPDATE tasks SET status=? WHERE task_id=?',
[('running', r[0]) for r in rows])
return [(r[0], json.loads(r[1])) for r in rows]
def run_parallel(con, fetch_fn, write_fn, workers=4, max_attempts=5):
init_tasks(con)
with ThreadPoolExecutor(max_workers=workers) as pool:
while True:
batch = claim(con, workers * 2)
if not batch:
break
futures = {pool.submit(fetch_fn, params): task_id for task_id, params in batch}
for fut in as_completed(futures):
task_id = futures[fut]
try:
rows = fut.result()
except Exception as e:
with con:
con.execute(
'UPDATE tasks SET attempts=attempts+1, last_error=?, '
'status=CASE WHEN attempts+1 >= ? THEN ? ELSE ? END WHERE task_id=?',
(str(e)[:500], max_attempts, 'failed', 'pending', task_id))
continue
with con: # данные и статус задания в одной транзакции
write_fn(con, rows)
con.execute('UPDATE tasks SET status=? WHERE task_id=?', ('done', task_id))
failed = con.execute('SELECT count(*) FROM tasks WHERE status=?', ('failed',)).fetchone()[0]
print('очередь пуста, заданий с ошибкой:', failed)La función fetch_fn realiza la petición a través del proxy y devuelve una lista de registros. Trabaja en un hilo y no toca la base. La función write_fn se llama en el hilo principal dentro de una transacción y hace la inserción idempotente del paso 3. La tarea fallida vuelve automáticamente a pending y será tomada de nuevo en el siguiente ciclo claim; tras agotar los intentos obtiene el estado failed, y tú te encargarás de ella manualmente.
Atención: no pases el objeto de conexión sqlite3 a los trabajadores. La conexión está vinculada al hilo en el que se creó, y usarla desde otro hilo provocará un error o, peor aún, una corrupción silenciosa de los datos. O cada hilo tiene su propia conexión, o, como en el ejemplo de arriba, ninguna.
Consejo: usa por cada trabajador un objeto requests.Session distinto con su propia dirección de proxy. Si en Proxeon tienes varios canales, distribúyelos entre los trabajadores por turnos: el trabajador 0 toma el canal 0, el trabajador 1 el canal 1, y así sucesivamente. Así la caída de conexión en un canal afectará solo a un hilo.
Resultado esperado
Verificación: pon en la cola 100 tareas, inicia cuatro trabajadores y mata el proceso a los medio minuto. Ejecuta SELECT status, count(*) FROM tasks GROUP BY status;. Verás varias done, algunas running y el resto pending. Vuelve a iniciar: las running deben desaparecer al arrancar, y al finalizar todas las tareas deben quedar en done, salvo las que fallaron honestamente y están en failed con el texto del error en last_error.
Paso 6: Reanudación tras una pausa larga
Objetivo de la etapa: continuar correctamente una exportación que se detuvo varias horas o días atrás, sin toparse con un estado caducado.
Qué caduca
Reanudar a los diez segundos y reanudar a la semana son tareas distintas. Tras una pausa larga, parte del estado guardado deja de ser válido.
- Sesiones y cookies. Las sesiones del servidor suelen durar de varias horas a un día. Después de eso, las cookies guardadas llevarán a respuestas 401 o a una redirección al formulario de inicio de sesión. Solución: al arrancar, realizar una autenticación completa de nuevo, y no restaurar las cookies desde un archivo.
- Tokens de acceso. Los tokens OAuth viven una hora, a veces menos. Si tienes un refresh token, actualiza el access token antes de arrancar y según lo previsto durante el trabajo, sin esperar al rechazo.
- Cursores de paginación. Muchas API limitan la vida útil del cursor a minutos u horas. Un cursor caducado devolverá un error 400 con un mensaje sobre cursor no válido. Precisamente por eso en el paso 2 guardamos un ancla de respaldo: el identificador del último registro. Si la fuente soporta el filtro por identificador o por fecha de modificación, construye una nueva petición desde ese ancla. Si no lo soporta, habrá que empezar el shard desde el principio, y la deduplicación del paso 4 descartará lo ya obtenido.
- El contenido del archivo en el servidor. Para la reanudación del paso 1 es crítico que el archivo no haya cambiado. Compara el ETag actual con el guardado antes de cada pasada; si no coinciden, empieza el archivo de nuevo.
- Configuración del proxy. En una semana en el panel de Proxeon pueden haber cambiado el puerto, la contraseña, o haberse agotado la vigencia del canal. Verifica el proxy con una petición de prueba antes de empezar a procesar la cola.
- El propio conjunto de datos. Si la exportación dura una semana y la fuente ha añadido y eliminado registros en ese tiempo, tu resultado será una mezcla de estados en distintos momentos. Para muchas tareas eso es aceptable. Si no lo es, guarda la hora de inicio en el checkpoint y, al terminar, haz una pasada incremental aparte por los registros modificados después de esa hora.
Comprobación previa al vuelo
Reúne todas las comprobaciones en una sola función que se ejecuta al arrancar, antes de cualquier trabajo real. O bien pone el estado en orden, o bien detiene el cargador con un mensaje comprensible.
def preflight(session, state, probe_url):
# 1. прокси жив и авторизован
r = session.head(probe_url, timeout=20)
if r.status_code == 407:
raise SystemExit('прокси отверг логин или пароль, проверьте данные в кабинете Proxeon')
# 2. токен доступа свежий
refresh_access_token(session)
# 3. курсор ещё действителен
if state['cursor']:
test = session.get(API_BASE + '/orders', params={'cursor': state['cursor'], 'limit': 1}, timeout=30)
if test.status_code == 400 and 'cursor' in test.text.lower():
print('курсор протух, переключаемся на якорь по last_id =', state['last_id'])
state['cursor'] = None
state['resume_after_id'] = state['last_id']
# 4. напоминание о возрасте выгрузки
print('выгрузка стартовала', state['started_at'], 'страниц записано', state['pages'])
return stateLa función refresh_access_token depende de tu fuente: normalmente es una petición POST con el refresh token, tras la cual actualizas el encabezado Authorization en la sesión. El campo resume_after_id se usa luego en la función de obtención de página como filtro «identificador mayor que el indicado».
Consejo: guarda el refresh token y la contraseña del proxy no en el checkpoint, sino en variables de entorno o en un archivo de secretos aparte con permisos restringidos. El checkpoint lo copiarás, lo enviarás a colegas y lo adjuntarás a informes de errores; los secretos ahí sobran.
Resultado esperado
Verificación: estropea manualmente el cursor en la tabla checkpoint con el comando UPDATE checkpoint SET cursor='broken'; y arranca el cargador. En la consola debe aparecer la línea sobre el cambio al ancla por last_id, y la exportación debe continuar sin caer. El número de registros tras finalizar debe coincidir con una ejecución de control sin estropear el cursor.
Paso 7: Esqueleto listo de un cargador resistente en Python
Objetivo de la etapa: reunir todo lo de los pasos anteriores en un solo archivo que se pueda ejecutar, interrumpir, volver a ejecutar y obtener el resultado completo sin duplicados.
Estructura del esqueleto
- Configuración desde variables de entorno: dirección del proxy de Proxeon, dirección de la API, token, nombre de la tarea.
- Clase Store: SQLite con tablas de registros y de checkpoint, una transacción por página.
- Función de clave de registro para la idempotencia.
- Función de obtención de página a través del proxy.
- Ciclo principal con reanudación por checkpoint y recreación de la sesión tras un corte.
Código completo
# resumable_loader.py
import hashlib
import json
import os
import sqlite3
import time
from datetime import datetime, timezone
import requests
PROXY_URL = os.environ['PROXY_URL']# http://USER:PASS@HOST:PORT из кабинета Proxeon
API_BASE = os.environ.get('API_BASE', 'https://api.example.com')
API_TOKEN = os.environ.get('API_TOKEN', '')
DB_PATH = os.environ.get('DB_PATH', 'export.sqlite')
JOB = os.environ.get('JOB', 'orders-2026')
PAGE_SIZE = 500
MAX_ATTEMPTS = 8
FIELDS = ('cursor', 'last_id', 'pages', 'rows', 'started_at')
def now():
return datetime.now(timezone.utc).isoformat()
def record_key(rec):
canon = json.dumps({'id': rec['id']}, sort_keys=True, separators=(',', ':'))
return hashlib.blake2b(canon.encode('utf-8'), digest_size=16).digest()
class Store:
def __init__(self, path):
self.con = sqlite3.connect(path)
self.con.execute('PRAGMA journal_mode=WAL')
self.con.executescript('''
CREATE TABLE IF NOT EXISTS records(
key BLOB PRIMARY KEY,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
) WITHOUT ROWID;
CREATE TABLE IF NOT EXISTS checkpoint(
job TEXT PRIMARY KEY,
cursor TEXT,
last_id TEXT,
pages INTEGER NOT NULL,
rows INTEGER NOT NULL,
started_at TEXT,
updated_at TEXT
);
''')
def load(self, job):
row = self.con.execute(
'SELECT cursor, last_id, pages, rows, started_at FROM checkpoint WHERE job=?',
(job,)).fetchone()
if row is None:
return {'cursor': None, 'last_id': None, 'pages': 0, 'rows': 0, 'started_at': now()}
return dict(zip(FIELDS, row))
def commit_page(self, job, rows, state):
ts = now()
with self.con:
self.con.executemany(
'INSERT OR IGNORE INTO records(key, payload, fetched_at) VALUES (?, ?, ?)',
[(record_key(r), json.dumps(r, ensure_ascii=False), ts) for r in rows])
self.con.execute(
'INSERT INTO checkpoint(job, cursor, last_id, pages, rows, started_at, updated_at) '
'VALUES (?, ?, ?, ?, ?, ?, ?) '
'ON CONFLICT(job) DO UPDATE SET cursor=excluded.cursor, last_id=excluded.last_id, '
'pages=excluded.pages, rows=excluded.rows, updated_at=excluded.updated_at',
(job, state['cursor'], state['last_id'], state['pages'], state['rows'],
state['started_at'], ts))
def unique_count(self):
return self.con.execute('SELECT count(*) FROM records').fetchone()[0]
def make_session():
s = requests.Session()
s.proxies = {'http': PROXY_URL, 'https': PROXY_URL}
s.headers['User-Agent'] = 'resumable-loader/1.0'
if API_TOKEN:
s.headers['Authorization'] = 'Bearer ' + API_TOKEN
return s
def fetch_page(session, cursor):
params = {'limit': PAGE_SIZE}
if cursor:
params['cursor'] = cursor
r = session.get(API_BASE + '/orders', params=params, timeout=(15, 90))
r.raise_for_status()
body = r.json()
return body['items'], body.get('next_cursor')
def run():
store = Store(DB_PATH)
state = store.load(JOB)
print('старт: страниц %d, строк %d, уникальных в базе %d'
% (state['pages'], state['rows'], store.unique_count()))
session = make_session()
attempts = 0
while True:
try:
items, next_cursor = fetch_page(session, state['cursor'])
attempts = 0
except (requests.ConnectionError, requests.Timeout, requests.HTTPError) as e:
attempts += 1
if attempts > MAX_ATTEMPTS:
print('попытки исчерпаны, состояние сохранено, запустите снова позже')
raise
print('обрыв (%s), попытка %d из %d' % (type(e).__name__, attempts, MAX_ATTEMPTS))
time.sleep(min(60, 2 ** attempts)) # выбор пауз описан в статье про 429 и ретраи
session = make_session() # новая сессия: соединение через прокси пересоздаётся
continue
if not items:
break
state['pages'] += 1
state['rows'] += len(items)
state['last_id'] = str(items[-1]['id'])
state['cursor'] = next_cursor
store.commit_page(JOB, items, state)
if state['pages'] % 20 == 0:
print('страниц %d, строк %d' % (state['pages'], state['rows']))
if next_cursor is None:
break
print('готово: страниц %d, строк получено %d, уникальных в базе %d'
% (state['pages'], state['rows'], store.unique_count()))
if __name__ == '__main__':
run()Cómo adaptarlo a tu fuente
- Sustituye la ruta
/ordersy los nombres de los campositems,next_cursor,idpor los que devuelva tu API. Son tres lugares en las funciones fetch_page y record_key. - Si tu fuente tiene paginación por offset, sustituye el parámetro cursor por page y calcula el siguiente valor como state['pages'] + 1. En el checkpoint guarda el número de página en lugar del cursor.
- Si la paginación es por clave, pasa un parámetro como
after_iddesde state['last_id'] y elimina el trabajo con el cursor. - Si necesitas paralelismo, mueve fetch_page a fetch_fn del paso 5, y usa commit_page como write_fn. Divide la exportación en shards y llena la cola de tareas.
- Añade la función preflight del paso 6 antes del ciclo principal.
Verificación del resultado: checklist
Antes de lanzar el cargador sobre un volumen real de varias horas, pásalo por esta lista. Cada punto lleva un par de minutos, y en conjunto garantizan que la exportación nocturna no perderá datos.
- Prueba de interrupción. Inicia el cargador, a los 30 segundos pulsa Ctrl+C. Vuelve a iniciarlo. La primera línea de salida debe mostrar un número de páginas distinto de cero, no «страниц 0».
- Prueba de duplicados. Reduce PAGE_SIZE a 10, interrumpe el cargador cinco veces seguidas en momentos aleatorios. Al finalizar, compara el número de registros únicos con el número de filas obtenidas: los únicos deben ser menor o igual, y con paginación por cursor limpia sin inserciones en la fuente, casi iguales.
- Prueba de consistencia. Tras cualquier interrupción, ejecuta dos consultas:
SELECT rows FROM checkpoint;ySELECT count(*) FROM records;. La diferencia entre ellas no debe superar un tamaño de página. Si lo supera, el checkpoint y los datos no se escriben en la misma transacción. - Prueba del proxy. Desactiva temporalmente la variable PROXY_URL o pon una contraseña incorrecta. El cargador debe caer en la primera petición con un error comprensible, y no colgarse ni empezar a ir directamente a la fuente.
- Prueba de cursor caducado. Estropea el cursor en la base, como se describe en el paso 6, y asegúrate de que se activa el cambio al ancla.
- Prueba de disco. Comprueba el tamaño del archivo export.sqlite tras mil páginas y multiplícalo por el número esperado de páginas. Asegúrate de que en el disco habrá espacio suficiente con un margen del 20 por ciento.
Verificación: se considera un resultado exitoso cuando, tras tres interrupciones intencionadas y tres reinicios, el número final de registros únicos coincide con el obtenido en una ejecución continua sobre la misma fuente, y en la consola no apareció ni una vez la línea sobre el inicio desde cero.
Funciones adicionales y optimización
- Progreso y estimación de tiempo. Si se conoce el número total de registros, muestra el porcentaje y la estimación del tiempo restante cada veinte páginas. Es útil tanto para ti como para distinguir un bloqueo de un funcionamiento lento.
- Compresión del payload. Para decenas de millones de registros, el JSON en texto ocupa mucho espacio. Comprime el campo payload con zlib.compress antes de escribir y guárdalo como BLOB. El ahorro suele ser de tres a ocho veces.
- Migración a una base de servidor. El esquema con el checkpoint en la misma transacción que los datos se traslada a PostgreSQL casi sin cambios. La construcción ON CONFLICT se soporta allí también, y la restricción de «un solo escritor» desaparece.
- Exportaciones incrementales. Guarda la hora de inicio de cada tarea y, tras la exportación completa, lanza una tarea aparte con el filtro «modificado después de». Así mantienes una copia actualizada sin recargarlo todo.
- Un proceso aparte por shard. En lugar de hilos, puedes lanzar varios ejemplares del script con distintos valores de JOB y distintos archivos DB_PATH, y unir los resultados al final. Es más sencillo de depurar y elimina por completo la cuestión de la escritura concurrente.
- Métricas de cortes. Registra cada corte con el tipo de excepción y la hora. En un día verás que los cortes se agrupan en torno a los intervalos de rotación de IP del canal de Proxeon, y podrás ajustar el intervalo de rotación a la duración de tus peticiones.
Errores típicos y soluciones
Abajo se reúnen situaciones con las que se topa casi todo el mundo en los primeros lanzamientos. Formato: problema, causa, solución.
- Problema: tras el reinicio, la exportación empieza cada vez desde cero. Causa: el checkpoint se escribe en memoria o en un archivo que no sobrevive a una caída, o el cargador no lo lee al arrancar. Solución: asegúrate de que la primera acción en la función run sea store.load, y que el estado se actualice tras cada página dentro de una transacción.
- Problema: en la base hay una vez y media más registros que en la fuente. Causa: la clave de deduplicación es inestable: incluye la hora de la petición, el número de página o un campo con orden aleatorio. Solución: calcula la clave solo por el identificador de la fuente o por una lista explícita de campos estables mediante el parámetro fields.
- Problema: en la base hay menos registros que en la fuente, aunque la exportación terminó sin errores. Causa: el checkpoint se actualizaba antes de escribir los datos, y tras una caída se omitió una página. O bien la paginación por offset, al eliminar registros en la fuente, desplazó las páginas hacia atrás. Solución: cambia el orden a «datos, luego checkpoint» en una misma transacción; para fuentes con eliminaciones pásate a la paginación por clave.
- Problema: la reanudación del archivo da un comprimido corrupto con el tamaño coincidente. Causa: el servidor respondió una vez 200 en lugar de 206, el archivo se sobrescribió parcialmente y luego se añadió el resto. Solución: guarda el ETag junto al archivo, elimina el archivo si cambia; verifica que el encabezado Content-Range corresponda al offset solicitado.
- Problema: error «database is locked» con trabajo paralelo. Causa: varios hilos escriben en SQLite a la vez, o la conexión se pasó entre hilos. Solución: un único punto de escritura en el hilo principal, los trabajadores solo descargan; modo WAL; la conexión se crea en el hilo donde se usa.
- Problema: tras una hora de trabajo, todas las peticiones empiezan a devolver 401. Causa: ha caducado el token de acceso. Solución: actualiza el token según lo previsto antes de que caduque, y al recibir 401 llama a la actualización y repite la petición una vez, sin considerarlo un corte.
- Problema: la memoria RAM crece hasta varios gigabytes. Causa: el conjunto de claves vistas o la lista de todos los registros se mantiene en la memoria del proceso. Solución: deduplicación mediante clave primaria en la base o tabla seen; escritura de datos por páginas, sin acumulación.
- Problema: los cortes ocurren estrictamente cada pocos minutos. Causa: coinciden con el intervalo de rotación de IP del canal del proxy. Solución: es una situación normal, el cargador debe sobrevivirla. Si las peticiones son largas, ajusta el intervalo de rotación en el panel de Proxeon para que sea notablemente mayor que el tiempo típico de una petición, o usa la rotación bajo demanda entre páginas.
FAQ: preguntas frecuentes sobre exportación resistente
¿Es obligatorio usar SQLite si el resultado se necesita en CSV?
No, pero es cómodo. SQLite aquí cumple el papel de almacén fiable del estado y del conjunto de claves vistas. El CSV final lo exportas con un solo comando desde la tabla records al terminar. Si quieres prescindir por completo de la base, usa la variante de archivos del paso 3 con un archivo por página y un checkpoint JSON con reemplazo atómico.
¿Con qué frecuencia guardar el checkpoint: tras cada página o menos?
Tras cada página. Una transacción de SQLite con varios cientos de filas se ejecuta en milisegundos, es insignificante comparado con una petición de red a través del proxy. El ahorro en checkpoints poco frecuentes no vale el riesgo de perder decenas de páginas.
¿Qué hacer si la API no devuelve ni cursor ni identificadores, solo números de página?
Trabaja por número de página, guárdalo en el checkpoint e incluye obligatoriamente la deduplicación por hash del contenido del registro. Acepta que, con cambios activos en la fuente, parte de los registros puede omitirse por el desplazamiento de páginas. Para datos críticos, haz una segunda pasada en orden inverso de páginas: lo omitido en la primera pasada con alta probabilidad aparecerá en la segunda.
¿Se puede reanudar la descarga de un archivo con varios hilos en distintos rangos?
Se puede, si el servidor soporta Range. Divide el archivo en fragmentos de 50-100 MB, cada fragmento es una tarea de la cola del paso 5 con su propio archivo temporal, y al terminar todos los fragmentos únelos en el orden correcto. Verifica cada fragmento por tamaño, y todo el archivo por suma de control, si la hay.
¿Cuántos trabajadores poner al trabajar a través de un proxy?
Empieza con tres o cuatro por canal de Proxeon y observa la proporción de errores en la tabla tasks. Si los errores son menos del uno por ciento, añade dos más. Si los errores crecen, redúcelos. Más de diez hilos por canal rara vez dan ventaja: chocas con el ancho de banda del canal o con la paciencia de la fuente.
¿Hay que guardar las cookies de sesión entre ejecuciones?
Normalmente no. La reautenticación al arrancar lleva segundos y es más fiable que restaurar cookies con una vida útil desconocida. Excepción: la fuente limita el número de inicios de sesión al día. Entonces guarda las cookies, pero ante el primer 401 o redirección al inicio de sesión, descártalas y autentícate de nuevo.
¿Cómo saber que la exportación terminó por completo y no se cortó silenciosamente?
Para archivos: el tamaño es igual a Content-Length y la suma de control coincide. Para API: se recibió una página sin next_cursor o una página vacía, y además el número de filas corresponde al total, si la fuente lo comunica. Guarda en el checkpoint un indicador explícito de finalización para que una nueva ejecución no empiece un recorrido nuevo.
¿Qué hacer con las tareas en estado failed?
Mira el campo last_error. Si son errores de red, simplemente devuelve las tareas a pending con el comando UPDATE y arranca el cargador de nuevo. Si son errores de análisis de datos, significa que en la fuente hay registros de forma no estándar: corrige el código y reinicia. Nunca elimines las failed en silencio, es la única evidencia de lo que falta en la exportación.
¿Se puede usar este enfoque no con requests, sino con una librería asíncrona?
Sí, los principios son los mismos: unidad de trabajo atómica, checkpoint junto con los datos, escritura idempotente, cola en disco. Solo cambia el transporte. El único matiz: deja la escritura en SQLite síncrona y secuencial, y mantén el paralelismo a nivel de peticiones de red.
Conclusión
Has recorrido el camino desde el archivo incompleto y el reinicio nervioso hasta un cargador al que le da igual un corte. Fijemos qué se ha hecho exactamente.
- Aclaramos la reanudación por HTTP: verificación de Accept-Ranges con HEAD, encabezado Range con el tamaño actual del archivo, distinción de los códigos 206 y 200, protección contra la sustitución del archivo mediante ETag e If-Range.
- Construimos checkpoints para exportaciones paginadas: cursor, número de página o identificador del último registro más contadores auxiliares, todo en una misma transacción con los datos.
- Hicimos la escritura idempotente mediante clave primaria y la construcción INSERT OR IGNORE o ON CONFLICT DO UPDATE, y para archivos mediante nombres deterministas y reemplazo atómico.
- Organizamos la deduplicación en disco, para que millones de claves no vivan en la memoria RAM y sobrevivan a los reinicios.
- Añadimos paralelismo con cola de tareas en SQLite, devolución automática de las tareas fallidas y un único punto de escritura.
- Previmos la reanudación tras una pausa larga: actualización de tokens, reautenticación, cambio de un cursor caducado a un ancla por identificador, verificación del proxy de Proxeon antes de arrancar.
- Reunimos todo en un esqueleto funcional que se adapta a una fuente concreta sustituyendo tres o cuatro líneas.
Qué hacer después
Toma el esqueleto del paso 7 y ejecútalo sobre un volumen real pequeño, digamos diez mil registros. Pasa la checklist de la sección de verificación. Solo después lanza la exportación completa de noche. Por la mañana verás o bien la línea «готоvo» con contadores coincidentes, o la línea sobre intentos agotados con el estado guardado, y entonces simplemente vuelve a ejecutar el script.
Hacia dónde crecer
El siguiente nivel son las exportaciones incrementales por hora de modificación en lugar de recorridos completos, el traslado del estado a una base de servidor para varias máquinas, y una estrategia sensata de reintentos teniendo en cuenta los códigos de respuesta, a la que está dedicado un artículo aparte sobre el 429 y los reintentos. La combinación de reintentos bien hechos de ese artículo y el estado resistente de este da un cargador que puedes dejar funcionando una semana sin abrir la terminal.
Y lo último. El corte de una exportación larga no es una avería, sino una situación de trabajo que ahora sabes manejar. Buenas exportaciones.