Giriş: neden uzun indirmeler neredeyse her zaman kesilir ve bu normaldir

Eğer bir kez bile proxy üzerinden milyonlarca kayıt veya onlarca gigabaytlık bir dosya indirdiyseniz, bu hisse aşinasınızdır. Betik altı saat çalıştı, yüzde 83 gösterdi ve sonra bağlantı hatasıyla çöktü. Elinizde kalan tek şey eksik bir dosya ve baştan başlamanız gerektiği gerçeği.

Kabul etmeniz gereken ilk şey: uzun indirmeler her zaman kesilir. "Bazen" değil, "kötü ağda" değil, yeterince uzun sürerse her zaman. Bunun onlarca nedeni var ve çoğu sizin kontrolünüz dışında:

  • Proxy dış IP adresini değiştirir. Proxeon mobil proxy'lerinde bu normal davranıştır: zamanlayıcıya veya isteğe göre rotasyon. IP değiştiği anda açık TCP bağlantısı kopar ve kaynak sunucu artık farklı bir istemci görür.
  • Kaynak sunucu kendi zaman aşımına göre bağlantıyı kapatır, yeniden başlatılır, güncelleme yayınlar veya sadece 5xx hatası döndürür.
  • Kendi süreciniz yeniden başlar: sistem güncellemesi, disk dolması, alışılmadık bir kayıtta kod hatası, yanlışlıkla Ctrl+C.
  • Yetkilendirme belirteci süresi dolar, oturum sona erer, sayfalama imleci geçersiz olur.
  • Dizüstü bilgisayar uyku moduna geçer, Wi-Fi başka bir erişim noktasına geçer, sağlayıcı rotayı değiştirir.

Bu nedenlerin her biriyle ayrı ayrı savaşmak anlamsızdır. Doğru yaklaşım farklıdır: indirmeyi, herhangi bir anda kesilmenin size altı saat değil, sadece bir sayfa veya bir dosya parçasına mal olacağı şekilde tasarlamak. Bu rehber tam olarak bunu ele alıyor.

Sonunda ne elde edeceksiniz

Kılavuzu tamamladıktan sonra şunlara sahip olacaksınız:

  • HTTP Range başlığı aracılığıyla bir dosyayı ortasından devam ettirebilen ve sunucunun bunu destekleyip desteklemediğini kontrol eden çalışan bir fonksiyon.
  • Sayfalı indirmeler için net bir kontrol noktası şeması: tam olarak neyin kaydedileceği ve durumun süreçle birlikte kaybolmaması için nerede saklanacağı.
  • Aynı sayfanın yeniden yüklenmesinin kopya oluşturmadığı idempotent sonuç yazımı.
  • Milyonlarca satırda tüm RAM'i tüketmeyen tekilleştirme.
  • İş kuyruğu, başarısız görevlerin tekrarı ve eşzamanlılık sınırlaması olan paralel bir indirici.
  • Bir saat içinde kendi kaynağınıza uyarlayabileceğiniz, Python'da sağlam bir indirici iskeleti.

Bu rehber kimler için

Python'dan HTTP istekleri yapmayı zaten bilen ve en az bir kez uzun bir indirmenin sonuçlarını kaybeden geliştiriciler ve analistler için. Seviye orta: temel şeyleri açıklıyoruz ama sıfırdan programlama öğretmiyoruz. İleri düzey okuyucular, bellek dışı hash saklama ve SQLite'a güvenli paralel yazma bölümlerini faydalı bulacaktır.

Önceden bilmeniz gerekenler

  • Fonksiyonlar, döngüler, sözlükler ve istisna işleme düzeyinde Python.
  • HTTP temelleri: metot, başlık, yanıt kodu, gövde nedir.
  • requests kütüphanesinde proxy'nin nasıl ayarlanacağına dair genel bir fikir.

Ayrıca belirtelim: yanıt kodları ve üstel geri çekilme ile yeniden deneme stratejileri burada ele alınmıyor. Bu konuya 429 hatası ve yeniden denemeler hakkındaki ayrı bir makale ayrılmıştır. Bu rehberde odak farklı: indirmenin durumu ve devam ettirilmesi. Yeniden denemeler "isteği ne zaman tekrarlamalıyım" sorusunu yanıtlar, biz ise "tekrarlar tükendikten ve süreç öldükten sonra işe hangi noktadan devam etmeliyim" sorusunu yanıtlıyoruz.

Ne kadar zaman alacak

Test kaynağında örnekleri okumak ve çalıştırmak: iki-üç saat. İskeleti gerçek API'nize veya dosya sunucunuza uyarlamak: sayfalamanın ne kadar standart dışı olduğuna bağlı olarak bir-iki saat daha. Toplamda geniş bir iş günü.

Ön hazırlık ve temel kavramlar

Araçlar ve erişimler

  1. Python 3.11 veya daha yenisini kurun. 2026'da 3.12 ve 3.13 sürümleri güncel, tüm örnekler bunlarda test edilmiştir. Sürümü kontrol etmek için: terminali açın ve python --version yazın. 3.11 veya üstünü görürseniz her şey yolunda.
  2. requests kütüphanesini kurun: pip install requests. 2.32 ve üstü sürüm yeterlidir. sqlite3 modülü Python standart kütüphanesine dahildir, ayrıca kurmanıza gerek yok.
  3. Proxy erişimi alın. Proxeon kontrol panelini açın, istediğiniz kanalı seçin ve dört değeri kopyalayın: host, port, kullanıcı adı ve şifre. Genellikle http://USER:PASS@HOST:PORT biçiminde tek bir satırda birleştirilirler. Bu satırı bundan sonra her yerde kullanacaksınız.
  4. Proxy satırını koda değil, ortam değişkenine koyun. Linux ve macOS'ta: export PROXY_URL=http://USER:PASS@HOST:PORT. Windows PowerShell'de: $env:PROXY_URL='http://USER:PASS@HOST:PORT'. Böylece şifreyi yanlışlıkla depoya işlemeyeceksiniz.
  5. Proxy'nin yanıt verdiğini kontrol edin. Terminalde şunu çalıştırın: curl -x $PROXY_URL -I https://api.example.com/, kaynağınızın adresini yazarak. Örneğin HTTP/2 200 gibi bir yanıt kodu satırı görmelisiniz. Proxy yetkilendirme hatası 407 görürseniz, kullanıcı adı ve şifreyi tekrar kontrol edin.

Sistem gereksinimleri

2 GB boş RAM ve indirme sonucunun sığacağı, üstüne SQLite indeksleri için yüzde 20 ek pay olan herhangi bir makine. Milyonlarca kayıt indirmeyi planlıyorsanız, bellekten çok disk önemlidir: tüm yaklaşım durumun süreç değişkenlerinde değil, diskte yaşaması üzerine kuruludur.

Yedeklemeler

Aşağıda oluşturacağınız durum dosyası (örneklerde export.sqlite) tüm işin en değerli eseri olacak. Kodla herhangi bir deneyden önce onu kopyalama alışkanlığı edinin: cp export.sqlite export.sqlite.bak. Bu bir kez bir günlük indirmenizi kurtaracak.

