Как выгружать большие объёмы данных через прокси и не начинать заново после обрыва
Содержание статьи
- Введение: почему длинная выгрузка почти всегда обрывается и это нормально
- Предварительная подготовка: инструменты, доступы и окружение
- Базовые понятия: словарь устойчивой выгрузки простыми словами
- Шаг 1: докачка по http через заголовок range
- Шаг 2: чекпоинты для постраничных выгрузок
- Шаг 3: идемпотентность, чтобы повтор не создавал дубли
- Шаг 4: дедупликация результатов без раздувания памяти
- Шаг 5: параллелизм без потерь
- Шаг 6: возобновление после долгого перерыва
- Шаг 7: готовый скелет устойчивого загрузчика на python
- Проверка результата: чек-лист устойчивой выгрузки
- Типичные ошибки и их решения
- Дополнительные возможности и оптимизация
- Faq: частые вопросы по устойчивой выгрузке
- Заключение: что вы теперь умеете и куда двигаться
Введение: почему длинная выгрузка почти всегда обрывается и это нормально
Если вы хоть раз запускали большую выгрузку данных, вы знаете это чувство. Процесс шёл несколько часов, дошёл до девяноста процентов и оборвался. Соединение упало, сервер вернул ошибку, ноутбук ушёл в сон. И всё приходится начинать заново. Этот гайд написан для того, чтобы такое больше не повторялось.
Что вы получите в итоге. Вы научитесь строить загрузчик, который переживает обрывы. Он докачивает файлы с середины, помнит, на какой странице остановился, не создаёт дубликаты при повторе и умеет возобновляться даже после долгого перерыва. Вы получите готовый скелет на Python, который можно адаптировать под свою задачу.
Для кого этот гайд. Для инженеров, аналитиков и разработчиков, которые выгружают данные из API, скачивают большие файлы или собирают постраничные результаты через прокси. Уровень средний. Вы должны понимать основы HTTP и уметь читать код на Python. Глубоких знаний сетевого программирования не требуется.
Что нужно знать заранее. Базовый Python, понятие HTTP-запроса и ответа, что такое заголовки и статус-коды. Если вы работали с библиотекой requests, этого достаточно.
Сколько времени потребуется. На чтение и понимание концепций уйдёт около сорока минут. На сборку рабочего загрузчика по нашему скелету от одного до трёх часов, в зависимости от вашего источника данных.
Важное уточнение по теме. Мы не будем разбирать статус-коды вроде 429 и стратегии повторов с задержкой (backoff). Об этом есть отдельный материал. Здесь фокус только на одном: состояние процесса и его возобновление. Как сохранять прогресс, как не потерять и не задублировать данные, как продолжить с того места, где остановились.
Совет: Держите под рукой блокнот или отдельный файл, куда будете выписывать параметры вашего источника: поддерживает ли он докачку, есть ли постраничная навигация, какой у него формат курсора. Эти заметки пригодятся на каждом шаге.
Предварительная подготовка: инструменты, доступы и окружение
Перед тем как писать код, соберём рабочее окружение. Это займёт десять минут, но сэкономит часы отладки.
Что нужно установить
- Установите Python версии 3.10 или новее. Проверьте версию командой
python --versionв терминале. - Создайте виртуальное окружение командой
python -m venv venv, чтобы зависимости проекта не смешивались с системными. - Активируйте окружение. На Windows это
venv\Scripts\activate, на macOS и Linux этоsource venv/bin/activate. - Установите библиотеку для HTTP-запросов командой
pip install requests. - Для более быстрой работы с базой состояния установите ничего дополнительно не нужно: модуль
sqlite3уже входит в стандартную библиотеку Python.
Что нужно из доступов
- Доступ к вашему источнику данных: URL, токен или ключ API, если он требуется.
- Прокси от Proxeon с адресом, портом и данными для авторизации. Без стабильного прокси устойчивая выгрузка теряет смысл, потому что именно прокси распределяет нагрузку и делает соединения предсказуемыми.
- Место на диске для файла состояния и для самих выгружаемых данных.
Проверка прокси Proxeon
- Возьмите строку подключения вида
http://логин:пароль@адрес:порт. - Проверьте её простым запросом. В терминале выполните
curl -x http://логин:пароль@адрес:порт https://api.ipify.orgи убедитесь, что вернулся IP-адрес прокси, а не ваш собственный.
Совет: Сохраните строку подключения к прокси в переменную окружения, а не в код. Так вы случайно не отправите пароль в систему контроля версий. В коде читайте её через os.environ.
⚠️ Внимание: Всегда работайте только с теми источниками данных, доступ к которым у вас есть на законных основаниях. Соблюдайте условия использования сервиса и лимиты, указанные его владельцем. Прокси Proxeon предназначен для легальной инженерной работы: распределения нагрузки, стабильности соединений и корректной выгрузки данных.
✅ Проверка: Окружение готово, если команда python -c "import requests, sqlite3" отработала без ошибок, а запрос через прокси вернул адрес прокси-сервера.
Базовые понятия: словарь устойчивой выгрузки простыми словами
Прежде чем писать код, разберём ключевые термины. Без них дальнейшие шаги будут звучать как заклинания.
Докачка
Докачка это продолжение загрузки файла с того байта, на котором она прервалась. Вместо того чтобы качать файл заново, вы просите сервер отдать только недостающий кусок. Работает через HTTP-заголовок Range.
Чекпоинт
Чекпоинт это сохранённая точка прогресса. Представьте сохранение в компьютерной игре. Если что-то пошло не так, вы возвращаетесь к последнему сохранению, а не к началу игры. В выгрузке чекпоинт хранит, на какой странице или записи вы остановились.
Курсор
Курсор это метка, которую API даёт вам, чтобы вы могли попросить следующую порцию данных. Часто это строка вроде eyJvZmZzZXQiOjEwMH0. Вы отправляете её обратно, и сервер понимает, откуда продолжить.
Идемпотентность
Идемпотентность это свойство операции, при котором её повтор не меняет результат. Если вы дважды записали одну и ту же строку с одним ключом, в итоге строка одна, а не две. Это защита от дубликатов при повторах.
Дедупликация
Дедупликация это отсеивание повторяющихся записей. Даже при аккуратной работе один и тот же объект может прийти дважды. Дедупликация гарантирует, что в вашем итоговом наборе он останется один.
Основной принцип
Устойчивый загрузчик строится на одной идее: прогресс нужно сохранять постоянно, а не только в конце. Любой шаг может стать последним перед обрывом. Значит, после каждого успешного куска работы состояние должно быть записано на диск. Тогда возобновление это просто чтение состояния и продолжение.
Совет: Запомните правило трёх вопросов для любой выгрузки. Первый: где я остановился? Второй: как не задублировать уже полученное? Третий: что протухнет, пока меня не было? Ответы на них и составляют устойчивость.
Шаг 1: Докачка по HTTP через заголовок Range
Цель этапа. Научиться скачивать большой файл так, чтобы после обрыва продолжить с недокачанного байта, а не с нуля.
Как это работает
HTTP позволяет запросить не весь файл, а его часть. Для этого в запрос добавляется заголовок Range. Например, Range: bytes=1048576- означает: отдай мне всё, начиная с байта номер 1048576. Но сначала нужно убедиться, что сервер это умеет.
- Отправьте на файл запрос методом HEAD или обычный GET и посмотрите заголовки ответа.
- Найдите заголовок
Accept-Ranges. Если там значениеbytes, сервер поддерживает докачку. - Если заголовка нет или там
none, докачка невозможна. В этом случае придётся качать файл целиком за один заход или искать альтернативный источник.
Проверка поддержки докачки
Вот код, который проверяет, умеет ли сервер отдавать части файла.
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", totalДокачка файла с середины
Теперь главный код. Он смотрит, сколько байт уже скачано локально, и просит у сервера только остаток.
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)Разберём важные моменты. Если сервер вернул статус 206, он честно отдал часть файла и дозапись пройдёт корректно. Если сервер вернул 200 несмотря на заголовок Range, значит он проигнорировал докачку и отдаёт файл целиком. В этом случае мы переключаемся в режим полной перезаписи, чтобы не склеить старый кусок с новым и не испортить файл.
⚠️ Внимание: Никогда не дописывайте данные в режиме ab, если не уверены, что сервер ответил кодом 206. Иначе вы получите битый файл, где начало это остаток от прошлой попытки, а продолжение это новый полный файл. Такой файл откроется с ошибкой, а вы потеряете время на поиск причины.
Совет: Скачивайте не сразу в целевой файл, а во временный с расширением .part. Когда загрузка полностью завершится, переименуйте его в финальное имя. Так вы никогда не спутаете готовый файл с недокачанным.
Контроль целостности
После полной загрузки хорошо бы проверить, что файл не побился. Если сервер отдавал заголовок Content-Length, сравните его с реальным размером файла на диске. Если размеры совпали, файл дошёл целиком.
def verify_size(dest, expected):
if expected is None:
return True
return os.path.getsize(dest) == int(expected)✅ Проверка: Прервите загрузку на середине, закрыв программу. Запустите её снова. В логах вы должны увидеть, что запрос ушёл с заголовком Range, а файл дописался, а не начался заново. Итоговый размер совпадает с ожидаемым.
Шаг 2: Чекпоинты для постраничных выгрузок
Цель этапа. Настроить сохранение прогресса для API, которые отдают данные страницами, чтобы после обрыва продолжить с нужной страницы.
Что именно сохранять
Файл докачивается по байтам, а постраничная выгрузка по совсем другой логике. Здесь нет байтов, есть страницы и записи. Значит, в чекпоинте нужно хранить другое.
- Курсор если API работает на курсорах. Это самый надёжный вариант, потому что курсор сам знает, откуда продолжить.
- Номер страницы или смещение если API работает на offset и limit. Храните номер последней успешно обработанной страницы.
- Идентификатор последней записи если можно сортировать по возрастающему ID или дате. Тогда следующий запрос просит записи с ID больше сохранённого.
- Счётчик обработанных записей для контроля и отчётности.
Где хранить состояние
У вас есть три основных варианта, от простого к надёжному.
- JSON-файл. Проще всего. Записываете словарь с курсором и счётчиком в файл после каждой страницы. Подходит для одиночных, не параллельных выгрузок.
- SQLite-база. Надёжнее. Даёт транзакции, поэтому состояние не побьётся при обрыве в момент записи. Хорош, когда данных много и нужна дедупликация.
- Внешняя база данных. Для больших распределённых выгрузок, когда несколько процессов делят одну работу.
Сохранение чекпоинта в 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)Обратите внимание на приём с временным файлом. Мы пишем в файл с суффиксом .tmp, а потом атомарно переименовываем его через os.replace. Это защищает от ситуации, когда программа упала прямо во время записи чекпоинта. Старый чекпоинт при этом остаётся целым, а не превращается в половину JSON, который невозможно прочитать.
⚠️ Внимание: Никогда не записывайте чекпоинт напрямую в тот же файл поверх старого без временного файла. Обрыв в середине записи оставит вам испорченный чекпоинт, и возобновление станет невозможным. Атомарная замена решает эту проблему полностью.
Совет: Сохраняйте чекпоинт только после того, как данные страницы реально записаны в хранилище. Порядок такой: получили страницу, записали данные, потом обновили чекпоинт. Если поменять порядок местами, при обрыве вы пропустите страницу и потеряете данные.
Основной цикл с чекпоинтами
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✅ Проверка: Запустите выгрузку, дайте ей обработать несколько страниц, прервите. Откройте файл чекпоинта и убедитесь, что в нём записаны актуальный курсор и счётчик. Запустите снова: выгрузка должна продолжиться с сохранённого курсора, а не с первой страницы.
Шаг 3: Идемпотентность, чтобы повтор не создавал дубли
Цель этапа. Сделать так, чтобы повторный запуск или повтор конкретного запроса не приводил к появлению одинаковых записей в вашем хранилище.
Почему возникают дубли
Представьте: вы получили страницу данных, записали её в файл, но программа упала до того, как обновился чекпоинт. При следующем запуске вы попросите ту же страницу снова. Данные придут повторно и запишутся вторично. Так рождаются дубликаты. Это неизбежное следствие обрывов, и бороться с ним нужно на уровне архитектуры.
Ключ дедупликации
Главный инструмент идемпотентности это ключ дедупликации. Это поле или комбинация полей, которые однозначно определяют запись. Правильный выбор ключа решает половину проблемы.
- Естественный ID. Если у записи есть уникальный идентификатор от источника, используйте его. Это идеальный ключ.
- Комбинация полей. Если единого ID нет, соберите ключ из нескольких стабильных полей. Например, email плюс дата регистрации.
- Хеш содержимого. Если стабильных полей нет вообще, посчитайте хеш от всей записи. Это крайний вариант, потому что любое изменение поля создаст новый ключ.
Запись без дублей через UPSERT
Если вы храните результат в SQLite или другой базе, используйте вставку с игнорированием конфликта. Тогда повторная запись с тем же ключом просто ничего не сделает.
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()Ключевая деталь здесь это PRIMARY KEY на поле dedup_key. База сама отклонит повторную вставку с тем же ключом, потому что INSERT OR IGNORE проглотит конфликт молча. Вам не нужно вручную проверять, есть ли уже такая запись. База делает это за вас и делает быстро.
Совет: Выбирайте ключ дедупликации один раз в начале проекта и фиксируйте его в документации. Смена ключа посреди выгрузки означает, что старые и новые записи перестанут сопоставляться, и дубли всё же появятся. Стабильность ключа важнее его красоты.
✅ Проверка: Запустите выгрузку дважды подряд на одном и том же диапазоне данных. Посчитайте количество строк в базе командой SELECT COUNT(*) FROM records. Число должно быть одинаковым после первого и после второго запуска.
Шаг 4: Дедупликация результатов без раздувания памяти
Цель этапа. Отсеивать повторяющиеся записи на миллионах строк, не загружая всю память компьютера множеством уже виденных ключей.
Наивный подход и его проблема
Простейшая дедупликация: держать в памяти множество set из всех виденных ключей. Для каждой новой записи проверять, есть ли ключ в множестве. Работает отлично на сотнях тысяч строк. Но на миллионах и десятках миллионов множество разрастается и съедает гигабайты оперативной памяти. Программа замедляется или падает.
Решение первое: полагаться на базу
Самый простой и надёжный способ на больших объёмах это не хранить виденное в памяти вообще, а доверить проверку базе через PRIMARY KEY, как мы сделали в прошлом шаге. База хранит индекс на диске, а не в памяти вашего процесса. Она справится с десятками миллионов ключей без нагрузки на вашу оперативку.
Решение второе: хеш записи
Когда естественного ключа нет, считайте компактный хеш от записи. Хеш занимает фиксированные и небольшой объём независимо от размера самой записи.
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()Параметр sort_keys=True здесь критичен. Он гарантирует, что одинаковые по содержанию записи дадут одинаковый хеш, даже если поля в них шли в разном порядке. Без этой сортировки два идентичных объекта могут получить разные хеши и проскочить как разные записи.
Решение третье: фильтр Блума для экономии памяти
Если вам всё же нужна быстрая проверка в памяти на огромных объёмах, применяют фильтр Блума. Это структура, которая занимает мало места и быстро отвечает, видели мы ключ или точно нет. У неё есть особенность: она может изредка ошибочно сказать, что ключ уже был, хотя его не было. Поэтому фильтр Блума используют как быстрый предварительный отсев, а окончательную проверку оставляют базе.
- Проверяем ключ фильтром Блума.
- Если фильтр говорит, что точно не видели, сразу пишем в базу.
- Если фильтр говорит, что возможно видели, делаем точную проверку в базе.
⚠️ Внимание: Не пытайтесь дедуплицировать десятки миллионов строк обычным множеством в памяти. На типичном ноутбуке это приведёт к исчерпанию памяти и краху процесса на середине выгрузки. Переносите нагрузку на диск через базу или используйте фильтр Блума.
Совет: Если выгружаете данные порциями и внутри одной порции возможны дубли, дедуплицируйте порцию в памяти обычным множеством перед записью в базу. Порция небольшая, память не пострадает, а на базу пойдёт меньше лишних вставок.
✅ Проверка: Запустите дедупликацию на большом тестовом наборе с намеренными повторами. Проверьте, что итоговое число уникальных записей верное, а потребление памяти процессом остаётся стабильным и не растёт линейно с числом строк.
Шаг 5: Параллелизм без потерь
Цель этапа. Ускорить выгрузку за счёт параллельных запросов, не потеряв при этом ни одного задания и корректно повторив упавшие.
Очередь заданий
Основа безопасного параллелизма это очередь заданий. Вы заранее разбиваете работу на независимые кусочки. Например, список страниц или диапазонов. Складываете их в очередь. Несколько воркеров берут задания из очереди, выполняют и складывают результат. Если воркер упал, его задание можно вернуть в очередь и отдать другому.
Ограничение одновременности
Нельзя запускать бесконечное число параллельных запросов. Это перегрузит источник и ваш прокси. Правильный подход это ограничить число одновременных воркеров разумным значением. Начните с небольшого количества и повышайте, наблюдая за стабильностью.
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, failedПовтор упавших заданий
Собранный список failed это не потерянные данные, а список того, что нужно повторить. После первого прохода вы прогоняете упавшие задания ещё раз. Обычно этого хватает, чтобы добить остаток.
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 remainingРоль прокси Proxeon в параллелизме. При параллельной работе прокси распределяет соединения, что делает выгрузку стабильнее и предсказуемее. Каждый воркер работает через своё соединение, а нагрузка не концентрируется в одной точке.
⚠️ Внимание: При параллельной записи в один файл или в один чекпоинт возникают гонки данных. Два воркера могут перезаписать состояние друг друга. Пишите результаты только в базу с транзакциями или используйте отдельный файл на каждого воркера, а сводный чекпоинт собирайте отдельным потоком.
Совет: Делайте задания мелкими и независимыми. Если одно задание охватывает слишком большой диапазон, его обрыв отбросит много работы. Мелкие задания повторяются дёшево и почти незаметно.
✅ Проверка: Запустите параллельную выгрузку, искусственно уроните часть воркеров. После повторных раундов список remaining должен стать пустым, а итоговый набор данных полным. Сравните число полученных записей с ожидаемым.
Шаг 6: Возобновление после долгого перерыва
Цель этапа. Корректно продолжить выгрузку, если между попытками прошло много времени, и понять, что могло протухнуть за этот период.
Что протухает со временем
Обрыв на минуту и пауза на сутки это разные ситуации. За долгий перерыв часть вашего состояния может стать недействительной.
- Сессия. Многие сервисы держат сессию ограниченное время. После долгой паузы сервер её забудет, и запросы начнут возвращать ошибку авторизации.
- Токен доступа. Токены API часто имеют срок жизни в минуты или часы. Просроченный токен нужно обновить перед продолжением.
- Курсор. Некоторые курсоры живут недолго. Если курсор протух, придётся начать с ближайшей стабильной точки, например по идентификатору последней записи.
- Сами данные. За время паузы в источнике могли появиться новые записи или измениться старые. Это влияет на смещения при постраничной навигации по offset.
Стратегия безопасного возобновления
- При старте проверьте возраст чекпоинта. Если он старый, будьте готовы к тому, что часть состояния устарела.
- Обновите токен доступа и создайте новую сессию перед первым запросом. Не полагайтесь на старые.
- Отдавайте предпочтение возобновлению по идентификатору последней записи, а не по номеру страницы. ID стабилен, а номер страницы сдвигается, если данные изменились.
- Сделайте пробный запрос с сохранённым курсором. Если он вернул ошибку недействительного курсора, переключитесь на возобновление по 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)Почему возобновление по ID надёжнее. Представьте, что вы остановились на странице 50 при сортировке по дате. Пока вы отсутствовали, добавились новые записи в начало. Теперь страница 50 содержит совсем другие данные, а часть записей вы пропустите. Возобновление по идентификатору последней записи от этого не страдает: вы просто просите всё, что больше сохранённого ID.
Совет: Всегда сохраняйте в чекпоинт и курсор, и идентификатор последней записи одновременно. Курсор быстрее, но ID это ваш страховочный трос на случай, если курсор протухнет за время долгой паузы.
✅ Проверка: Остановите выгрузку, подождите достаточно долго, чтобы токен или курсор устарели, и запустите снова. Загрузчик должен обновить токен, обнаружить протухший курсор и продолжить по идентификатору без потери и без дублирования записей.
Шаг 7: Готовый скелет устойчивого загрузчика на Python
Цель этапа. Собрать всё изученное в единый работающий каркас, который вы адаптируете под свой источник данных.
Ниже собран каркас, объединяющий чекпоинты, дедупликацию через базу, обновление токена и возобновление. Функции получения страницы вы подставляете свои, под конкретный API.
import os
import json
import sqlite3
import requests
class ResilientLoader:
def __init__(self, cp_path, db_path, proxies):
self.cp_path = cp_path
self.proxies = proxies
self.conn = sqlite3.connect(db_path)
self.conn.execute(
"CREATE TABLE IF NOT EXISTS records ("
"dedup_key TEXT PRIMARY KEY, payload TEXT)"
)
self.conn.commit()
def load_cp(self):
if not os.path.exists(self.cp_path):
return {"cursor": None, "last_id": None, "count": 0}
with open(self.cp_path) as f:
return json.load(f)
def save_cp(self, cp):
tmp = self.cp_path + ".tmp"
with open(tmp, "w") as f:
json.dump(cp, f)
os.replace(tmp, self.cp_path)
def save_records(self, items):
rows = [(str(i["id"]), json.dumps(i)) for i in items]
self.conn.executemany(
"INSERT OR IGNORE INTO records "
"(dedup_key, payload) VALUES (?, ?)", rows
)
self.conn.commit()
def run(self, fetch_page):
cp = self.load_cp()
while True:
items, next_cursor = fetch_page(
cp["cursor"], cp["last_id"], self.proxies
)
if not items:
break
self.save_records(items)
cp["count"] += len(items)
cp["last_id"] = items[-1]["id"]
cp["cursor"] = next_cursor
self.save_cp(cp)
if next_cursor is None:
break
return cp["count"]Пример функции получения страницы под ваш источник. Здесь вы реализуете логику запроса и разбор ответа.
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")Запуск всего механизма выглядит просто.
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)Совет: Добавьте в цикл логирование каждой сотни записей: время, счётчик, текущий курсор. Так вы будете видеть прогресс и легко поймёте, если выгрузка застряла на одном месте.
✅ Проверка: Запустите каркас на реальном источнике, прервите на середине, запустите снова. Итоговое число записей после докачки совпадёт с полным числом записей источника, а повторный запуск не увеличит счётчик уникальных строк.
Проверка результата: чек-лист устойчивой выгрузки
Пройдитесь по этому списку. Если все пункты выполняются, ваш загрузчик действительно устойчив.
- Файловая докачка продолжается с недокачанного байта, а не с нуля.
- Загрузчик корректно обрабатывает случай, когда сервер игнорирует заголовок Range.
- Чекпоинт сохраняется после каждой обработанной страницы, а не только в конце.
- Чекпоинт пишется атомарно через временный файл и переименование.
- Повторный запуск не создаёт дубликаты в итоговом наборе.
- Дедупликация не растёт по памяти линейно с числом строк.
- Параллельные воркеры не теряют упавшие задания и повторяют их.
- После долгой паузы токен обновляется, а протухший курсор заменяется возобновлением по ID.
Как протестировать
- Запустите полную выгрузку небольшого набора и запомните число записей.
- Запустите ещё раз на том же наборе и убедитесь, что число не изменилось.
- Прервите выгрузку в разных местах: в начале, в середине, ближе к концу.
- После каждого прерывания перезапустите и проверьте, что итог одинаковый.
Показатели успеха. Число уникальных записей стабильно между запусками. Потребление памяти не растёт бесконтрольно. Возобновление всегда продолжает с сохранённой точки. Битых файлов после докачки нет.
Типичные ошибки и их решения
Проблема: файл после докачки не открывается. Причина: данные дописывались в режиме ab, хотя сервер ответил кодом 200 и отдал файл целиком. Решение: проверяйте статус ответа, при 200 переключайтесь на полную перезапись файла с нуля.
Проблема: после обрыва выгрузка начинается с первой страницы. Причина: чекпоинт сохранялся только в конце работы или не сохранялся вовсе. Решение: сохраняйте чекпоинт после каждой обработанной страницы, сразу за записью данных.
Проблема: чекпоинт не читается, JSON битый. Причина: программа упала во время записи прямо в целевой файл. Решение: пишите во временный файл и заменяйте атомарно через os.replace.
Проблема: в итоговом наборе дубликаты. Причина: нет ключа дедупликации или он нестабилен. Решение: задайте PRIMARY KEY на надёжном ключе и используйте INSERT OR IGNORE.
Проблема: процесс падает по нехватке памяти на больших объёмах. Причина: все виденные ключи держатся в множестве в памяти. Решение: перенесите проверку уникальности в базу или примените фильтр Блума.
Проблема: после долгой паузы запросы возвращают ошибку авторизации. Причина: токен или сессия протухли за время перерыва. Решение: обновляйте токен и создавайте новую сессию при каждом старте.
Проблема: после паузы часть записей пропущена или задвоена. Причина: возобновление шло по номеру страницы, а данные в источнике изменились. Решение: возобновляйтесь по идентификатору последней записи, а не по offset.
Проблема: параллельные воркеры теряют часть данных. Причина: несколько воркеров писали в один чекпоинт и перезаписывали друг друга. Решение: пишите результаты в базу с транзакциями, а не в общий файл состояния.
Дополнительные возможности и оптимизация
Пакетная запись
Не записывайте в базу по одной строке. Собирайте пакет из нескольких сотен записей и вставляйте разом через executemany. Это в разы ускоряет запись на больших объёмах.
Периодическая фиксация
Вызывайте commit не на каждой записи, а раз в несколько сотен строк. Слишком частый commit замедляет базу. Слишком редкий рискует потерять больше данных при обрыве. Найдите баланс под вашу нагрузку.
Отчёт о прогрессе
Добавьте оценку оставшегося времени. Зная скорость обработки страниц и общее число записей, вы прикинете, сколько ещё ждать. Это удобно для долгих выгрузок.
Раздельные хранилища сырых и обработанных данных
Храните сырые ответы отдельно от разобранных записей. Если позже вы измените логику разбора, вам не придётся заново скачивать данные. Достаточно перегнать сырые ответы через новый парсер.
Совет: Настройте прокси Proxeon так, чтобы соединения были стабильными на протяжении всей выгрузки. Стабильное соединение сокращает число обрывов, а значит, ваш загрузчик реже уходит в возобновление и работает быстрее.
FAQ: частые вопросы по устойчивой выгрузке
Как понять, поддерживает ли сервер докачку файла? Отправьте запрос HEAD и посмотрите заголовок Accept-Ranges. Значение bytes означает поддержку. Отсутствие заголовка или none означает, что докачка невозможна.
Что делать, если API не отдаёт курсор, только страницы? Сохраняйте номер страницы и, если возможно, идентификатор последней записи. Возобновляйтесь предпочтительно по ID, потому что номера страниц сдвигаются при изменении данных.
Как часто сохранять чекпоинт? После каждой успешно обработанной и записанной страницы. Так при обрыве вы теряете максимум одну страницу работы, а не всю выгрузку.
Можно ли дедуплицировать без базы данных? На малых объёмах да, обычным множеством в памяти. На миллионах строк это опасно из-за памяти. Лучше использовать базу с PRIMARY KEY или фильтр Блума.
Что выбрать в качестве ключа дедупликации? Естественный уникальный ID источника, если он есть. Если нет, комбинацию стабильных полей. В крайнем случае хеш всей записи с сортировкой ключей.
Почему возобновление по ID надёжнее, чем по номеру страницы? Потому что данные в источнике могут меняться. Новые записи сдвигают страницы, и по номеру вы пропустите или задвоите данные. ID от этого не зависит.
Сколько параллельных воркеров ставить? Начните с небольшого числа и повышайте, наблюдая за стабильностью и лимитами источника. Чрезмерная параллельность вредит больше, чем помогает.
Как хранить строку подключения к прокси безопасно? В переменной окружения, а не в коде. Читайте её через os.environ. Так пароль не попадёт в систему контроля версий.
Что делать, если курсор протух за время долгой паузы? Поймайте ошибку недействительного курсора и переключитесь на возобновление по сохранённому идентификатору последней записи.
Нужно ли проверять целостность скачанного файла? Да. Сравните реальный размер файла с заголовком Content-Length. Если сервер отдаёт контрольную сумму, проверьте и её.
Заключение: что вы теперь умеете и куда двигаться
Вы прошли путь от хрупкой выгрузки, которая рушится при первом обрыве, до устойчивого загрузчика. Теперь у вас есть все инструменты, чтобы обрыв перестал быть катастрофой и стал обычной рабочей ситуацией.
Что вы освоили. Докачку файлов через заголовок Range с проверкой Accept-Ranges. Сохранение прогресса через атомарные чекпоинты. Идемпотентность на уровне ключа дедупликации. Дедупликацию на миллионах строк без раздувания памяти. Параллелизм с повтором упавших заданий. Возобновление после долгой паузы с обновлением токенов и заменой протухших курсоров. И самое главное, готовый каркас на Python, который объединяет всё это.
Что делать дальше. Возьмите свой реальный источник данных и адаптируйте под него функцию получения страницы. Начните с малого объёма, отладьте возобновление на прерываниях, а потом масштабируйте. Настройте стабильный прокси Proxeon, чтобы соединения были предсказуемыми на всём протяжении выгрузки.