Jak pobierać duże ilości danych przez proxy bez restartu od zera po zerwaniu połączenia: przewodnik krok po kroku
Spis treści
- Wprowadzenie: dlaczego długie pobieranie prawie zawsze się rwie i to jest normalne
- Przygotowanie wstępne i podstawowe pojęcia
- Krok 1: wznawianie pliku po http przez nagłówek range
- Krok 2: checkpointy dla pobierania stronicowego
- Krok 3: idempotentność zapisu wyników
- Krok 4: deduplikacja wyników bez rozdymania pamięci
- Krok 5: równoległość bez strat
- Krok 6: wznawianie po długiej przerwie
- Krok 7: gotowy szkielet stabilnego loadera w pythonie
- Typowe błędy i rozwiązania
- Faq: częste pytania o stabilne pobieranie
- Podsumowanie
Wprowadzenie: dlaczego długie pobieranie prawie zawsze się rwie i to jest normalne
Jeśli choć raz pobierałeś przez proxy kilka milionów rekordów albo plik o rozmiarze kilkudziesięciu gigabajtów, znasz to uczucie. Skrypt działał sześć godzin, pokazywał 83 procent, a potem padł z błędem połączenia. I wszystko, co masz, to niekompletny plik i świadomość, że trzeba zacząć od nowa.
Pierwsze, co trzeba zaakceptować: długie pobieranie zawsze się rwie. Nie „czasami”, nie „przy słabym łączu”, ale zawsze, jeśli trwa wystarczająco długo. Przyczyn są dziesiątki i większość z nich jest poza twoją kontrolą:
- Proxy zmienia zewnętrzny adres IP. W proxy mobilnych Proxeon to normalne zachowanie: rotacja po czasie albo na żądanie. W momencie zmiany IP otwarte połączenie TCP zostaje zerwane, a serwer źródłowy widzi już innego klienta.
- Serwer źródłowy zamyka połączenie po swoim timeoucie, restartuje się, wdraża aktualizację albo po prostu zwraca błąd 5xx.
- Twój własny proces się restartuje: aktualizacja systemu, zapełniony dysk, błąd w kodzie na nietypowym rekordzie, przypadkowy Ctrl+C.
- Wygasa token autoryzacji, kończy się sesja, przedawnia się kursor paginacji.
- Laptop zasypia, Wi-Fi przełącza się na inny punkt, provider zmienia trasę.
Walka z każdą z tych przyczyn osobno nie ma sensu. Właściwe podejście jest inne: zaprojektować pobieranie tak, żeby zerwanie w dowolnym momencie kosztowało cię nie sześć godzin, a jedną stronę albo jeden fragment pliku. Właśnie temu poświęcony jest ten przewodnik.
Co dostaniesz na końcu
Po przejściu instrukcji będziesz mieć:
- Działającą funkcję wznawiania pliku od środka po protokole HTTP przez nagłówek Range z weryfikacją, czy serwer to obsługuje.
- Przejrzysty schemat checkpointów dla pobierania stronicowego: co dokładnie zapisywać i gdzie trzymać stan, żeby nie przepadł razem z procesem.
- Idempotentny zapis wyników, przy którym ponowne pobranie tej samej strony nie tworzy duplikatów.
- Deduplikację, która nie zjada całej pamięci RAM na milionach wierszy.
- Równoległy loader z kolejką zadań, ponawianiem nieudanych zadań i limitem jednoczesności.
- Gotowy szkielet stabilnego loadera w Pythonie, który dostosujesz do swojego źródła w godzinę.
Dla kogo jest ten przewodnik
Dla programistów i analityków, którzy już umieją robić zapytania HTTP z Pythona i choć raz stracili wyniki długiego pobierania. Poziom średni: podstawy wyjaśniamy, ale nie uczymy programowania od zera. Zaawansowani czytelnicy znajdą rozdziały o trzymaniu hashy poza pamięcią i o bezpiecznym równoległym zapisie do SQLite.
Co trzeba wiedzieć wcześniej
- Python na poziomie funkcji, pętli, słowników i obsługi wyjątków.
- Podstawy HTTP: co to metoda, nagłówek, kod odpowiedzi, ciało.
- Ogólne pojęcie, jak ustawić proxy w bibliotece requests.
Osobno zaznaczmy: kodów odpowiedzi i strategii ponawiania z wykładniczą pauzą tu nie omawiamy. Temu poświęcony jest osobny artykuł o błędzie 429 i retry. W tym przewodniku skupiamy się na czymś innym: na stanie pobierania i jego wznawianiu. Retry odpowiadają na pytanie „kiedy powtórzyć zapytanie”, a my odpowiadamy na pytanie „od którego miejsca kontynuować pracę, gdy ponowienia się wyczerpią i proces umrze”.
Ile czasu to zajmie
Czytanie i uruchomienie przykładów na źródle testowym: dwie-trzy godziny. Dostosowanie szkieletu do twojego rzeczywistego API albo serwera plików: jeszcze jedna-dwie godziny w zależności od tego, jak niestandardowa jest paginacja. Razem dzień pracy z zapasem.
Przygotowanie wstępne i podstawowe pojęcia
Narzędzia i dostępy
- Zainstaluj Pythona w wersji 3.11 lub nowszej. W 2026 roku aktualne są gałęzie 3.12 i 3.13, wszystkie przykłady są na nich sprawdzone. Sprawdź wersję: otwórz terminal i wpisz
python --version. Jeśli zobaczysz 3.11 lub wyżej, wszystko w porządku. - Zainstaluj bibliotekę requests:
pip install requests. Wystarczy wersja 2.32 i nowsza. Moduł sqlite3 wchodzi w skład standardowej biblioteki Pythona, nic nie trzeba instalować osobno. - Uzyskaj dostęp do proxy. Otwórz panel klienta Proxeon, wybierz właściwy kanał i skopiuj cztery wartości: host, port, login i hasło. Zwykle są zebrane w jeden ciąg w postaci
http://USER:PASS@HOST:PORT. Tego ciągu będziesz używać wszędzie dalej. - Wrzuć ciąg proxy do zmiennej środowiskowej, a nie do kodu. W Linuksie i macOS:
export PROXY_URL=http://USER:PASS@HOST:PORT. W Windows PowerShell:$env:PROXY_URL='http://USER:PASS@HOST:PORT'. Tak nie zacommitujesz hasła do repozytorium przez przypadek. - Sprawdź, czy proxy odpowiada. Wykonaj w terminalu:
curl -x $PROXY_URL -I https://api.example.com/, wstawiając adres swojego źródła. Powinieneś zobaczyć linię z kodem odpowiedzi, na przykładHTTP/2 200. Jeśli widzisz błąd autoryzacji proxy 407, sprawdź jeszcze raz login i hasło.
Wymagania systemowe
Dowolna maszyna z 2 GB wolnej pamięci RAM i dyskiem, na którym zmieści się wynik pobierania plus 20 procent zapasu na indeksy SQLite. Jeśli planujesz pobierać miliony rekordów, dysk jest ważniejszy niż pamięć: całe podejście opiera się na tym, że stan żyje na dysku, a nie w zmiennych procesu.
Kopie zapasowe
Plik stanu, który utworzysz poniżej (w przykładach to export.sqlite), stanie się najcenniejszym artefaktem całej pracy. Nabierz nawyku kopiowania go przed każdym eksperymentem z kodem: cp export.sqlite export.sqlite.bak. Raz to uratuje ci dobę pobierania.
Uwaga: nigdy nie edytuj pliku SQLite ręcznie w czasie pracy loadera. Nawet odczyt z zewnętrznego programu w niewłaściwym trybie może zablokować zapis i wywalić proces. Jeśli chcesz zobaczyć stan, zatrzymaj loader albo użyj trybu WAL, o którym opowiemy w rozdziale o równoległości.
Kluczowe terminy prostym językiem
- Checkpoint — zapisany na dysku znacznik „do tego miejsca wszystko pobrane i zapisane”. Po zerwaniu loader czyta checkpoint i kontynuuje od niego.
- Kursor — nieprzejrzysty ciąg, który API zwraca razem ze stroną i który trzeba przekazać, żeby dostać następną stronę. Sam kursora nie tworzysz i nie rozbierasz.
- Pobieranie stronicowe po offsecie — gdy prosisz o „stronę 37 po 500 rekordów”. Prosty schemat, ale przy dodawaniu nowych rekordów w źródle strony się przesuwają i pojawiają się duplikaty albo luki.
- Pobieranie stronicowe po kluczu (keyset) — gdy prosisz o „wszystko z identyfikatorem większym niż 184203, posortowane po identyfikatorze”. Najstabilniejszy schemat do wznawiania, jeśli źródło go obsługuje.
- Idempotentność — właściwość operacji, przy której jej powtórzenie daje ten sam wynik co jednokrotne wykonanie. Zapisaliśmy stronę dwa razy, a w bazie leży raz.
- Klucz deduplikacji — wartość, po której dwa rekordy są uznawane za ten sam. Idealnie, jeśli to identyfikator ze źródła; jeśli go nie ma, klucz liczy się jako hash stabilnych pól.
- Nagłówek Range — sposób, żeby poprosić serwer HTTP o oddanie nie całego pliku, a jego części, na przykład bajtów od 1048576 do końca.
- Semantyka at-least-once — gwarancja, że każdy rekord zostanie pobrany co najmniej raz. Możliwe są powtórzenia, ale nie luki. Właśnie ją uzyskasz po tym przewodniku, a duplikaty usunie deduplikacja.
Główna zasada
Wszystkie siedem kroków poniżej sprowadza się do jednej idei: każda jednostka pracy musi być atomowa i powtarzalna. Jednostka pracy to albo fragment pliku, albo strona API, albo zadanie z kolejki. Atomowa — czyli wynik i znacznik jego zakończenia zapisują się razem. Powtarzalna — czyli jeśli jednostkę wykonasz dwa razy, nic się nie zepsuje. Gdy obie właściwości są spełnione, zerwanie w dowolnym punkcie staje się bezpieczne.
Krok 1: Wznawianie pliku po HTTP przez nagłówek Range
Cel etapu: nauczyć się pobierać duży plik przez proxy tak, żeby po zerwaniu pobieranie kontynuowało się od tego bajtu, na którym się zatrzymało, a nie od zera.
Jak to działa
Protokół HTTP pozwala klientowi zażądać części zasobu. W tym celu do zapytania dodaje się nagłówek Range: bytes=POCZĄTEK-, gdzie POCZĄTEK to offset w bajtach. Jeśli serwer obsługuje zapytania częściowe, odpowiada kodem 206 Partial Content i nagłówkiem Content-Range: bytes POCZĄTEK-KONIEC/RAZEM. Jeśli nie obsługuje, zignoruje Range i odda cały plik z kodem 200. Twoim zadaniem jest rozróżniać te przypadki.
Zawczasu sprawdzić, czy serwer obsługuje wznawianie, pomaga zapytanie HEAD: zwraca tylko nagłówki bez ciała. Patrz na Accept-Ranges: bytes. Wartość none albo brak nagłówka zwykle oznacza, że wznawiania nie ma, choć niektóre serwery i tak obsługują Range poprawnie, dlatego ostateczne sprawdzenie robimy po kodzie odpowiedzi.
Instrukcja krok po kroku
- Zrób zapytanie HEAD przez proxy i zapisz nagłówki Accept-Ranges, Content-Length, ETag i Last-Modified. ETag przyda się, żeby zrozumieć, czy plik na serwerze zmienił się między twoimi próbami.
- Sprawdź, ile bajtów leży już w lokalnym pliku. Jeśli pliku nie ma, przyjmij zero.
- Jeśli lokalny rozmiar jest już równy Content-Length, plik jest kompletny, nic nie trzeba robić.
- Jeśli lokalny rozmiar jest większy od zera i serwer deklaruje obsługę zakresów, dodaj do zapytania nagłówek Range z bieżącym rozmiarem. Dodaj też
If-Rangez zapisanym ETagiem: wtedy serwer odda odpowiedź częściową tylko wtedy, gdy plik się nie zmienił, a w przeciwnym razie zwróci cały plik z kodem 200. - Koniecznie wyślij
Accept-Encoding: identity. Bez tego serwer może zastosować kompresję w locie i offsety bajtowe przestaną się zgadzać z twoim plikiem. - Otwórz lokalny plik w trybie dopisywania
ab, jeśli dostałeś 206, albo w trybie nadpisywaniawb, jeśli dostałeś 200. - Czytaj ciało strumieniowo kawałkami po 256 KB i zapisuj na dysk. Nie ładuj całej odpowiedzi do pamięci.
- Po zakończeniu porównaj końcowy rozmiar z Content-Length. Jeśli się nie zgadza, połączenie zerwało się cicho i potrzebne jest jeszcze jedno podejście.
Działający kod
import os
import requests
PROXY_URL = os.environ['PROXY_URL'] # ciąg z panelu klienta 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')Zwróć uwagę na funkcję download_until_done: nie zawiera logiki oczekiwania między próbami. To celowe. Wstaw tam swoją strategię pauz z artykułu o retry, tutaj ważny jest tylko cykl „sprawdził rozmiar, dociągnął, sprawdził ponownie”.
Wskazówka: jeśli plik jest udostępniany jako archiwum, nie rozpakowuj go w locie w czasie wznawiania. Najpierw pobierz pełny plik, sprawdź rozmiar i, jeśli serwer podaje sumę kontrolną, zweryfikuj ją. Dopiero potem rozpakowuj. Częściowo pobrany gzip wygląda jak uszkodzony i stracisz czas na szukanie nieistniejącego błędu.
Oczekiwany wynik
Sprawdzenie: uruchom skrypt na pliku o rozmiarze co najmniej 200 MB, po dziesięciu sekundach przerwij go kombinacją Ctrl+C. Sprawdź rozmiar lokalnego pliku, na przykład 41 943 040 bajtów. Uruchom skrypt ponownie. W konsoli nie powinna pojawić się linia „zaczynamy od zera”, a rozmiar pliku powinien dalej rosnąć, a nie się resetować. Po zakończeniu końcowy rozmiar musi dokładnie zgadzać się z Content-Length z zapytania HEAD.
Możliwe problemy
- Serwer zawsze zwraca 200 zamiast 206. Znaczy to, że wznawianie nie jest obsługiwane. Jedyne wyjście dla takiego źródła to pobranie całego pliku za jednym podejściem z dużym timeoutem albo poszukanie u źródła alternatywnego formatu pobierania w częściach, na przykład podziału po datach.
- Brak nagłówka Content-Length. Serwer oddaje plik w trybie chunked bez deklarowania rozmiaru. Nie można sprawdzić kompletności po rozmiarze, nie można też wznawiać: Range wymaga znanych offsetów. Dogadaj się ze źródłem albo użyj sumy kontrolnej, jeśli jest publikowana.
- Odpowiedź 416 już przy pierwszym podejściu. Lokalny plik jest większy niż plik na serwerze. Plik na serwerze się zmienił i stał się krótszy. Usuń lokalny plik i zacznij od nowa.
- Rozmiar się zgadza, ale plik jest uszkodzony. Najprawdopodobniej gdzieś w środku było zerwanie z kodem 200 i zapisem od zera, a potem dopisanie. Odtwórz plik od nowa. Żeby to się nie powtarzało, trzymaj ETag w osobnym pliku obok i porównuj przed każdym podejściem.
Krok 2: Checkpointy dla pobierania stronicowego
Cel etapu: zapisywać stan pobierania tak, żeby po każdym padnięciu proces kontynuował od ostatniej poprawnie zapisanej strony.
Co zapisywać
Minimalny checkpoint zależy od typu paginacji w źródle. Przeanalizujmy trzy przypadki.
- Paginacja po kursorze. API zwraca razem z danymi pole w rodzaju
next_cursor. Zapisuj właśnie jego. To najprostszy przypadek: kursor już zawiera wszystko, czego serwer potrzebuje do kontynuacji. - Paginacja po numerze strony albo offsecie. Zapisuj numer ostatniej w pełni zapisanej strony i rozmiar strony. Pamiętaj, że przy dodawaniu rekordów do źródła w czasie pobierania offsety się przesuwają, dlatego deduplikacja z kroku 4 jest obowiązkowa.
- Paginacja po kluczu. Zapisuj identyfikator ostatniego zapisanego rekordu. Przy wznawianiu żądasz wszystkiego, co większe od tego identyfikatora. Schemat nie boi się ani wstawek, ani długich pauz.
Niezależnie od typu paginacji do checkpointu warto dodać pola pomocnicze: identyfikator ostatniego rekordu (nawet dla schematu kursorowego, to zapasowa kotwica na wypadek, gdyby kursor wygasł), liczniki stron i wierszy do kontroli postępu, czas startu pobierania i czas ostatniej aktualizacji.
Gdzie trzymać stan
Są dwa działające warianty i oba są lepsze niż zmienne w pamięci.
Wariant A: plik JSON z atomowym zastąpieniem
Pasuje, jeśli wyniki zapisujesz do osobnych plików, a nie do bazy. Główna pułapka: jeśli piszesz stan bezpośrednio do docelowego pliku i proces padnie w połowie zapisu, dostaniesz obcięty JSON, który się nie odczyta. Rozwiązanie — pisać do pliku tymczasowego obok i zmieniać jego nazwę na docelową. Operacja zmiany nazwy w jednym systemie plików jest atomowa.
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)Wariant B: tabela w SQLite obok danych
Preferowany wariant, jeśli rekordy składasz w bazie. Checkpoint aktualizuje się w tej samej transakcji, co wstawienie wierszy strony. Albo zapisały się i dane, i znacznik, albo nic. Rozdziału między „dane są, znacznika nie ma” nie ma w zasadzie wcale.
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))Kolejność operacji
Zapamiętaj zasadę: najpierw dane, potem checkpoint, i najlepiej w jednej transakcji. Jeśli transakcja jest niedostępna (na przykład dane zapisują się do plików), kolejność jest właśnie taka: zapisaliśmy plik strony, zsynchronizowaliśmy na dysk, potem zaktualizowaliśmy stan. Przy padnięciu między tymi dwiema czynnościami dostaniesz ponowne pobranie jednej strony, a to jest bezpieczne dzięki krokowi 3. Odwrotna kolejność da pominięcie strony, a to już utrata danych.
Wskazówka: trzymaj w checkpoincie nie bieżący kursor, a kursor następnej strony, który zwrócił serwer. Wtedy przy wznawianiu od razu żądasz tego, czego jeszcze nie masz, bez zbędnego zapytania o już pobraną stronę.
Oczekiwany wynik
Sprawdzenie: uruchom pobieranie, poczekaj na dziesięć stron i zakończ proces wymuszenie. Otwórz bazę komendą sqlite3 export.sqlite i wykonaj SELECT pages, last_id FROM checkpoint;. Powinieneś zobaczyć liczbę 10 i identyfikator. Potem wykonaj SELECT count(*) FROM records; i upewnij się, że liczba rekordów równa się dziesięciu rozmiarom strony. Uruchom loader ponownie: pierwszym komunikatem w konsoli powinno być coś w rodzaju „start: stron 10”.
Możliwe problemy
- Błąd „database is locked”. Inny proces trzyma połączenie. Zamknij wszystkie okna sqlite3 i inne narzędzia, które otwierały plik. Do pracy wielowątkowej włącz tryb WAL, patrz krok 5.
- Kursor się zapisał, a dane nie. Zaktualizowałeś checkpoint poza transakcją z danymi. Wróć do kodu powyżej i upewnij się, że obie operacje są wewnątrz jednego bloku
with con:. - JSON stanu okazał się pusty albo uszkodzony. Pisałeś bezpośrednio do pliku bez pliku tymczasowego i zastąpienia. Użyj funkcji save_state w całości.
Krok 3: Idempotentność zapisu wyników
Cel etapu: sprawić, żeby ponowne przetworzenie dowolnej strony nie tworzyło duplikatów i nie psuło danych.
Dlaczego powtórzenie jest nieuniknione
Po kroku 2 widziałeś już scenariusz, w którym strona zapisuje się dwa razy: proces padł po wstawieniu danych, ale przed aktualizacją checkpointu. Poza tym powtórzenia przychodzą z paginacji po offsecie przy zmianach w źródle, z równoległych workerów, które po restarcie dostały to samo zadanie, i po prostu z ręcznego restartu „na wszelki wypadek”. Walczyć z powtórzeniami po stronie zapytania nie ma sensu. Właściwie trzeba zrobić sam zapis tak, żeby powtórzenie było nieszkodliwe.
Wybór klucza deduplikacji
- Źródło ma identyfikator. Użyj go. To pole
id,uuid,order_numberalbo podobne, które źródło gwarantuje jako unikalne. Jeśli pobierasz z kilku źródeł do jednej tabeli, zrób klucz złożony: nazwa źródła plus identyfikator. - Identyfikatora nie ma, ale jest zestaw pól, które razem określają rekord. Na przykład dla wiersza cennika to artykuł plus magazyn plus data. Złóż klucz z tych pól, normalizując je: sprowadź ciągi do jednego rozmiaru liter, usuń spacje z brzegów, daty przekształć do jednolitego formatu.
- Nie ma nic stabilnego. Wtedy kluczem staje się hash całego rekordu po kanonizacji. Ten przypadek jest szczegółowo opisany w kroku 4. Uwzględnij, że jeśli źródło zmienia rekord (aktualizuje cenę), hash się zmieni i dostaniesz obie wersje. Czasem to jest właśnie to, czego chcesz, a czasem nie.
Uwaga: nie używaj jako klucza numeru porządkowego wiersza w odpowiedzi ani numeru strony. Te wartości zmieniają się przy każdej zmianie źródła i deduplikacja zamieni się w generator duplikatów.
Idempotentny zapis do bazy
W SQLite i większości relacyjnych baz jest konstrukcja, która albo ignoruje konflikt po kluczu głównym, albo aktualizuje istniejący wiersz. Pierwszy wariant INSERT OR IGNORE widziałeś już w kroku 2. Pasuje, gdy rekordy są niezmienne. Drugi wariant jest potrzebny, jeśli źródło może aktualizować rekordy i chcesz mieć świeżą wersję:
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])Idempotentny zapis do plików
Jeśli wynik ma leżeć w plikach, a nie w bazie, stosuj tę samą zasadę: jedna strona równa się jednemu plikowi z deterministyczną nazwą. Nazwa zależy od parametrów strony, a nie od czasu czy licznika. Przed pobraniem sprawdzasz, czy docelowy plik istnieje; jeśli tak, stronę pomijasz. Piszesz pod nazwą tymczasową i zmieniasz nazwę po zakończeniu, jak w funkcji 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 TruePliki z rozszerzeniem .part, które zostały po padnięciu, można śmiało usuwać przy starcie: z definicji są niekompletne.
Wskazówka: przy pobieraniu po offsecie nie polegaj tylko na „plik istnieje, więc strona jest gotowa”. Dodatkowo sprawdzaj, że liczba wierszy w pliku równa się rozmiarowi strony (poza ostatnią). Pusty albo krótki plik przy istniejącej nazwie lepiej pobrać ponownie.
Oczekiwany wynik
Sprawdzenie: wywołaj funkcję zapisu tej samej strony trzy razy pod rząd. Potem wykonaj SELECT count(*) FROM records;. Liczba powinna równać się rozmiarowi jednej strony, a nie potrójnemu. Dla wariantu plikowego w katalogu powinien być dokładnie jeden plik strony i ani jednego pliku .part.
Krok 4: Deduplikacja wyników bez rozdymania pamięci
Cel etapu: odsiewać powtarzające się rekordy w strumieniu milionów wierszy, nie trzymając wszystkich kluczy w pamięci RAM.
Hash rekordu
Gdy rekord nie ma identyfikatora, kluczem staje się hash jego zawartości. Żeby identyczne rekordy dawały identyczny hash, zawartość trzeba skanonizować: posortować klucze słownika, usunąć zbędne spacje, ustalić separatory. Inaczej ten sam rekord, który przyszedł z inną kolejnością pól, dostanie inny 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()Funkcja zwraca 16 bajtów. To wystarczy: prawdopodobieństwo przypadkowego trafienia dla setek milionów rekordów jest pomijalnie małe. Parametr fields pozwala liczyć hash tylko po stabilnych polach, wykluczając na przykład czas ostatniej aktualizacji, który zmienia się przy każdym zapytaniu.
Dlaczego zbiór w pamięci nie działa na milionach
Pierwsze, co przychodzi do głowy: założyć seen = set() i wrzucać tam klucze. Policzmy. Jeden obiekt bytes o długości 16 bajtów zajmuje w Pythonie około 49 bajtów plus same dane, razem mniej więcej 65 bajtów. Slot w zbiorze z uwzględnieniem współczynnika wypełnienia dodaje jeszcze około 30 bajtów. Dostajemy rzędu 95 bajtów na klucz. Na 10 milionach rekordów to około 950 MB, na 50 milionach prawie 5 GB. I co najważniejsze: po restarcie procesu zbiór jest pusty i cała deduplikacja zaczyna się od czystej karty.
Trzy sposoby, żeby nie rozdymać pamięci
- Trzymać klucze w samej bazie. Najprostsza i najpewniejsza droga. Jeśli klucz jest kluczem głównym tabeli records, deduplikacja jest już zrobiona konstrukcją INSERT OR IGNORE z kroku 3. Indeks żyje na dysku, przeżywa restarty, a SQLite sam cachuje gorące strony indeksu. Dla 10 milionów 16-bajtowych kluczy indeks zajmie około 400-500 MB na dysku, ale nie w pamięci.
- Osobna tabela widzianych kluczy bez rowid. Potrzebna, jeśli same dane zapisujesz nie do SQLite, a na przykład do plików. Wtedy SQLite służy tylko jako kompaktowy zbiór na dysku.
- Skompresowany zbiór w pamięci jako prefiltr. Wariant zaawansowany: obciąć hash do 8 bajtów i trzymać jako liczbę całkowitą w posortowanej tablicy albo użyć filtra Blooma. Pamięć zmniejsza się kilkukrotnie, ale pojawia się prawdopodobieństwo fałszywego trafienia. Dlatego taki prefiltr stosuje się tylko po to, żeby szybko odciąć z góry nowe rekordy, a ostateczne sprawdzenie i tak robi się po bazie.
Implementacja zbioru na dysku
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])Wywołuj filter_new przed zapisem strony, a remember wewnątrz tej samej transakcji, co zapis danych i checkpointu. Wtedy po padnięciu zbiór widzianych, dane i znacznik postępu są zawsze ze sobą zgodne.
Wskazówka: zapytanie z IN na 500 wartości wykonuje się po indeksie w milisekundach. Nie sprawdzaj kluczy pojedynczo w pętli: to dziesiątki razy wolniejsze z powodu narzutu na każde zapytanie.
Oczekiwany wynik
Sprawdzenie: zbuduj testową stronę z 500 rekordów, gdzie 100 powtarza się dwa razy w obrębie strony, a kolejne 100 leży już w tabeli seen. Funkcja filter_new powinna zwrócić dokładnie 300 rekordów. Po restarcie procesu te same 500 rekordów powinno dać zero nowych.
Możliwe problemy
- Duplikaty i tak przechodzą. Sprawdź kanonizację: najprawdopodobniej w rekordach jest pole z czasem zapytania albo losową kolejnością elementów w liście. Wyklucz je przez parametr fields albo posortuj zagnieżdżone listy przed hashowaniem.
- Wolne wstawianie po kilku milionach wierszy. Indeks przestał mieścić się w cache. Zwiększ cache SQLite komendą
PRAGMA cache_size=-200000(to 200 MB) i upewnij się, że wstawienia idą paczkami w jednej transakcji na stronę, a nie po jednym wierszu.
Krok 5: Równoległość bez strat
Cel etapu: przyspieszyć pobieranie kilkoma jednoczesnymi workerami przez proxy tak, żeby padnięcie któregokolwiek nie gubiło zadań i nie psuło bazy.
Kiedy można paralelizować, a kiedy nie
Paginacja po kursorze jest z natury sekwencyjna: następny kursor jest znany dopiero po pobraniu poprzedniej strony. Nie da się jej paralelizować bezpośrednio. Ale prawie zawsze można podzielić pobieranie na niezależne shardy: po dniach, po kategoriach, po regionach, po pierwszych znakach identyfikatora. Każdy shard pobiera się sekwencyjnie ze swoim checkpointem, a shardy idą równolegle. Paginacja po offsecie i po kluczu ze znanymi granicami paralelizuje się bezpośrednio: zadania w rodzaju „strony od 1 do 100” albo „identyfikatory od 0 do 100000”.
Kolejka zadań na dysku
Kolejka w pamięci umiera razem z procesem. Dlatego zadania żyją w tabeli ze statusami:
pending— czeka na wykonanie;running— wzięte przez workera;done— wykonane i zapisane;failed— wyczerpane próby, wymaga uwagi człowieka.
Przy starcie loader najpierw przekłada wszystkie zadania z running z powrotem na pending: jeśli wiszą w tym statusie, znaczy to, że poprzedni proces umarł w połowie pracy. Potem workery rozbierają pending.
Jeden punkt zapisu
SQLite dopuszcza wielu jednoczesnych czytelników, ale tylko jednego pisarza. Najprostszy i najbezpieczniejszy wzorzec: workery tylko pobierają i zwracają dane, a cały zapis do bazy robi wątek główny. Żadnych blokad w kodzie, żadnych „database is locked”. Dodatkowo włącz tryb WAL, żeby odczyt stanu z innego procesu nie przeszkadzał zapisowi.
Limit jednoczesności
Liczbę workerów ograniczają dwie rzeczy. Pierwsza — możliwości proxy: jeśli w panelu klienta Proxeon masz kilka kanałów, rozsądnie jest trzymać po jednym-dwa workery na kanał, żeby rotacja IP na jednym kanale nie rwała połączeń wszystkich wątków naraz. Druga — uprzejmość wobec źródła: nawet bez formalnych limitów dziesięć równoległych wątków na małym API stworzy obciążenie, przez które dostaniesz odmowy. Zaczynaj od trzech-czterech workerów i podnoś, obserwując udział błędów.
Kod równoległego loadera
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)Funkcja fetch_fn wykonuje zapytanie przez proxy i zwraca listę rekordów. Działa w wątku i nie dotyka bazy. Funkcja write_fn jest wywoływana w wątku głównym wewnątrz transakcji i robi idempotentne wstawienie z kroku 3. Padnięte zadanie automatycznie wraca do pending i zostanie wzięte ponownie w następnym cyklu claim; po wyczerpaniu prób dostaje status failed i zajmiesz się nim ręcznie.
Uwaga: nie przekazuj obiektu połączenia sqlite3 do workerów. Połączenie jest przywiązane do wątku, w którym powstało, a próba użycia go z innego wątku doprowadzi do błędu albo, co gorsza, do cichego uszkodzenia danych. Każdemu wątkowi albo własne połączenie, albo, jak w przykładzie powyżej, żadne.
Wskazówka: używaj na każdego workera osobnego obiektu requests.Session z własnym adresem proxy. Jeśli w Proxeon masz kilka kanałów, rozdziel je po workerach po kolei: worker 0 bierze kanał 0, worker 1 kanał 1 i tak dalej. Wtedy zerwanie połączenia na jednym kanale dotknie tylko jeden wątek.
Oczekiwany wynik
Sprawdzenie: postaw w kolejce 100 zadań, uruchom cztery workery i zabij proces po pół minuty. Wykonaj SELECT status, count(*) FROM tasks GROUP BY status;. Zobaczysz kilka done, kilka running i resztę pending. Uruchom ponownie: running powinny zniknąć przy starcie, a po zakończeniu wszystkie zadania powinny znaleźć się w done, poza tymi, które uczciwie padły i leżą w failed z tekstem błędu w last_error.
Krok 6: Wznawianie po długiej przerwie
Cel etapu: poprawnie kontynuować pobieranie, które zatrzymałeś na kilka godzin albo dni, i nie natrafić na przedawniony stan.
Co się przedawnia
Wznowienie po dziesięciu sekundach i po tygodniu to dwie różne sprawy. Po długiej przerwie część zapisanego stanu przestaje być aktualna.
- Sesje i cookies. Sesje serwerowe żyją zwykle od kilku godzin do doby. Zapisane cookies po tym czasie doprowadzą do odpowiedzi 401 albo przekierowania na formularz logowania. Rozwiązanie: przy starcie wykonywać pełną ponowną autoryzację, a nie odtwarzać cookies z pliku.
- Tokeny dostępu. Tokeny OAuth żyją godzinę, czasem krócej. Jeśli masz refresh token, odświeżaj access token przed startem i zgodnie z harmonogramem podczas pracy, nie czekając na odmowę.
- Kursory paginacji. Wiele API ogranicza czas życia kursora do minut albo godzin. Przedawniony kursor zwróci błąd 400 z komunikatem o nieprawidłowym kursorze. Właśnie dlatego w kroku 2 zapisywaliśmy zapasową kotwicę: identyfikator ostatniego rekordu. Jeśli źródło obsługuje filtr po identyfikatorze albo po dacie zmiany, budujesz nowe zapytanie od tej kotwicy. Jeśli nie obsługuje, trzeba zacząć shard od początku, a deduplikacja z kroku 4 odsieje już pobrane.
- Zawartość pliku na serwerze. Dla wznawiania z kroku 1 kluczowe jest, że plik się nie zmienił. Porównaj bieżący ETag z zapisanym przed każdym podejściem; przy niezgodności zacznij plik od nowa.
- Ustawienia proxy. Przez tydzień w panelu klienta Proxeon mogły się zmienić port, hasło albo skończył się okres kanału. Sprawdzaj proxy zapytaniem testowym, zanim zaczniesz rozbierać kolejkę.
- Sam zbiór danych. Jeśli pobieranie trwa tydzień, a źródło w tym czasie dodało i usunęło rekordy, twój wynik będzie mieszanką stanów z różnych momentów. Dla wielu zadań jest to akceptowalne. Jeśli nie, trzymaj czas startu w checkpoincie i po zakończeniu rób osobny przyrostowy przebieg po rekordach zmienionych po tym czasie.
Sprawdzenie przedstartowe
Zbierz wszystkie sprawdzenia w jedną funkcję, która wykonuje się przy starcie przed jakąkolwiek realną pracą. Albo doprowadza stan do porządku, albo zatrzymuje loader z jasnym komunikatem.
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 stateFunkcja refresh_access_token zależy od twojego źródła: zwykle to zapytanie POST z refresh tokenem, po którym aktualizujesz nagłówek Authorization w sesji. Pole resume_after_id jest potem używane w funkcji pobierania strony jako filtr „identyfikator większy od podanego”.
Wskazówka: trzymaj refresh token i hasło proxy nie w checkpoincie, a w zmiennych środowiskowych albo w osobnym pliku sekretów z ograniczonymi uprawnieniami. Checkpoint będziesz kopiować, przesyłać kolegom i załączać do raportów o błędach; sekrety są tam zbędne.
Oczekiwany wynik
Sprawdzenie: ręcznie zepsuj kursor w tabeli checkpoint komendą UPDATE checkpoint SET cursor='broken'; i uruchom loader. W konsoli powinna pojawić się linia o przełączeniu na kotwicę po last_id, a pobieranie powinno kontynuować się bez padnięcia. Liczba rekordów po zakończeniu powinna zgadzać się z kontrolnym uruchomieniem bez psucia kursora.
Krok 7: Gotowy szkielet stabilnego loadera w Pythonie
Cel etapu: złożyć wszystko z poprzednich kroków w jeden plik, który można uruchomić, przerwać, uruchomić ponownie i dostać pełny wynik bez duplikatów.
Struktura szkieletu
- Konfiguracja ze zmiennych środowiskowych: adres proxy Proxeon, adres API, token, nazwa zadania.
- Klasa Store: SQLite z tabelami rekordów i checkpointu, jedna transakcja na stronę.
- Funkcja klucza rekordu dla idempotentności.
- Funkcja pobierania strony przez proxy.
- Główna pętla ze wznawianiem po checkpoincie i odtwarzaniem sesji po zerwaniu.
Pełny kod
# 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 z panelu klienta 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()Jak dostosować do swojego źródła
- Zamień ścieżkę
/ordersi nazwy pólitems,next_cursor,idna te, które zwraca twoje API. To trzy miejsca w funkcjach fetch_page i record_key. - Jeśli źródło ma paginację po offsecie, zamień parametr cursor na page i licz następną wartość jako state['pages'] + 1. W checkpoincie zamiast kursora zapisuj numer strony.
- Jeśli paginacja jest po kluczu, przekazuj parametr w rodzaju
after_idz state['last_id'] i usuń obsługę kursora. - Jeśli potrzebujesz równoległości, wyciągnij fetch_page do fetch_fn z kroku 5, a commit_page użyj jako write_fn. Podziel pobieranie na shardy i wypełnij kolejkę zadań.
- Dodaj funkcję preflight z kroku 6 przed główną pętlą.
Sprawdzenie wyniku: lista kontrolna
Zanim uruchomisz loader na realnym kilkugodzinnym wolumenie, przepuść go przez tę listę. Każdy punkt zajmuje parę minut, a razem gwarantują, że nocne pobieranie nie straci danych.
- Test na przerwanie. Uruchom loader, po 30 sekundach naciśnij Ctrl+C. Uruchom ponownie. Pierwsza linia wyjścia powinna pokazywać niezerową liczbę stron, a nie „stron 0”.
- Test na duplikaty. Zmniejsz PAGE_SIZE do 10, przerwij loader pięć razy pod rząd w losowych momentach. Po zakończeniu porównaj liczbę unikalnych rekordów z liczbą pobranych wierszy: unikalnych powinno być mniej lub równo, a przy czystej paginacji kursorowej bez wstawek do źródła prawie równo.
- Test na spójność. Po każdym przerwaniu wykonaj dwa zapytania:
SELECT rows FROM checkpoint;iSELECT count(*) FROM records;. Różnica między nimi nie powinna przekraczać jednego rozmiaru strony. Jeśli przekracza, checkpoint i dane nie są zapisywane w jednej transakcji. - Test na proxy. Na chwilę wyłącz zmienną PROXY_URL albo podaj błędne hasło. Loader powinien paść na pierwszym zapytaniu z jasnym błędem, a nie zawisnąć i nie zacząć chodzić do źródła bezpośrednio.
- Test na przedawniony kursor. Zepsuj kursor w bazie, jak opisano w kroku 6, i upewnij się, że działa przełączenie na kotwicę.
- Test na dysk. Sprawdź rozmiar pliku export.sqlite po tysiącu stron i pomnóż przez oczekiwaną liczbę stron. Upewnij się, że na dysku wystarczy miejsca z 20-procentowym zapasem.
Sprawdzenie: wskaźnikiem udanego wykonania jest sytuacja, w której po trzech celowych przerwaniach i trzech restartach końcowa liczba unikalnych rekordów zgadza się z liczbą uzyskaną w jednym ciągłym przebiegu na tym samym źródle, a w konsoli ani raz nie pojawiła się linia o starcie od zera.
Dodatkowe możliwości i optymalizacja
- Postęp i szacowanie czasu. Jeśli znana jest całkowita liczba rekordów, wyświetlaj procent i szacunek pozostałego czasu raz na dwadzieścia stron. To przydaje się i tobie, i do odróżniania zawieszenia od wolnej pracy.
- Kompresja payload. Dla dziesiątek milionów rekordów JSON w postaci tekstowej zajmuje dużo miejsca. Kompresuj pole payload funkcją zlib.compress przed zapisem i trzymaj jako BLOB. Oszczędność zwykle od trzech do ośmiu razy.
- Przeniesienie na bazę serwerową. Schemat z checkpointem w jednej transakcji z danymi przenosi się na PostgreSQL prawie bez zmian. Konstrukcja ON CONFLICT jest tam obsługiwana, a ograniczenie „jeden pisarz” znika.
- Pobieranie przyrostowe. Zapisuj czas startu każdego zadania i po pełnym pobraniu uruchamiaj osobne zadanie z filtrem „zmienione po”. Tak utrzymujesz aktualną kopię bez pełnego ponownego pobierania.
- Osobny proces na shard. Zamiast wątków możesz uruchomić kilka instancji skryptu z różnymi wartościami JOB i różnymi plikami DB_PATH, a wyniki scalić na końcu. To prostsze w debugowaniu i całkowicie usuwa kwestię konkurencyjnego zapisu.
- Metryki zerwań. Loguj każde zerwanie z typem wyjątku i czasem. Po dobie zobaczysz, że zerwania grupują się wokół interwałów rotacji IP na kanale Proxeon, i będziesz mógł dopasować interwał rotacji do długości swoich zapytań.
Typowe błędy i rozwiązania
Poniżej zebrano sytuacje, z którymi spotyka się prawie każdy na pierwszych uruchomieniach. Format: problem, przyczyna, rozwiązanie.
- Problem: po restarcie pobieranie za każdym razem zaczyna się od zera. Przyczyna: checkpoint jest zapisywany do pamięci albo do pliku, który nie przeżywa padnięcia, ewentualnie loader nie czyta go przy starcie. Rozwiązanie: upewnij się, że pierwszą czynnością w funkcji run jest store.load, a stan aktualizuje się po każdej stronie wewnątrz transakcji.
- Problem: w bazie jest półtora raza więcej rekordów niż w źródle. Przyczyna: klucz deduplikacji jest niestabilny: trafił do niego czas zapytania, numer strony albo pole z losową kolejnością. Rozwiązanie: licz klucz tylko po identyfikatorze źródła albo po jawnej liście stabilnych pól przez parametr fields.
- Problem: w bazie jest mniej rekordów niż w źródle, choć pobieranie zakończyło się bez błędów. Przyczyna: checkpoint aktualizował się przed zapisem danych i po padnięciu strona została pominięta. Albo paginacja po offsecie przy usuwaniu rekordów w źródle przesunęła strony do tyłu. Rozwiązanie: zmień kolejność na „dane, potem checkpoint” w jednej transakcji; dla źródeł z usunięciami przejdź na paginację po kluczu.
- Problem: wznawianie pliku daje uszkodzone archiwum przy zgodnym rozmiarze. Przyczyna: serwer raz odpowiedział 200 zamiast 206, plik został częściowo nadpisany, a potem dopisany. Rozwiązanie: trzymaj ETag obok pliku, przy zmianie usuwaj plik; sprawdzaj nagłówek Content-Range pod kątem zgodności z żądanym offsetem.
- Problem: błąd „database is locked” przy pracy równoległej. Przyczyna: kilka wątków pisze do SQLite jednocześnie albo połączenie zostało przekazane między wątkami. Rozwiązanie: jeden punkt zapisu w wątku głównym, workery tylko pobierają; tryb WAL; połączenie tworzone w tym wątku, w którym jest używane.
- Problem: po godzinie pracy wszystkie zapytania zaczęły zwracać 401. Przyczyna: wygasł token dostępu. Rozwiązanie: odświeżaj token zgodnie z harmonogramem przed wygaśnięciem, a przy otrzymaniu 401 wywołaj odświeżenie i powtórz zapytanie raz, nie traktując tego jako zerwania.
- Problem: pamięć RAM rośnie do kilku gigabajtów. Przyczyna: zbiór widzianych kluczy albo lista wszystkich rekordów trzyma się w pamięci procesu. Rozwiązanie: deduplikacja przez klucz główny w bazie albo tabelę seen; zapis danych stronicowo, bez kumulowania.
- Problem: zerwania występują ściśle co kilka minut. Przyczyna: pokrywają się z interwałem rotacji IP na kanale proxy. Rozwiązanie: to normalna sytuacja, loader musi ją przeżywać. Jeśli zapytania są długie, dobierz interwał rotacji w panelu klienta Proxeon tak, żeby był wyraźnie większy od typowego czasu jednego zapytania, albo użyj rotacji na żądanie między stronami.
FAQ: częste pytania o stabilne pobieranie
Czy trzeba używać SQLite, jeśli wynik ma być w CSV?
Nie, ale to wygodne. SQLite pełni tu rolę niezawodnego magazynu stanu i zbioru widzianych kluczy. Końcowy CSV eksportujesz jedną komendą z tabeli records po zakończeniu. Jeśli chcesz całkiem bez bazy, użyj wariantu plikowego z kroku 3 z jednym plikiem na stronę i checkpointem JSON z atomowym zastąpieniem.
Jak często zapisywać checkpoint: po każdej stronie czy rzadziej?
Po każdej stronie. Jedna transakcja SQLite z kilkuset wierszami wykonuje się w milisekundach, to nic wobec zapytania sieciowego przez proxy. Oszczędność na rzadkich checkpointach nie jest warta ryzyka utraty dziesiątek stron.
Co zrobić, jeśli API nie zwraca ani kursora, ani identyfikatorów, tylko numery stron?
Pracuj po numerze strony, trzymaj go w checkpoincie i koniecznie włącz deduplikację po hashu zawartości rekordu. Przyjmij, że przy aktywnych zmianach w źródle część rekordów może zostać pominięta z powodu przesuwania stron. Dla krytycznych danych zrób drugi przebieg w odwrotnej kolejności stron: pominięte przy pierwszym przebiegu z dużym prawdopodobieństwem trafią do drugiego.
Czy można wznawiać plik kilkoma wątkami różnymi zakresami?
Można, jeśli serwer obsługuje Range. Podziel plik na kawałki po 50-100 MB, każdy kawałek to zadanie z kolejki kroku 5 z własnym plikiem tymczasowym, a po zakończeniu wszystkich kawałków sklej je w prawidłowej kolejności. Sprawdzaj każdy kawałek po rozmiarze, a cały plik po sumie kontrolnej, jeśli jest dostępna.
Ile workerów ustawić przy pracy przez proxy?
Zacznij od trzech-czterech na jeden kanał Proxeon i obserwuj udział błędów w tabeli tasks. Jeśli błędów jest mniej niż jeden procent, dodaj jeszcze dwa. Jeśli błędy rosną, zmniejszaj. Więcej niż dziesięć wątków na jeden kanał rzadko daje zysk: opierasz się albo o przepustowość kanału, albo o cierpliwość źródła.
Czy trzeba zapisywać same cookies sesji między uruchomieniami?
Zwykle nie. Ponowna autoryzacja przy starcie zajmuje sekundy i jest pewniejsza niż odtwarzanie cookies o nieznanym czasie życia. Wyjątek: źródło ogranicza liczbę logowań na dobę. Wtedy zapisuj cookies, ale przy pierwszym 401 albo przekierowaniu na login wyrzuć je i zaloguj się od nowa.
Jak poznać, że pobieranie zakończyło się w całości, a nie zerwało cicho?
Dla plików: rozmiar równa się Content-Length i suma kontrolna się zgadza. Dla API: pobrano stronę bez next_cursor albo pustą stronę, a przy tym liczba wierszy zgadza się z całkowitą liczbą, jeśli źródło ją podaje. Zapisuj w checkpoincie jawną flagę zakończenia, żeby ponowne uruchomienie nie zaczynało nowego obchodu.
Co zrobić z zadaniami w statusie failed?
Spójrz na pole last_error. Jeśli to błędy sieciowe, po prostu przywróć zadania do pending komendą UPDATE i uruchom loader ponownie. Jeśli to błędy parsowania danych, znaczy to, że w źródle są rekordy o niestandardowej formie: popraw kod i uruchom ponownie. Nigdy nie usuwaj failed po cichu, to jedyny dowód na to, czego brakuje w pobieraniu.
Czy można użyć tego podejścia nie z requests, a z biblioteką asynchroniczną?
Tak, zasady są te same: atomowa jednostka pracy, checkpoint razem z danymi, idempotentny zapis, kolejka na dysku. Zmienia się tylko transport. Jedyny niuans: zapis do SQLite zostaw synchroniczny i sekwencyjny, a równoległość trzymaj na poziomie zapytań sieciowych.
Podsumowanie
Przeszedłeś drogę od niekompletnego pliku i nerwowego restartu do loadera, dla którego zerwanie jest obojętne. Ustalmy, co dokładnie zostało zrobione.
- Rozpisaliśmy wznawianie po HTTP: sprawdzanie Accept-Ranges przez HEAD, nagłówek Range z bieżącym rozmiarem pliku, rozróżnianie kodów 206 i 200, ochrona przed podmianą pliku przez ETag i If-Range.
- Zbudowaliśmy checkpointy dla pobierania stronicowego: kursor, numer strony albo identyfikator ostatniego rekordu plus liczniki pomocnicze, wszystko w jednej transakcji z danymi.
- Zrobiliśmy zapis idempotentnym przez klucz główny i konstrukcję INSERT OR IGNORE albo ON CONFLICT DO UPDATE, a dla plików przez deterministyczne nazwy i atomowe zastąpienie.
- Zorganizowaliśmy deduplikację na dysku, żeby miliony kluczy nie żyły w pamięci RAM i przeżywały restarty.
- Dodaliśmy równoległość z kolejką zadań w SQLite, automatycznym przywracaniem padniętych zadań i jedynym punktem zapisu.
- Przewidzieliśmy wznawianie po długiej przerwie: odświeżanie tokenów, ponowna autoryzacja, przełączanie z przedawnionego kursora na kotwicę po identyfikatorze, sprawdzanie proxy Proxeon przed startem.
- Złożyliśmy wszystko w jeden działający szkielet, który dostosowuje się do konkretnego źródła zamianą trzech-czterech linii.
Co dalej
Weź szkielet z kroku 7 i uruchom go na niewielkim realnym wolumenie, powiedzmy na dziesięciu tysiącach rekordów. Przepuść listę kontrolną z sekcji sprawdzeń. Dopiero potem uruchom pełne pobieranie na noc. Rano albo zobaczysz linię „gotowe” ze zgodnymi licznikami, albo linię o wyczerpanych próbach z zapisanym stanem, i wtedy po prostu uruchom skrypt jeszcze raz.
Kierunki rozwoju
Następny poziom to pobieranie przyrostowe po czasie zmiany zamiast pełnych obchodów, przeniesienie stanu na bazę serwerową dla kilku maszyn, a także przemyślana strategia ponawiania z uwzględnieniem kodów odpowiedzi, której poświęcony jest osobny artykuł o 429 i retry. Połączenie sprawnych retry z tamtego artykułu i stabilnego stanu z tego daje loader, który można zostawić działający na tydzień i nie otwierać terminala.
I na koniec. Zerwanie długiego pobierania to nie awaria, a normalna sytuacja, którą teraz umiesz obsłużyć. Udanych pobierań.