Dikkat: indirici çalışırken SQLite dosyasını asla elle düzenlemeyin. Yanlış modda başka bir programdan okumak bile yazmayı bloklayıp süreci düşürebilir. Durumu görmek istiyorsanız indiriciyi durdurun veya paralellik bölümünde anlatacağımız WAL modunu kullanın.

Terimler sade bir dille

  • Kontrol noktası — diskte saklanan "buraya kadar her şey indirildi ve yazıldı" işareti. Kesintiden sonra indirici kontrol noktasını okur ve oradan devam eder.
  • İmleç (cursor) — API'nin sayfayla birlikte döndürdüğü ve bir sonraki sayfayı almak için iletilmesi gereken opak dize. İmleci kendiniz oluşturmaz veya çözmezsiniz.
  • Ofset ile sayfalı indirme — "37. sayfa, 500 kayıt" istediğiniz durum. Basit şema, ancak kaynağa yeni kayıtlar eklendiğinde sayfalar kayar ve kopyalar veya atlamalar oluşur.
  • Anahtar ile sayfalı indirme (keyset) — "kimliği 184203'ten büyük olan her şeyi, kimliğe göre sıralanmış" istediğiniz durum. Kaynak destekliyorsa devam ettirme için en sağlam şema.
  • İdempotency — bir işlemin tekrarlanmasının yapılmamış gibi aynı sonucu vermesi özelliği. Sayfayı iki kez yazdınız ama veritabanında bir kez duruyor.
  • Tekilleştirme anahtarı — iki kaydın aynı sayıldığı değer. Mükemmeli kaynaktan gelen kimliktir; yoksa anahtar kararlı alanların hash'i olarak hesaplanır.
  • Range başlığı — HTTP sunucusundan tüm dosyayı değil, bir kısmını, örneğin 1048576'dan sonuna kadar olan baytları isteme yöntemi.
  • At-least-once semantiği — her kaydın en az bir kez alınacağı garantisi. Tekrarlar olabilir ama atlamalar olmaz. Bu rehberden sonra elde edeceğiniz tam olarak budur, kopyaları tekilleştirme temizler.

Ana ilke

Aşağıdaki yedi adımın tamamı tek bir fikre indirgenir: her iş birimi atomik ve tekrarlanabilir olmalıdır. İş birimi ya bir dosya parçası, ya bir API sayfası, ya da kuyruktan bir görevdir. Atomik — yani sonuç ve tamamlandığına dair işaret birlikte kaydedilir. Tekrarlanabilir — yani bir birim iki kez yürütülürse hiçbir şey bozulmaz. Her iki özellik de sağlandığında, herhangi bir noktada kesilme güvenli hale gelir.

Adım 1: HTTP Range başlığı ile dosya devam ettirme

Bu aşamanın amacı: büyük bir dosyayı proxy üzerinden indirmeyi öğrenmek, öyle ki kesintiden sonra indirme sıfırdan değil, kaldığı bayttan devam etsin.

Nasıl çalışır

HTTP protokolü istemcinin kaynağın bir kısmını talep etmesine izin verir. Bunun için isteğe Range: bytes=BAŞLANGIÇ- başlığı eklenir, burada BAŞLANGIÇ bayt cinsinden ofsettir. Sunucu kısmi istekleri destekliyorsa 206 Partial Content kodu ve Content-Range: bytes BAŞLANGIÇ-SON/TOPLAM başlığı ile yanıt verir. Desteklemiyorsa Range'i yok sayar ve tüm dosyayı 200 kodu ile döndürür. Sizin göreviniz bu durumları ayırt etmektir.

Sunucunun devam ettirmeyi destekleyip desteklemediğini önceden öğrenmek için HEAD isteği yardımcı olur: sadece gövdesiz başlıkları döndürür. Accept-Ranges: bytes değerine bakın. none değeri veya başlığın olmaması genellikle devam ettirmenin olmadığı anlamına gelir, ancak bazı sunucular yine de Range'i doğru işler, bu yüzden son kontrolü yanıt koduna göre yaparız.

Adım adım talimatlar

  1. Proxy üzerinden bir HEAD isteği yapın ve Accept-Ranges, Content-Length, ETag ve Last-Modified başlıklarını kaydedin. ETag, dosyanın sunucuda denemeleriniz arasında değişip değişmediğini anlamak için gereklidir.
  2. Yerel dosyada zaten kaç bayt olduğuna bakın. Dosya yoksa sıfır kabul edin.
  3. Yerel boyut zaten Content-Length'e eşitse dosya tamdır, hiçbir şey yapmaya gerek yok.
  4. Yerel boyut sıfırdan büyükse ve sunucu aralık desteği bildiriyorsa, isteğe mevcut boyutla Range başlığını ekleyin. Ayrıca kaydedilen ETag ile If-Range ekleyin: böylece sunucu kısmi yanıtı yalnızca dosya değişmediyse döndürür, aksi halde tüm dosyayı 200 kodu ile verir.
  5. Mutlaka Accept-Encoding: identity gönderin. Aksi halde sunucu anlık sıkıştırma uygulayabilir ve bayt ofsetleri dosyanızla eşleşmez.
  6. 206 aldıysanız yerel dosyayı ab (ekleme) modunda, 200 aldıysanız wb (üzerine yazma) modunda açın.
  7. Gövdeyi 256 KB'lik parçalar halinde akış olarak okuyun ve diske yazın. Tüm yanıtı belleğe yüklemeyin.
  8. Tamamlandıktan sonra nihai boyutu Content-Length ile karşılaştırın. Eşleşmiyorsa bağlantı sessizce kopmuştur ve bir tur daha gerekir.

Çalışan kod

import os
import requests

PROXY_URL = os.environ['PROXY_URL'] # Proxeon kontrol panelinden alınan satır
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('dosya zaten tam:', have, 'bayt')
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('sunucu yanlış aralık verdi: ' + r.headers.get('Content-Range', ''))
mode = 'ab'
elif r.status_code == 200:
print('sunucu dosyayı tamamen veriyor, sıfırdan başlıyoruz')
mode = 'wb'
have = 0
elif r.status_code == 416:
raise IOError('istenen aralık dosya dışında, yerel boyutu kontrol edin')
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('kesinti: %d / %d alındı' % (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. tur kesildi: %s' % (i + 1, type(e).__name__))
# sonraki tura kadar bekleme: strateji 429 ve yeniden denemeler makalesinde açıklanmıştır
raise RuntimeError('dosya %d turda tamamlanamadı' % max_rounds)


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

download_until_done fonksiyonuna dikkat edin: denemeler arasındaki bekleme mantığını içermiyor. Bu kasıtlıdır. Oraya yeniden denemeler makalesindeki kendi duraklatma stratejinizi ekleyin, burada yalnızca "boyutu kontrol et, devam ettir, tekrar kontrol et" döngüsü önemlidir.

İpucu: dosya arşiv olarak sunuluyorsa, devam ettirme sırasında onu anında açmayın. Önce tam dosyayı alın, boyutu kontrol edin ve sunucu sağlama toplamı veriyorsa onu doğrulayın. Ancak ondan sonra açın. Kısmen indirilmiş bir gzip bozuk görünür ve var olmayan bir hatayı ararken zaman kaybedersiniz.

Beklenen sonuç

Kontrol: betiği en az 200 MB'lık bir dosyada çalıştırın, on saniye sonra Ctrl+C ile kesin. Yerel dosyanın boyutuna bakın, örneğin 41 943 040 bayt. Betiği tekrar çalıştırın. Konsolda "sıfırdan başlıyoruz" satırı görünmemeli ve dosya boyutu sıfırlanmak yerine büyümeye devam etmelidir. Tamamlandığında nihai boyut, HEAD isteğindeki Content-Length ile tam olarak eşleşmelidir.

Olası sorunlar

  • Sunucu her zaman 206 yerine 200 döndürüyor. Yani devam ettirme desteklenmiyor. Böyle bir kaynak için tek çıkış yolu, dosyayı tek seferde büyük bir zaman aşımıyla indirmek veya kaynağın tarihlere göre bölme gibi alternatif bir indirme biçimi aramaktır.
  • Content-Length başlığı yok. Sunucu dosyayı boyut bildirmeden chunked modda veriyor. Boyuta göre tamlığı kontrol edemezsiniz, devam da ettiremezsiniz: Range bilinen ofsetler gerektirir. Kaynakla anlaşın veya yayınlanıyorsa sağlama toplamı kullanın.
  • İlk turda 416 yanıtı. Yerel dosya sunucudaki dosyadan büyük. Sunucudaki dosya değişti ve kısaldı. Yerel dosyayı silin ve baştan başlayın.
  • Boyut eşleşti ama dosya bozuk. Muhtemelen ortada bir yerde 200 kodu ve sıfırdan yazma, ardından ekleme oldu. Dosyayı yeniden oluşturun. Bunun tekrarlanmaması için ETag'i yanında ayrı bir dosyada saklayın ve her turdan önce karşılaştırın.

Adım 2: Sayfalı indirmeler için kontrol noktaları

Bu aşamanın amacı: indirme durumunu, herhangi bir çökmeden sonra sürecin son başarıyla yazılan sayfadan devam edeceği şekilde kaydetmek.

Ne kaydedilmeli

Minimum kontrol noktası, kaynağın sayfalama türüne bağlıdır. Üç durumu ele alalım.

  1. İmleç ile sayfalama. API verilerle birlikte next_cursor gibi bir alan döndürür. Tam olarak onu kaydedin. Bu en basit durumdur: imleç zaten sunucunun devam etmek için ihtiyaç duyduğu her şeyi içerir.
  2. Sayfa numarası veya ofset ile sayfalama. Tamamen yazılan son sayfanın numarasını ve sayfa boyutunu kaydedin. İndirme sırasında kaynağa kayıt eklendiğinde ofsetlerin kayacağını unutmayın, bu nedenle 4. adımdaki tekilleştirme zorunludur.
  3. Anahtar ile sayfalama. Son yazılan kaydın kimliğini kaydedin. Devam ederken bu kimlikten büyük olan her şeyi istersiniz. Şema ne eklemelerden ne de uzun duraklamalardan korkar.

Sayfalama türünden bağımsız olarak kontrol noktasına yardımcı alanlar eklemeye değer: son kaydın kimliği (imleç şeması için bile, imleç geçersiz olursa yedek çapa olarak), ilerleme kontrolü için sayfa ve satır sayaçları, indirmenin başlangıç zamanı ve son güncelleme zamanı.

Durum nerede saklanmalı

İki çalışan seçenek var ve ikisi de bellekteki değişkenlerden daha iyidir.

Seçenek A: Atomik değiştirme ile JSON dosyası

Sonuçlar veritabanına değil, ayrı dosyalara yazılıyorsa uygundur. Ana tuzak: durumu doğrudan hedef dosyaya yazarsanız ve süreç yazma sırasında çökerse, okunamayan kesik bir JSON elde edersiniz. Çözüm: yanına geçici bir dosya yazın ve ana dosyanın üzerine yeniden adlandırın. Aynı dosya sistemindeki yeniden adlandırma işlemi atomiktir.

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)

Seçenek B: Verilerin yanında SQLite tablosu

Kayıtları bir veritabanına koyuyorsanız tercih edilen seçenek. Kontrol noktası, sayfa satırlarının eklenmesiyle aynı işlem içinde güncellenir. Ya hem veriler hem de işaret kaydedilir, ya da hiçbiri. "Veri var, işaret yok" arasında boşluk prensipte olmaz.

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: # sayfa başına bir işlem
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))

İşlem sırası

Kuralı hatırlayın: önce veriler, sonra kontrol noktası ve tercihen aynı işlem içinde. İşlem mümkün değilse (örneğin, veriler dosyalara yazılıyorsa), sıra tam olarak böyledir: sayfa dosyasını yazın, diske senkronize edin, sonra durumu güncelleyin. Bu iki eylem arasında bir çökme olursa, bir sayfanın yeniden yüklenmesini alırsınız ve bu 3. adım sayesinde güvenlidir. Ters sıra bir sayfanın atlanmasına yol açar ve bu zaten veri kaybıdır.

İpucu: kontrol noktasında mevcut imleci değil, sunucunun döndürdüğü bir sonraki sayfanın imlecini saklayın. Böylece devam ederken henüz sahip olmadığınız şeyi hemen talep edersiniz, zaten alınmış sayfa için ekstra istek olmaz.

Beklenen sonuç

Kontrol: indirmeyi başlatın, on sayfaya ulaşmasını bekleyin ve süreci zorla sonlandırın. Veritabanını sqlite3 export.sqlite komutuyla açın ve SELECT pages, last_id FROM checkpoint; çalıştırın. 10 sayısını ve kimliği görmelisiniz. Ardından SELECT count(*) FROM records; çalıştırın ve kayıt sayısının on sayfa boyutuna eşit olduğundan emin olun. İndiriciyi tekrar başlatın: konsoldaki ilk mesaj "başlangıç: sayfalar 10" gibi bir şey olmalıdır.

Olası sorunlar

  • "database is locked" hatası. Başka bir süreç bağlantıyı tutuyor. Dosyayı açan tüm sqlite3 pencerelerini ve diğer araçları kapatın. Çok iş parçacıklı çalışma için WAL modunu etkinleştirin, bkz. 5. adım.
  • İmleç kaydedildi ama veriler kaydedilmedi. Kontrol noktasını verilerle aynı işlem dışında güncellediniz. Yukarıdaki koda dönün ve her iki işlemin de aynı with con: bloğu içinde olduğundan emin olun.
  • JSON durumu boş veya bozuk çıktı. Geçici dosya ve değiştirme olmadan doğrudan dosyaya yazdınız. save_state fonksiyonunu bir bütün olarak kullanın.

Adım 3: Sonuçların idempotent yazımı

Bu aşamanın amacı: herhangi bir sayfanın yeniden işlenmesinin kopya üretmemesi ve verileri bozmaması.

Tekrar neden kaçınılmaz

2. adımdan sonra, bir sayfanın iki kez yazıldığı senaryoyu zaten gördünüz: süreç verileri ekledikten sonra ama kontrol noktasını güncellemeden önce çöktü. Bunun dışında, kaynak değiştiğinde ofset ile sayfalamadan, yeniden başlatmadan sonra aynı görevi alan paralel çalışanlardan ve sadece "her ihtimale karşı" yapılan manuel yeniden başlatmalardan tekrarlar gelir. İstek tarafında tekrarlarla savaşmak işe yaramaz. Doğrusu, yazmanın kendisini tekrarın zararsız olacağı şekilde yapmaktır.

Tekilleştirme anahtarının seçimi

  1. Kaynağın kimliği var. Onu kullanın. Bu id, uuid, order_number veya kaynağın benzersiz olduğunu garanti ettiği benzer bir alandır. Birden fazla kaynaktan tek bir tabloya indiriyorsanız, bileşik anahtar yapın: kaynak adı artı kimlik.
  2. Kimlik yok ama birlikte kaydı tanımlayan bir alanlar kümesi var. Örneğin, bir fiyat listesi satırı için bu, ürün kodu artı depo artı tarihtir. Bu alanlardan anahtar oluşturun ve normalleştirin: dizeleri tek bir büyük/küçük harfe getirin, baştaki ve sondaki boşlukları kaldırın, tarihleri tek bir biçime çevirin.
  3. Kararlı hiçbir şey yok. O zaman anahtar, kanonikleştirmeden sonra tüm kaydın hash'i olur. Bu durum 4. adımda ayrıntılı olarak ele alınmıştır. Kaynağın kaydı değiştirdiğini (fiyatı güncellediğini) unutmayın, hash değişir ve her iki sürümü de alırsınız. Bazen tam da istenen budur, bazen değildir.

Dikkat: anahtar olarak yanıttaki sıra numarasını veya sayfa numarasını kullanmayın. Bu değerler kaynaktaki herhangi bir değişiklikte değişir ve tekilleştirme kopya üreticisine dönüşür.

Veritabanına idempotent ekleme

SQLite ve çoğu ilişkisel veritabanında, birincil anahtar çakışmasını ya yok sayan ya da mevcut satırı güncelleyen bir yapı vardır. İlk varyant INSERT OR IGNORE'u 2. adımda zaten gördünüz. Kayıtlar değişmezse uygundur. İkinci varyant, kaynak kayıtları güncelleyebiliyorsa ve en güncel sürüme ihtiyacınız varsa gereklidir:

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

Dosyalara idempotent yazma

Sonuç veritabanında değil dosyalarda olacaksa, aynı ilkeyi uygulayın: bir sayfa, belirlenimci bir ada sahip bir dosyaya eşittir. Ad, sayfa parametrelerine bağlıdır, zamana veya sayaca değil. Yüklemeden önce nihai dosyanın var olup olmadığını kontrol edin; varsa sayfayı atlayın. Geçici bir ada yazın ve tamamlandığında yeniden adlandırın, tıpkı save_state fonksiyonunda olduğu gibi.

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 # sayfa zaten var, yeniden yazmaya gerek yok
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + '.part'
with open(tmp, 'w', encoding='utf-8') as f:
for r in rows:
f.write(json.dumps(r, ensure_ascii=False))
f.write(chr(10))
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
return True

Çökmeden sonra kalan .part uzantılı dosyalar başlangıçta güvenle silinebilir: tanım gereği eksiktirler.

İpucu: ofset ile indirmede yalnızca "dosya var, sayfa hazır"a güvenmeyin. Ek olarak dosyadaki satır sayısının sayfa boyutuna eşit olduğunu (sonuncu hariç) kontrol edin. Boş veya kısa bir dosyayı mevcut adla yeniden indirmek daha iyidir.

Beklenen sonuç

Kontrol: aynı sayfayı arka arkaya üç kez yazma fonksiyonunu çağırın. Ardından SELECT count(*) FROM records; çalıştırın. Sayı, üç katına değil bir sayfa boyutuna eşit olmalıdır. Dosya varyantı için dizinde tam olarak bir sayfa dosyası ve hiç .part dosyası olmamalıdır.

Adım 4: Belleği şişirmeden sonuçları tekilleştirme

Bu aşamanın amacı: milyonlarca satırlık bir akışta, tüm anahtarları RAM'de tutmadan yinelenen kayıtları elemek.

Kayıt hash'i

Bir kaydın kimliği yoksa, anahtar içeriğinin hash'i olur. Aynı kayıtların aynı hash'i vermesi için içerik kanonikleştirilmelidir: sözlük anahtarlarını sıralayın, fazla boşlukları kaldırın, ayırıcıları sabitleyin. Aksi halde, farklı alan sırasıyla gelen aynı kayıt farklı bir hash alır.

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

Fonksiyon 16 bayt döndürür. Bu yeterlidir: yüz milyonlarca kayıt için rastgele çakışma olasılığı ihmal edilebilir. fields parametresi, hash'i yalnızca kararlı alanlara göre hesaplamanıza olanak tanır; örneğin her istekte değişen son güncelleme zamanını hariç tutar.

Bellekteki küme neden milyonlarda çalışmaz

İlk akla gelen: seen = set() oluşturup anahtarları oraya koymak. Hesaplayalım. 16 bayt uzunluğunda bir bytes nesnesi Python'da yaklaşık 49 bayt artı verinin kendisi, toplamda yaklaşık 65 bayt kaplar. Doluluk faktörü hesaba katıldığında kümedeki bir yuva yaklaşık 30 bayt daha ekler. Anahtar başına yaklaşık 95 bayt elde ederiz. 10 milyon kayıtta bu yaklaşık 950 MB, 50 milyonda neredeyse 5 GB. Ve en önemlisi: süreç yeniden başladıktan sonra küme boş olur ve tüm tekilleştirme sıfırdan başlar.

Belleği şişirmemenin üç yolu

  1. Anahtarları veritabanının kendisinde saklayın. En basit ve en güvenilir yol. Anahtar records tablosunun birincil anahtarıysa, tekilleştirme zaten 3. adımdaki INSERT OR IGNORE yapısıyla yapılmıştır. İndeks diskte yaşar, yeniden başlatmaları atlatır ve SQLite sıcak indeks sayfalarını kendi önbelleğe alır. 10 milyon 16 baytlık anahtar için indeks diskte yaklaşık 400-500 MB yer kaplar, ancak bellekte değil.
  2. rowid'siz ayrı bir görülen anahtarlar tablosu. Verilerin kendisini SQLite'a değil, örneğin dosyalara yazıyorsanız gereklidir. O zaman SQLite yalnızca diskte kompakt bir küme olarak kullanılır.
  3. Ön filtre olarak bellekte sıkıştırılmış küme. İleri düzey varyant: hash'i 8 bayta kısaltın ve sıralı bir dizide tam sayı olarak saklayın veya bir Bloom filtresi kullanın. Bellek birkaç kat azalır, ancak yanlış pozitif olasılığı ortaya çıkar. Bu nedenle böyle bir ön filtre yalnızca kesinlikle yeni kayıtları hızlıca elemek için kullanılır ve nihai kontrol yine de veritabanından yapılır.

Diskteki kümenin uygulanması

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]
# sayfanın kendi içinde tekilleştirme
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])

Sayfayı yazmadan önce filter_new'i, veri yazımı ve kontrol noktasıyla aynı işlem içinde remember'ı çağırın. Böylece bir çökmeden sonra görülenler kümesi, veriler ve ilerleme işareti her zaman tutarlı olur.

İpucu: 500 değerli IN sorgusu indeks üzerinden milisaniyeler içinde yürütülür. Anahtarları döngüde tek tek kontrol etmeyin: her sorgunun ek yükü nedeniyle onlarca kat daha yavaştır.

Beklenen sonuç

Kontrol: 500 kayıtlık bir test sayfası oluşturun, 100'ü sayfa içinde iki kez tekrarlansın ve 100'ü de seen tablosunda zaten olsun. filter_new fonksiyonu tam olarak 300 kayıt döndürmelidir. Süreç yeniden başlatıldıktan sonra aynı 500 kayıt sıfır yeni vermelidir.

Olası sorunlar

  • Kopyalar hâlâ geçiyor. Kanonikleştirmeyi kontrol edin: muhtemelen kayıtlarda istek zamanı veya listedeki öğelerin rastgele sırası olan bir alan var. Bunu fields parametresiyle hariç tutun veya hash'lemeden önce iç içe listeleri sıralayın.
  • Birkaç milyon satırdan sonra yavaş ekleme. İndeks önbelleğe sığmayı bıraktı. SQLite önbelleğini PRAGMA cache_size=-200000 komutuyla artırın (bu 200 MB'tır) ve eklemelerin satır satır değil, sayfa başına bir işlemde toplu olarak yapıldığından emin olun.

Adım 5: Kayıpsız paralellik

Bu aşamanın amacı: proxy üzerinden birkaç eşzamanlı çalışanla indirmeyi hızlandırmak, öyle ki herhangi birinin çökmesi görevleri kaybetmesin ve veritabanını bozmasın.

Paralelleştirme ne zaman mümkün, ne zaman değil

İmleç ile sayfalama doğası gereği sıralıdır: bir sonraki imleç ancak önceki sayfa alındıktan sonra bilinir. Doğrudan paralelleştirmek imkansızdır. Ancak neredeyse her zaman indirmeyi bağımsız parçalara bölebilirsiniz: günlere, kategorilere, bölgelere, kimliğin ilk karakterlerine göre. Her parça kendi kontrol noktasıyla sıralı olarak indirilir ve parçalar paralel ilerler. Ofset ve bilinen sınırlarla anahtar ile sayfalama doğrudan paralelleştirilir: "1'den 100'e sayfalar" veya "0'dan 100000'e kimlikler" gibi görevler.

Diskte iş kuyruğu

Bellekteki kuyruk süreçle birlikte ölür. Bu nedenle görevler durum içeren bir tabloda yaşar:

  • pending — yürütülmeyi bekliyor;
  • running — bir çalışan tarafından alındı;
  • done — yürütüldü ve kaydedildi;
  • failed — denemeler tükendi, insan dikkati gerektiriyor.

Başlangıçta indirici önce tüm running görevleri pending'e geri çevirir: bu durumda asılı kalmışlarsa, önceki süreç çalışma ortasında ölmüştür. Sonra çalışanlar pending'leri alır.

Tek yazma noktası

SQLite birçok eşzamanlı okuyucuya izin verir, ancak yalnızca bir yazıcıya. En basit ve güvenli kalıp: çalışanlar yalnızca indirir ve veri döndürür, tüm veritabanı yazımını ana iş parçacığı yapar. Kodda kilit yok, "database is locked" yok. Ek olarak, başka bir süreçten durumu okumanın yazmayı engellememesi için WAL modunu etkinleştirin.

Eşzamanlılık sınırlaması

Çalışan sayısını iki şeyle sınırlayın. Birincisi — proxy'nin olanakları: Proxeon kontrol panelinde birkaç kanalınız varsa, kanal başına bir-iki çalışan tutmak mantıklıdır, böylece bir kanaldaki IP rotasyonu tüm iş parçacıklarının bağlantısını aynı anda koparmaz. İkincisi — kaynağa nezaket: resmi sınırlar olmasa bile, küçük bir API'ye on paralel iş parçacığı, reddedilmeler alacağınız bir yük oluşturur. Üç-dört çalışanla başlayın ve hata oranını gözlemleyerek artırın.

Paralel indirici kodu

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: # veriler ve görev durumu aynı işlemde
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('kuyruk boş, hatalı görev sayısı:', failed)

fetch_fn fonksiyonu proxy üzerinden isteği yürütür ve kayıt listesi döndürür. İş parçacığında çalışır ve veritabanına dokunmaz. write_fn fonksiyonu ana iş parçacığında bir işlem içinde çağrılır ve 3. adımdaki idempotent eklemeyi yapar. Başarısız görev otomatik olarak pending'e döner ve bir sonraki claim döngüsünde tekrar alınır; denemeler tükendikten sonra failed durumunu alır ve onunla manuel olarak ilgilenirsiniz.

Dikkat: sqlite3 bağlantı nesnesini çalışanlara geçirmeyin. Bağlantı oluşturulduğu iş parçacığına bağlıdır ve başka bir iş parçacığından kullanmaya çalışmak hataya veya daha kötüsü sessiz veri bozulmasına yol açar. Ya her iş parçacığına kendi bağlantısı, ya da yukarıdaki örnekte olduğu gibi hiçbiri.

İpucu: her çalışan için kendi proxy adresiyle ayrı bir requests.Session nesnesi kullanın. Proxeon'da birkaç kanalınız varsa, bunları çalışanlara sırayla dağıtın: çalışan 0 kanal 0'ı, çalışan 1 kanal 1'i alır ve böyle devam eder. Böylece bir kanaldaki bağlantı kopması yalnızca bir iş parçacığını etkiler.

Beklenen sonuç

Kontrol: kuyruğa 100 görev koyun, dört çalışan başlatın ve yarım dakika sonra süreci öldürün. SELECT status, count(*) FROM tasks GROUP BY status; çalıştırın. Birkaç done, birkaç running ve geri kalanı pending göreceksiniz. Tekrar başlatın: başlangıçta running'ler kaybolmalı ve tamamlandığında dürüstçe başarısız olanlar ve last_error'da hata metniyle failed'da duranlar dışında tüm görevler done olmalıdır.

Adım 6: Uzun bir aradan sonra devam ettirme

Bu aşamanın amacı: birkaç saat veya gün durdurulan bir indirmeyi doğru şekilde sürdürmek ve bayatlamış duruma çarpmamak.

Neler bayatlar

On saniye sonra devam ettirmek ile bir hafta sonra devam ettirmek farklı görevlerdir. Uzun bir aradan sonra kaydedilen durumun bir kısmı geçerliliğini yitirir.

  1. Oturumlar ve çerezler. Sunucu oturumları genellikle birkaç saatten bir güne kadar yaşar. Bu süreden sonra kaydedilen çerezler 401 yanıtlarına veya giriş formuna yönlendirmeye yol açar. Çözüm: başlangıçta çerezleri dosyadan geri yüklemek yerine tam bir yeniden yetkilendirme yapmak.
  2. Erişim belirteçleri. OAuth belirteçleri bir saat, bazen daha az yaşar. Bir yenileme belirteciniz varsa, başlangıçtan önce ve çalışma sırasında programlı olarak erişim belirtecini yenileyin, reddedilmeyi beklemeden.
  3. Sayfalama imleçleri. Birçok API imlecin yaşam süresini dakikalar veya saatlerle sınırlar. Bayat imleç, geçersiz imleç mesajıyla 400 hatası döndürür. Bu yüzden 2. adımda yedek çapa kaydettik: son kaydın kimliği. Kaynak kimliğe veya değişiklik tarihine göre filtrelemeyi destekliyorsa, yeni isteği bu çapadan oluşturun. Desteklemiyorsa, parçayı baştan başlatmanız gerekir ve 4. adımdaki tekilleştirme zaten alınanları eler.
  4. Sunucudaki dosya içeriği. 1. adımdaki devam ettirme için dosyanın değişmemesi kritiktir. Her turdan önce mevcut ETag'i kaydedilenle karşılaştırın; uyuşmazlık durumunda dosyayı baştan başlatın.
  5. Proxy ayarları. Bir hafta içinde Proxeon kontrol panelinde port, şifre değişmiş veya kanalın süresi dolmuş olabilir. Kuyruğu çözmeye başlamadan önce proxy'yi bir test isteğiyle kontrol edin.
  6. Veri kümesinin kendisi. İndirme bir hafta sürüyorsa ve kaynak bu sürede kayıt ekleyip sildiyse, sonucunuz farklı zaman noktalarındaki durumların bir karışımı olur. Birçok görev için bu kabul edilebilir. Değilse, kontrol noktasında başlangıç zamanını saklayın ve tamamlandıktan sonra bu zamandan sonra değişen kayıtlar için ayrı bir artımlı geçiş yapın.

Uçuş öncesi kontrol

Tüm kontrolleri, başlangıçta herhangi bir gerçek işten önce yürütülen tek bir fonksiyonda toplayın. Ya durumu düzene sokar ya da indiriciyi anlaşılır bir mesajla durdurur.

def preflight(session, state, probe_url):
# 1. proxy canlı ve yetkili
r = session.head(probe_url, timeout=20)
if r.status_code == 407:
raise SystemExit('proxy kullanıcı adı veya şifreyi reddetti, Proxeon panelindeki bilgileri kontrol edin')
# 2. erişim belirteci taze
refresh_access_token(session)
# 3. imleç hâlâ geçerli
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('imleç bayatladı, last_id =', state['last_id'], 'çapasına geçiliyor')
state['cursor'] = None
state['resume_after_id'] = state['last_id']
# 4. indirme yaşı hatırlatması
print('indirme', state['started_at'], 'tarihinde başladı, kaydedilen sayfa sayısı', state['pages'])
return state

refresh_access_token fonksiyonu kaynağınıza bağlıdır: genellikle yenileme belirteciyle bir POST isteğidir ve ardından oturumdaki Authorization başlığını güncellersiniz. resume_after_id alanı daha sonra sayfa alma fonksiyonunda "kimlik belirtilenden büyük" filtresi olarak kullanılır.

İpucu: yenileme belirtecini ve proxy şifresini kontrol noktasında değil, ortam değişkenlerinde veya sınırlı erişim izinlerine sahip ayrı bir sırlar dosyasında saklayın. Kontrol noktasını kopyalayacak, meslektaşlarınıza iletecek ve hata raporlarına ekleyeceksiniz; sırlar orada fazlalıktır.

Beklenen sonuç

Kontrol: checkpoint tablosundaki imleci UPDATE checkpoint SET cursor='broken'; komutuyla elle bozun ve indiriciyi başlatın. Konsolda last_id çapasına geçiş satırı görünmeli ve indirme çökmeden devam etmelidir. Tamamlandıktan sonraki kayıt sayısı, imleç bozulmadan yapılan kontrol çalıştırmasıyla eşleşmelidir.

Adım 7: Python'da hazır sağlam indirici iskeleti

Bu aşamanın amacı: önceki adımlardaki her şeyi, çalıştırılabilecek, kesilebilecek, tekrar çalıştırılabilecek ve kopyasız tam sonuç elde edilebilecek tek bir dosyada birleştirmek.

İskeletin yapısı

  • Ortam değişkenlerinden yapılandırma: Proxeon proxy adresi, API adresi, belirteç, görev adı.
  • Store sınıfı: kayıtlar ve kontrol noktası tablolarına sahip SQLite, sayfa başına bir işlem.
  • İdempotency için kayıt anahtarı fonksiyonu.
  • Proxy üzerinden sayfa alma fonksiyonu.
  • Kontrol noktasından devam etme ve kesintiden sonra oturumu yeniden oluşturma ile ana döngü.

Tam 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']# Proxeon panelinden http://USER:PASS@HOST:PORT
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('başlangıç: sayfalar %d, satırlar %d, veritabanındaki benzersiz %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('denemeler tükendi, durum kaydedildi, daha sonra tekrar başlatın')
raise
print('kesinti (%s), deneme %d / %d' % (type(e).__name__, attempts, MAX_ATTEMPTS))
time.sleep(min(60, 2 ** attempts)) # duraklama seçimi 429 ve yeniden denemeler makalesinde açıklanmıştır
session = make_session() # yeni oturum: proxy üzerinden bağlantı yeniden oluşturulur
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('sayfalar %d, satırlar %d' % (state['pages'], state['rows']))
if next_cursor is None:
break
print('tamamlandı: sayfalar %d, alınan satırlar %d, veritabanındaki benzersiz %d'
 % (state['pages'], state['rows'], store.unique_count()))


if __name__ == '__main__':
run()

Kendi kaynağınıza nasıl uyarlanır

  1. /orders yolunu ve items, next_cursor, id alan adlarını API'nizin döndürdükleriyle değiştirin. Bu, fetch_page ve record_key fonksiyonlarındaki üç yerdir.
  2. Kaynakta ofset ile sayfalama varsa, cursor parametresini page ile değiştirin ve bir sonraki değeri state['pages'] + 1 olarak hesaplayın. Kontrol noktasında imleç yerine sayfa numarasını saklayın.
  3. Anahtar ile sayfalama varsa, state['last_id']'den after_id gibi bir parametre iletin ve imleçle çalışmayı kaldırın.
  4. Paralellik gerekiyorsa, fetch_page'i 5. adımdaki fetch_fn'e taşıyın ve commit_page'i write_fn olarak kullanın. İndirmeyi parçalara bölün ve görev kuyruğunu doldurun.
  5. Ana döngüden önce 6. adımdaki preflight fonksiyonunu ekleyin.

Sonucu kontrol etme: kontrol listesi

İndiriciyi gerçek çok saatlik bir hacimde çalıştırmadan önce bu listeden geçirin. Her madde birkaç dakika alır ve birlikte gece indirmesinin veri kaybetmeyeceğini garanti ederler.

  1. Kesme testi. İndiriciyi başlatın, 30 saniye sonra Ctrl+C'ye basın. Tekrar başlatın. Çıktının ilk satırı "sayfalar 0" değil, sıfırdan farklı bir sayı göstermelidir.
  2. Kopya testi. PAGE_SIZE'ı 10'a düşürün, indiriciyi rastgele anlarda arka arkaya beş kez kesin. Tamamlandığında benzersiz kayıt sayısını alınan satır sayısıyla karşılaştırın: benzersiz olan, kaynağa ekleme yapılmadığı sürece neredeyse eşit olmalıdır.
  3. Tutarlılık testi. Herhangi bir kesintiden sonra iki sorgu çalıştırın: SELECT rows FROM checkpoint; ve SELECT count(*) FROM records;. Aralarındaki fark bir sayfa boyutunu geçmemelidir. Geçiyorsa, kontrol noktası ve veriler aynı işlemde yazılmıyor.
  4. Proxy testi. Geçici olarak PROXY_URL değişkenini devre dışı bırakın veya yanlış bir şifre verin. İndirici ilk istekte anlaşılır bir hatayla çökmeli, askıda kalmamalı ve doğrudan kaynağa gitmeye başlamamalıdır.
  5. Bayat imleç testi. 6. adımda açıklandığı gibi veritabanındaki imleci bozun ve çapaya geçişin tetiklendiğinden emin olun.
  6. Disk testi. Bin sayfadan sonra export.sqlite dosyasının boyutunu kontrol edin ve beklenen sayfa sayısıyla çarpın. Diskte yüzde 20 payla yeterli yer olduğundan emin olun.

Kontrol: Başarılı yürütmenin göstergesi, üç kasıtlı kesintiden ve üç yeniden başlatmadan sonra nihai benzersiz kayıt sayısının aynı kaynakta tek bir kesintisiz çalıştırmada elde edilenle eşleşmesi ve konsolda hiçbir zaman sıfırdan başlama satırının görünmemesidir.

Ek özellikler ve optimizasyon

  • İlerleme ve süre tahmini. Toplam kayıt sayısı biliniyorsa, her yirmi sayfada bir yüzde ve kalan süre tahmini çıkarın. Bu hem sizin için hem de takılmayı yavaş çalışmadan ayırt etmek için faydalıdır.
  • Payload sıkıştırma. On milyonlarca kayıt için metin biçimindeki JSON çok yer kaplar. Payload alanını yazmadan önce zlib.compress ile sıkıştırın ve BLOB olarak saklayın. Tasarruf genellikle üç ila sekiz kat arasındadır.
  • Sunucu veritabanına geçiş. Kontrol noktasını verilerle aynı işlemde tutma şeması PostgreSQL'e neredeyse değişiklik yapmadan taşınır. ON CONFLICT yapısı orada da desteklenir ve "tek yazıcı" kısıtlaması ortadan kalkar.
  • Artımlı indirmeler. Her görevin başlangıç zamanını kaydedin ve tam indirmeden sonra "şu tarihten sonra değişen" filtresiyle ayrı bir görev çalıştırın. Böylece tam yeniden indirme olmadan güncel bir kopya tutarsınız.
  • Parça başına ayrı süreç. İş parçacıkları yerine, farklı JOB değerleri ve farklı DB_PATH dosyalarıyla birkaç betik örneği başlatabilir ve sonuçları sonda birleştirebilirsiniz. Hata ayıklaması daha kolaydır ve eşzamanlı yazma sorununu tamamen ortadan kaldırır.
  • Kesinti metrikleri. Her kesintiyi istisna türü ve zamanla günlüğe kaydedin. Bir gün sonra kesintilerin Proxeon kanalındaki IP rotasyon aralıkları etrafında toplandığını görecek ve rotasyon aralığını isteklerinizin süresine göre ayarlayabileceksiniz.

Tipik hatalar ve çözümleri

Aşağıda ilk çalıştırmalarda neredeyse herkesin karşılaştığı durumlar toplanmıştır. Biçim: sorun, neden, çözüm.

  1. Sorun: yeniden başlatmadan sonra indirme her seferinde sıfırdan başlıyor. Neden: kontrol noktası belleğe veya çökmeyi atlatmayan bir dosyaya yazılıyor ya da indirici başlangıçta onu okumuyor. Çözüm: run fonksiyonundaki ilk eylemin store.load olduğundan ve durumun her sayfadan sonra işlem içinde güncellendiğinden emin olun.
  2. Sorun: veritabanında kaynaktan bir buçuk kat fazla kayıt var. Neden: tekilleştirme anahtarı kararsız: içine istek zamanı, sayfa numarası veya rastgele sıralı bir alan girmiş. Çözüm: anahtarı yalnızca kaynağın kimliğine veya fields parametresiyle açık kararlı alanlar listesine göre hesaplayın.
  3. Sorun: indirme hatasız tamamlanmasına rağmen veritabanında kaynaktan daha az kayıt var. Neden: kontrol noktası verilerden önce güncellendi ve çökmeden sonra sayfa atlandı. Veya kaynaktaki silmeler ofset ile sayfalamada sayfaları geriye kaydırdı. Çözüm: sırayı "önce veriler, sonra kontrol noktası" olarak aynı işlem içinde değiştirin; silmelerin olduğu kaynaklar için anahtar ile sayfalamaya geçin.
  4. Sorun: dosya devam ettirme, boyut eşleşirken bozuk arşiv veriyor. Neden: sunucu bir kez 206 yerine 200 yanıtladı, dosya kısmen üzerine yazıldı ve sonra eklendi. Çözüm: ETag'i dosyanın yanında saklayın, değiştiğinde dosyayı silin; Content-Range başlığının istenen ofsetle eşleştiğini kontrol edin.
  5. Sorun: paralel çalışmada "database is locked" hatası. Neden: birkaç iş parçacığı aynı anda SQLite'a yazıyor veya bağlantı iş parçacıkları arasında geçirilmiş. Çözüm: ana iş parçacığında tek yazma noktası, çalışanlar yalnızca indirir; WAL modu; bağlantı kullanıldığı iş parçacığında oluşturulur.
  6. Sorun: bir saatlik çalışmadan sonra tüm istekler 401 döndürmeye başladı. Neden: erişim belirtecinin süresi doldu. Çözüm: belirteci sona ermeden önce programlı olarak yenileyin ve 401 alındığında yenilemeyi çağırın ve isteği bir kez tekrarlayın, bunu kesinti saymayın.
  7. Sorun: RAM birkaç gigabayta kadar büyüyor. Neden: görülen anahtarlar kümesi veya tüm kayıtların listesi süreç belleğinde tutuluyor. Çözüm: veritabanındaki birincil anahtar veya seen tablosu üzerinden tekilleştirme; verileri biriktirmeden sayfa sayfa yazma.
  8. Sorun: kesintiler tam olarak birkaç dakikada bir oluyor. Neden: proxy kanalındaki IP rotasyon aralığıyla çakışıyor. Çözüm: bu normal bir durumdur, indirici bunu atlatmalıdır. İstekler uzunsa, Proxeon kontrol panelindeki rotasyon aralığını tipik bir istek süresinden belirgin şekilde daha uzun olacak şekilde ayarlayın veya sayfalar arasında istek üzerine rotasyon kullanın.

SSS: sağlam indirme hakkında sık sorulan sorular

Sonuç CSV olarak gerekiyorsa SQLite kullanmak zorunlu mu?

Hayır, ama kullanışlı. SQLite burada güvenilir bir durum deposu ve görülen anahtarlar kümesi rolünü üstlenir. Nihai CSV'yi tamamlandıktan sonra records tablosundan tek bir komutla dışa aktarırsınız. Tamamen veritabanı olmadan tercih ederseniz, 3. adımdaki dosya varyantını sayfa başına bir dosya ve atomik değiştirme ile JSON kontrol noktasıyla kullanın.

Kontrol noktası ne sıklıkla kaydedilmeli: her sayfadan sonra mı yoksa daha seyrek mi?

Her sayfadan sonra. Birkaç yüz satırlık tek bir SQLite işlemi milisaniyeler içinde yürütülür, bu proxy üzerinden bir ağ isteğine kıyasla önemsizdir. Seyrek kontrol noktalarından tasarruf, onlarca sayfayı kaybetme riskine değmez.

API ne imleç, ne kimlik, sadece sayfa numaraları veriyorsa ne yapmalı?

Sayfa numarasıyla çalışın, onu kontrol noktasında saklayın ve mutlaka kayıt içeriğinin hash'i ile tekilleştirmeyi ekleyin. Kaynaktaki aktif değişikliklerde sayfa kaymaları nedeniyle bazı kayıtların atlanabileceğini kabul edin. Kritik veriler için sayfaları ters sırayla ikinci bir geçiş yapın: ilk geçişte atlananlar büyük olasılıkla ikinciye yakalanır.

Bir dosyayı birden fazla iş parçacığıyla farklı aralıklarla devam ettirebilir miyim?

Sunucu Range'i destekliyorsa evet. Dosyayı 50-100 MB'lık parçalara bölün, her parça 5. adımdaki kuyruktan kendi geçici dosyasıyla bir görev olsun ve tüm parçalar tamamlandığında doğru sırayla birleştirin. Her parçayı boyuta göre, tüm dosyayı varsa sağlama toplamına göre doğrulayın.

Proxy üzerinden çalışırken kaç çalışan koymalı?

Proxeon'da kanal başına üç-dört ile başlayın ve tasks tablosundaki hata oranını izleyin. Hata oranı yüzde birin altındaysa iki tane daha ekleyin. Hatalar artıyorsa azaltın. Kanal başına ondan fazla iş parçacığı nadiren kazanç sağlar: ya kanalın bant genişliğine ya da kaynağın sabrına takılırsınız.

Oturum çerezlerini çalıştırmalar arasında saklamalı mıyım?

Genellikle hayır. Başlangıçta yeniden yetkilendirme saniyeler alır ve bilinmeyen yaşam süresine sahip çerezleri geri yüklemekten daha güvenilirdir. İstisna: kaynak günlük giriş sayısını sınırlıyorsa. O zaman çerezleri saklayın, ancak ilk 401 veya girişe yönlendirmede onları atın ve yeniden yetkilendirin.

İndirmenin tamamen tamamlandığını, sessizce kesilmediğini nasıl anlarım?

Dosyalar için: boyut Content-Length'e eşit ve sağlama toplamı uyuşuyor. API için: next_cursor olmayan bir sayfa veya boş sayfa alındı ve kaynak toplam sayıyı bildiriyorsa satır sayısı buna karşılık geliyor. Kontrol noktasına açık bir tamamlanma bayrağı yazın, böylece tekrar başlatma yeni bir tur başlatmaz.

failed durumundaki görevlerle ne yapmalı?

last_error alanına bakın. Ağ hatalarıysa, görevleri UPDATE komutuyla pending'e döndürün ve indiriciyi tekrar başlatın. Veri ayrıştırma hatalarıysa, kaynakta standart dışı biçimde kayıtlar var demektir: kodu düzeltin ve yeniden başlatın. failed'ları asla sessizce silmeyin, indirmede neyin eksik olduğunun tek kanıtı onlardır.

Bu yaklaşımı requests yerine asenkron bir kütüphaneyle kullanabilir miyim?

Evet, ilkeler aynı: atomik iş birimi, verilerle birlikte kontrol noktası, idempotent yazma, diskte kuyruk. Yalnızca taşıma değişir. Tek nüans: SQLite yazımını senkron ve sıralı bırakın, paralelliği ağ istekleri düzeyinde tutun.

Sonuç

Eksik bir dosyadan ve sinir bozucu yeniden başlatmalardan, kesintinin umurunda olmadığı bir indiriciye kadar geldiniz. Tam olarak ne yapıldığını özetleyelim.

  • HTTP ile devam ettirmeyi ele aldık: HEAD ile Accept-Ranges kontrolü, mevcut dosya boyutuyla Range başlığı, 206 ve 200 kodlarının ayırt edilmesi, ETag ve If-Range ile dosya değişimine karşı koruma.
  • Sayfalı indirmeler için kontrol noktaları oluşturduk: imleç, sayfa numarası veya son kaydın kimliği artı yardımcı sayaçlar, hepsi verilerle aynı işlemde.
  • Yazımı birincil anahtar ve INSERT OR IGNORE veya ON CONFLICT DO UPDATE yapısıyla idempotent hale getirdik; dosyalar için belirlenimci adlar ve atomik değiştirme kullandık.
  • Milyonlarca anahtarın RAM'de yaşamaması ve yeniden başlatmaları atlatması için diskte tekilleştirme düzenledik.
  • SQLite'ta iş kuyruğu, başarısız görevlerin otomatik iadesi ve tek yazma noktası ile paralellik ekledik.
  • Uzun bir aradan sonra devam ettirmeyi öngördük: belirteç yenileme, yeniden yetkilendirme, bayat imleçten kimliğe göre çapaya geçiş, başlamadan önce Proxeon proxy kontrolü.
  • Her şeyi, üç-dört satır değiştirerek belirli bir kaynağa uyarlanabilen tek bir çalışan iskelette topladık.

Bundan sonra ne yapmalı

7. adımdaki iskeleti alın ve küçük bir gerçek hacimde, örneğin on bin kayıtta çalıştırın. Kontrol bölümündeki kontrol listesinden geçirin. Ancak ondan sonra gece boyunca tam indirmeyi başlatın. Sabah ya sayaçları eşleşen "tamamlandı" satırını ya da durumu kaydedilmiş "denemeler tükendi" satırını göreceksiniz, o zaman betiği tekrar çalıştırmanız yeterli.

Nereye ilerlemeli

Sonraki seviye, tam geçişler yerine değişiklik zamanına göre artımlı indirmeler, durumu birkaç makine için sunucu veritabanına taşıma ve yanıt kodlarını hesaba katan anlamlı bir yeniden deneme stratejisidir; bu konuya 429 ve yeniden denemeler hakkındaki ayrı bir makale ayrılmıştır. O makaledeki isabetli yeniden denemelerin bu makaledeki sağlam durumla birleşimi, bir hafta çalışmaya bırakıp terminali açmamanız gereken bir indirici verir.

Ve son olarak. Uzun bir indirmenin kesilmesi bir kaza değil, artık nasıl başa çıkacağınızı bildiğiniz bir çalışma durumudur. İyi indirmeler.