如何通过代理稳定导出海量数据:断线后不用从头重来,分步指南
引言:为什么长时间导出几乎总会中断,而且这是正常的
如果你曾经通过代理导出过几百万条记录,或者下载过几十 GB 的文件,你一定熟悉这种感受:脚本跑了六个小时,进度显示 83%,然后突然因为连接错误崩溃。你手头只剩下一个不完整的文件,以及“一切都得重来”的绝望。
首先要接受一个事实:长时间导出总会中断。不是“偶尔”,不是“网络差的时候”,而是只要它跑得足够久,必然会断。原因有几十种,而且大多数你都控制不了:
- 代理切换了外部 IP。对于 Proxeon 的移动代理来说,这是正常行为:按时间或按请求轮换。IP 切换的瞬间,已打开的 TCP 连接会被断开,源站看到的已经是一个不同的客户端。
- 源站按自己的超时关闭连接、重启、上线新版本,或者干脆返回 5xx 错误。
- 你自己的进程重启:系统更新、磁盘写满、代码在不常见的数据记录上报错、不小心按了 Ctrl+C。
- 授权 token 过期、会话失效、分页游标失效。
- 笔记本进入睡眠,Wi-Fi 切换到另一个热点,运营商切换了路由。
逐个对抗这些原因毫无意义。正确的做法是换一种思路:把导出设计成“任意时刻断线,代价不是六个小时,而只是一页数据或一个文件片段”。这篇指南就是讲这个的。
读完后你能得到什么
跟着本教程走完,你将获得:
- 一个可用的函数:通过 HTTP 的 Range 请求头从文件中间断点续传,并且能判断服务器是否支持。
- 分页导出的清晰检查点方案:保存哪些字段、把状态放哪里,才能不随进程一起丢失。
- 幂等写入结果:重复下载同一页不会产生重复数据。
- 去重方案:处理数百万行时不会吃光内存。
- 带任务队列的并发下载器,支持失败重试和并发上限。
- 一份可直接使用的 Python 稳健下载器骨架,你花一个小时就能改成适配自己的数据源。
适合谁阅读
适合已经会用 Python 发 HTTP 请求、并且至少经历过一次长导出结果全部丢失的开发者和分析师。难度中等:基础概念会解释,但不会从零教你编程。进阶读者可以关注“把哈希存在内存之外”和“安全并发写 SQLite”两部分。
需要提前掌握的知识
- Python 的函数、循环、字典和异常处理。
- HTTP 基础:方法、请求头、状态码、响应体。
- 大致了解如何在 requests 里设置代理。
特别说明:响应状态码和指数退避重试策略这里不展开。那部分有专门讲 429 错误和重试的文章。本指南的重点不同:关注导出状态和如何恢复。重试回答的是“什么时候该重发请求”,我们回答的是“重试全部用尽、进程死掉之后,从哪里继续”。
需要多少时间
阅读并跑通测试数据源示例:两到三个小时。把骨架改成适配你实际 API 或文件服务器:再花一到两个小时,取决于分页方式有多不标准。算下来一整个工作日绰绰有余。
准备工作与基础概念
工具与权限
- 安装 Python 3.11 或更高版本。2026 年主流是 3.12 和 3.13,本文所有示例都在这些版本上验证过。检查版本:打开终端输入
python --version。如果显示 3.11 或更高,就没问题。 - 安装 requests 库:
pip install requests。2.32 及以上版本即可。sqlite3 模块是 Python 标准库自带的,不需要单独安装。 - 获取代理访问权限。打开 Proxeon 个人中心,选择需要的通道,复制四个值:主机、端口、用户名、密码。通常它们会拼成一行,形如
http://USER:PASS@HOST:PORT。后面到处都会用到这行字符串。 - 把代理字符串放进环境变量,而不是代码里。Linux 和 macOS 下:
export PROXY_URL=http://USER:PASS@HOST:PORT。Windows PowerShell 下:$env:PROXY_URL='http://USER:PASS@HOST:PORT'。这样你不会不小心把密码提交到代码仓库。 - 确认代理能通。在终端执行:
curl -x $PROXY_URL -I https://api.example.com/,把地址换成你的数据源。你应该看到一行响应状态码,例如HTTP/2 200。如果看到 407 代理认证错误,请重新检查用户名和密码。
系统要求
任何有 2 GB 空闲内存的机器都可以,磁盘要能放下导出结果加上 20% 的余量给 SQLite 索引。如果你打算导出几百万条记录,磁盘比内存更重要:整个方案的思路是把状态放在磁盘上,而不是进程的变量里。
备份
你接下来要创建的状态文件(示例中是 export.sqlite)会成为整个工作里最宝贵的产物。养成习惯:在对代码做任何实验之前先复制一份,cp export.sqlite export.sqlite.bak。只要用上一次,就能救回你一整天的工作。
注意:下载器运行期间,绝不要手动编辑 SQLite 文件。即使用外部程序以错误的模式读取,也可能锁住写入并让进程崩溃。如果需要查看状态,请先停止下载器,或者使用我们后面在并发部分会讲到的 WAL 模式。
关键术语大白话
- 检查点(Checkpoint) —— 保存到磁盘的标记,表示“到这里为止都已导出并写入”。断线后下载器读取检查点,从那里继续。
- 游标(Cursor) —— API 随页返回的不透明字符串,下一次请求要带上它才能拿到下一页。你不需要自己构造或解析它。
- 按偏移量分页 —— 请求“第 37 页,每页 500 条”。简单,但源数据新增时页码会偏移,导致重复或漏数据。
- 按键分页(keyset) —— 请求“ID 大于 184203 的全部记录,按 ID 排序”。如果数据源支持,这是最稳健的续传方案。
- 幂等(Idempotency) —— 一个操作重复执行和只执行一次结果相同。同一页写两次,数据库里也只有一份。
- 去重键 —— 用来判断两条记录是同一条的值。最好直接用源的 ID;如果没有,就用稳定字段的哈希。
- Range 请求头 —— 请求 HTTP 服务器只返回文件的一部分,比如从第 1048576 字节到末尾。
- at-least-once 语义 —— 保证每条记录至少被获取一次。可能有重复,但不会有遗漏。读完本指南你得到的就是这个语义,重复部分由去重处理。
核心原则
下面七个步骤归根结底就是一句话:每个工作单元都必须是原子的、可重复的。工作单元可以是一个文件片段、一页 API 数据,或者队列里的一个任务。原子性意味着结果和完成标记一起保存。可重复意味着同一个单元做两次也不会坏。两个条件都满足,任何位置断线都是安全的。
第 1 步:通过 HTTP Range 请求头断点续传文件
本步目标:学会通过代理下载大文件,断线后从断掉的那个字节继续,而不是从零开始。
原理
HTTP 协议允许客户端请求资源的一部分。做法是加一个请求头 Range: bytes=起始-,其中“起始”是字节偏移量。如果服务器支持分片请求,会返回 206 Partial Content,并带一个响应头 Content-Range: bytes 起始-结束/总大小。如果不支持,它会忽略 Range,返回完整的文件,状态码 200。你需要区分这两种情况。
提前了解服务器是否支持续传,可以用 HEAD 请求:它只返回响应头不带正文。看 Accept-Ranges: bytes。值为 none 或没有这个头,通常表示不支持续传,不过有些服务器虽然不声明也能正确处理 Range,所以最终判断要看响应状态码。
分步操作
- 通过代理发一个 HEAD 请求,保存 Accept-Ranges、Content-Length、ETag 和 Last-Modified 这几个头。ETag 后面用来判断文件在两次尝试之间有没有被改动。
- 看看本地文件已经有多少字节。文件不存在就算 0。
- 如果本地大小已经等于 Content-Length,文件是完整的,什么都不用做。
- 如果本地大小大于 0 且服务器声明支持范围请求,就在请求里加上 Range 头,值为当前大小。同时加上带保存的 ETag 的
If-Range头:这样只有当文件没改过,服务器才会返回分片响应;否则返回完整文件,状态码 200。 - 务必发送
Accept-Encoding: identity。没有它,服务器可能实时压缩,字节偏移量就跟你本地文件对不上了。 - 收到 206 时以追加模式
ab打开本地文件,收到 200 时以覆盖模式wb打开。 - 按 256 KB 的块流式读取正文并写入磁盘。不要把整个响应读进内存。
- 完成后,把最终大小跟 Content-Length 对比。不相等说明连接悄悄断了,需要再来一轮。
可直接运行的代码
import os
import requests
PROXY_URL = os.environ['PROXY_URL'] # строка из личного кабинета Proxeon
PROXIES = {'http': PROXY_URL, 'https': PROXY_URL}
def probe(url):
r = requests.head(url, proxies=PROXIES, allow_redirects=True, timeout=30,
headers={'Accept-Encoding': 'identity'})
r.raise_for_status()
return {
'ranges': r.headers.get('Accept-Ranges', 'none').lower(),
'length': int(r.headers.get('Content-Length', 0) or 0),
'etag': r.headers.get('ETag'),
}
def download_resumable(url, path):
meta = probe(url)
have = os.path.getsize(path) if os.path.exists(path) else 0
if meta['length'] and have >= meta['length']:
print('файл уже полный:', have, 'байт')
return True
headers = {'Accept-Encoding': 'identity'}
if have > 0 and meta['ranges'] == 'bytes':
headers['Range'] = 'bytes=%d-' % have
if meta['etag']:
headers['If-Range'] = meta['etag']
with requests.get(url, headers=headers, proxies=PROXIES, stream=True,
timeout=(30, 120)) as r:
if r.status_code == 206:
expected = 'bytes %d-' % have
if not r.headers.get('Content-Range', '').startswith(expected):
raise IOError('сервер отдал не тот диапазон: ' + r.headers.get('Content-Range', ''))
mode = 'ab'
elif r.status_code == 200:
print('сервер отдаёт файл целиком, начинаем с нуля')
mode = 'wb'
have = 0
elif r.status_code == 416:
raise IOError('запрошенный диапазон вне файла, проверьте локальный размер')
else:
r.raise_for_status()
with open(path, mode) as f:
for chunk in r.iter_content(chunk_size=256 * 1024):
if chunk:
f.write(chunk)
have += len(chunk)
if meta['length'] and have != meta['length']:
print('обрыв: получено %d из %d' % (have, meta['length']))
return False
return True
def download_until_done(url, path, max_rounds=50):
for i in range(max_rounds):
try:
if download_resumable(url, path):
return
except (requests.ConnectionError, requests.Timeout, IOError) as e:
print('заход %d прерван: %s' % (i + 1, type(e).__name__))
# пауза перед следующим заходом: стратегия описана в статье про 429 и ретраи
raise RuntimeError('не удалось докачать файл за %d заходов' % max_rounds)
if __name__ == '__main__':
download_until_done('https://files.example.com/export-2026.csv.gz', 'export-2026.csv.gz')注意 download_until_done 函数:它不包含重试之间的等待逻辑。这是有意的。把你从重试那篇文章里学到的退避策略放进去,这里重要的只是“检查大小、续传、再检查”的循环。
小贴士:如果文件是压缩包,续传过程中不要边下边解压。先拿到完整文件,检查大小,如果服务器提供了校验和就核对一下。然后再解压。下载到一半的 gzip 看起来就像损坏的文件,你会浪费时间去找一个根本不存在的错误。
预期结果
验证:用一个至少 200 MB 的文件运行脚本,十秒后用 Ctrl+C 中断。看一下本地文件的大小,比如 41 943 040 字节。再运行一次脚本。控制台不应出现“从头开始”的字样,文件大小应该继续增长而不是被重置。全部完成后,最终大小必须和 HEAD 请求里的 Content-Length 完全一致。
可能遇到的问题
- 服务器始终返回 200 而不是 206。说明不支持续传。这种数据源只能加大超时、一次性下载完整文件,或者找源站有没有按部分导出的替代方式,比如按日期拆分。
- 没有 Content-Length 头。服务器以 chunked 模式传输,不声明大小。既没法通过大小判断完整性,也没法续传:Range 需要已知偏移量。跟数据源协商,或者用发布的校验和。
- 第一次请求就返回 416。本地文件比服务器上的文件大。服务器上的文件变了、变短了。删掉本地文件重新下载。
- 大小对上了,但文件损坏。多半是中间某次断线后返回 200、以覆盖模式重写、然后又追加写入导致的。重建文件。为了防止再发生,把 ETag 存到单独的文件里,每次请求前比较。
第 2 步:分页导出的检查点
本步目标:保存导出状态,让进程无论在哪崩溃,都能从最后一次成功写入的页继续。
要保存什么
最小检查点取决于数据源的分页类型。分三种情况。
- 游标分页。API 随数据返回类似
next_cursor的字段。就保存它。这是最简单的情况:游标本身已经包含了服务器继续所需的全部信息。 - 按页码或偏移量分页。保存最后一页完整写入的页码和页大小。注意,如果导出期间源数据有新增,偏移量会移动,所以第 4 步的去重是必须的。
- 按键分页。保存最后写入记录的 ID。续传时请求所有大于该 ID 的记录。这种方案不怕插入,也不怕长时间暂停。
不管哪种分页方式,建议在检查点里加几个辅助字段:最后一条记录的 ID(即使是游标方案也留着,作为游标失效时的备用锚点)、页数和行数计数器用于监控进度、导出开始时间和最后更新时间。
状态放哪里
有两种可行方案,都比放在内存变量里好。
方案 A:JSON 文件 + 原子替换
适合结果写在单独文件里而不是数据库的场景。主要陷阱:如果直接往目标文件写状态,进程写到一半崩溃,你会得到一个截断的 JSON,根本读不出来。解决办法是写到旁边的临时文件,再重命名覆盖原文件。同一文件系统内的重命名操作是原子的。
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)方案 B:跟数据放一起的 SQLite 表
如果你把记录存进数据库,这是首选方案。检查点跟页行插入在同一个事务里更新。要么数据和标记都写进去,要么都不写。根本不会出现“数据有了、标记没有”的断层。
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))操作顺序
记住这个规则:先写数据,再更新检查点,最好在同一个事务里。如果事务不可用(比如数据写入文件),顺序依然是这样:写完页文件、同步到磁盘,然后更新状态。这两步之间崩溃只会导致重载一页,而第 3 步保证了这是安全的。反过来先更新检查点则会漏页,那就是真的丢数据了。
小贴士:检查点里存的不要是当前游标,而是服务器返回的下一页游标。这样续传时直接请求你还没有的数据,不用多请求一次已经拿到的页。
预期结果
验证:启动导出,等到十页后强制结束进程。打开数据库执行 sqlite3 export.sqlite,然后执行 SELECT pages, last_id FROM checkpoint;。你应该看到数字 10 和对应的 ID。再执行 SELECT count(*) FROM records;,确认记录数等于 10 倍页大小。重新启动下载器:控制台第一条消息应该是类似“启动:页数 10”。
可能遇到的问题
- 报错“database is locked”。有另一个进程占着连接。关掉所有 sqlite3 窗口和其他打开过该文件的工具。多线程场景请启用 WAL 模式,见第 5 步。
- 游标保存了,数据没保存。你在事务外更新了检查点。回到上面的代码,确认两个操作都在同一个
with con:块里。 - 状态 JSON 空了或损坏了。你直接往文件里写而没有用临时文件加替换。直接用完整的 save_state 函数。
第 3 步:写入结果的幂等性
本步目标:让任何一页重复处理都不会产生重复数据或破坏数据。
为什么重复不可避免
第 2 步你已经见过一种页重复写入的场景:进程在插入数据后、更新检查点前崩溃。除此之外,按偏移量分页在源数据变化时、并发 worker 在重启后拿到同一个任务时、以及你出于“以防万一”手动重跑时,都会产生重复。在请求侧对抗重复是徒劳的。正确做法是让写入本身对重复免疫。
选择去重键
- 源有 ID。直接用。
id、uuid、order_number之类的字段,只要源保证唯一。如果从多个源导出到同一张表,用组合键:源名称加 ID。 - 没有 ID,但有一组字段共同确定一条记录。比如价格表的一行是货号加仓库加日期。用这些字段拼出键,并做规范化:字符串统一大小写、去掉首尾空白、日期统一格式。
- 没有稳定字段。那就用记录规范化后整体的哈希作为键。第 4 步会详细讲这个情况。注意,如果源会更新记录(改价格),哈希会变,你会得到两个版本。有时这正是你想要的,有时不是。
注意:不要用响应里的行号或页码作为键。这些值在源数据任何变化时都会变,去重会变成重复数据生成器。
幂等写入数据库
SQLite 和大多数关系型数据库都有一种构造:要么忽略主键冲突,要么更新已有行。第一种 INSERT OR IGNORE 你在第 2 步见过了。适合记录不变的情况。第二种适合源会更新记录、你需要最新版本的场景:
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])幂等写入文件
如果结果要放文件而不是数据库,用同样的原则:一页对应一个文件名确定的文件。文件名取决于页的参数,而不是时间或计数器。下载前检查最终文件是否存在;存在就跳过这一页。写到临时文件名,完成后重命名,就像 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 True崩溃后留下的 .part 文件启动时可以放心删掉:它们天生就是不完整的。
小贴士:按偏移量分页导出时,别只靠“文件存在就算页完成了”。额外检查一下文件里的行数是否等于页大小(最后一页除外)。存在但为空或过短的文件,最好重新下载。
预期结果
验证:连续三次调用同一页的写入函数。然后执行 SELECT count(*) FROM records;。数量应该等于一页大小,而不是三倍。文件方案下,目录里应该正好一个页文件,没有 .part 文件。
第 4 步:不撑爆内存的结果去重
本步目标:在数百万行的流式处理中剔除重复记录,不让所有键都待在内存里。
记录哈希
当记录没有 ID 时,就用它内容的哈希作为键。要让相同记录产生相同哈希,内容必须规范化:字典键排序、去掉多余空白、固定分隔符。否则同一条记录换个字段顺序就会得到不同哈希。
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()函数返回 16 字节。这足够了:对于几亿条记录,随机碰撞的概率可以忽略不计。参数 fields 让你只对稳定字段算哈希,排除比如每次请求都变的时间字段。
为什么内存集合处理不了几百万条
第一个念头通常是:弄个 seen = set() 把键丢进去。算一下。一个 16 字节的 bytes 对象在 Python 里占约 49 字节加上数据本身,总共约 65 字节。集合槽位按填充系数摊上,再加约 30 字节。大概每个键 95 字节。1000 万条记录约 950 MB,5000 万条接近 5 GB。而且最关键的是:进程重启后集合是空的,去重又从零开始。
三种不撑爆内存的办法
- 把键存进数据库本身。最简单最可靠。如果键是 records 表的主键,第 3 步的 INSERT OR IGNORE 已经完成了去重。索引在磁盘上,重启后还在,SQLite 自己会缓存热索引页。1000 万个 16 字节键的索引大约占磁盘 400-500 MB,但不占内存。
- 独立的已见键表、无 rowid。适合数据本身不写 SQLite、而是写文件的场景。那时 SQLite 只作为磁盘上的紧凑集合使用。
- 内存里的压缩集合作为预过滤。进阶方案:把哈希截断到 8 字节,以整数形式存在有序数组里,或用布隆过滤器。内存能省好几倍,但会有误判概率。所以这种预过滤只用来快速排除肯定新的记录,最终判断还是查数据库。
磁盘集合的实现
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])写页之前调用 filter_new,remember 放在写数据和检查点的同一个事务里。这样崩溃后,已见集合、数据和进度标记始终一致。
小贴士:带 500 个值的 IN 查询走索引只要几毫秒。不要循环逐条查键:每次查询的开销会让它慢几十倍。
预期结果
验证:构造一页 500 条测试记录,其中 100 条在页内重复两次,另外 100 条已经在 seen 表里。filter_new 应该正好返回 300 条。进程重启后,同样的 500 条应该返回零条新记录。
可能遇到的问题
- 重复数据还是混进来了。检查规范化:多半是记录里有个请求时间字段,或者列表元素顺序随机。用 fields 参数排除它,或者在哈希前对嵌套列表排序。
- 几百万行后插入变慢。索引不再装得进缓存。用
PRAGMA cache_size=-200000加大 SQLite 缓存(即 200 MB),并确认插入是按页成批在同一个事务里,而不是一行一行。
第 5 步:不丢数据的并发
本步目标:通过代理用多个并发 worker 加速导出,同时任何一个 worker 崩溃都不会丢任务或搞坏数据库。
什么时候能并发,什么时候不能
游标分页天生是顺序的:下一个游标只有拿到上一页才知道。没法直接并发。但几乎总是可以把导出拆成独立的分片:按天、按类别、按区域、按 ID 首字母。每个分片顺序导出、有自己的检查点,分片之间并行。按偏移量和按键分页可以直接并发:任务形如“第 1 到 100 页”或“ID 从 0 到 100000”。
磁盘上的任务队列
内存队列随进程一起死。所以任务存在带状态的表里:
pending—— 等待执行;running—— 已被 worker 取走;done—— 已完成并写入;failed—— 尝试次数用尽,需要人工处理。
启动时下载器先把所有 running 的任务改回 pending:如果它们卡在这个状态,说明上个进程半路死了。然后 worker 开始消费 pending。
单一写入点
SQLite 允许多个并发读者,但只有一个写者。最简单最安全的模式是:worker 只负责下载并返回数据,所有写库操作由主线程完成。代码里没有锁,也不会出现“database is locked”。另外启用 WAL 模式,这样另一进程读状态不会阻塞写入。
并发上限
worker 数量受两方面限制。第一是代理能力:如果你在 Proxeon 个人中心有多个通道,每个通道放一两个 worker 比较合理,这样某个通道 IP 轮换时不会同时断掉所有线程的连接。第二是对数据源的礼貌:即使没有正式限流,十个并发线程打在一个小 API 上也会造成压力,招来拒绝。先从三四个 worker 开始,观察错误率再往上加。
并发下载器代码
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)函数 fetch_fn 通过代理发请求并返回记录列表。它在子线程里运行,不碰数据库。函数 write_fn 在主线程的事务里调用,执行第 3 步的幂等插入。失败的任务会自动回到 pending,下一轮 claim 会被再次取走;尝试次数用尽后状态变为 failed,你手动处理。
注意:不要把 sqlite3 连接对象传给 worker。连接绑定在创建它的线程上,在别的线程使用会报错,更糟的是可能悄悄损坏数据。每个线程要么有自己的连接,要么像上面例子一样根本不碰连接。
小贴士:给每个 worker 单独建一个 requests.Session,配置各自的代理地址。如果你有多个 Proxeon 通道,就轮着分配给 worker:worker 0 用通道 0,worker 1 用通道 1,以此类推。这样某个通道连接挂掉只影响一个线程。
预期结果
验证:往队列里放 100 个任务,启动四个 worker,半分钟后杀掉进程。执行 SELECT status, count(*) FROM tasks GROUP BY status;。你会看到几个 done、几个 running,其余 pending。再次启动:running 应该在启动时消失,全部完成后所有任务都在 done,除了真的失败并在 last_error 里留下错误信息的那些在 failed。
第 6 步:长时间中断后如何续传
本步目标:正确继续一个停了几小时甚至几天的导出,不踩到失效状态的坑。
什么会过期
隔十秒续传和隔一周续传是两回事。长时间中断后,部分保存的状态不再有效。
- 会话和 cookie。服务端会话通常活几小时到一天。过期后保存的 cookie 会导致 401 或跳转到登录页。解决办法:启动时完整重新登录,而不是从文件恢复 cookie。
- 访问令牌。OAuth 令牌寿命一小时,有时更短。如果有 refresh token,启动前和运行中按计划刷新 access token,别等到被拒。
- 分页游标。很多 API 限制游标寿命只有几分钟到几小时。过期游标会返回 400 和“游标无效”的报错。所以第 2 步我们保存了备用锚点:最后一条记录的 ID。如果数据源支持按 ID 或修改时间过滤,就从锚点构造新请求。如果不支持,分片只能从头开始,第 4 步的去重会剔除已有的数据。
- 服务器上的文件内容。第 1 步的续传关键在文件没变。每次请求前把当前 ETag 跟保存的比较;不一致就重新下载整个文件。
- 代理配置。一周内 Proxeon 个人中心可能改了端口、密码,或者通道到期了。在处理队列之前先用测试请求验证代理。
- 数据本身。如果导出跑了一周,而源站这期间增删过记录,你的结果是不同时间点的混合。很多任务可以接受。如果不能,就把开始时间存进检查点,完成后单独跑一次增量遍历,拉取该时间之后修改过的记录。
飞行前检查
把所有检查集中到一个函数里,启动时先执行,再开始任何实际工作。它要么把状态收拾好,要么以清晰的提示停下下载器。
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 state函数 refresh_access_token 取决于你的数据源:通常是带 refresh token 的 POST 请求,之后更新会话的 Authorization 头。字段 resume_after_id 随后在获取页函数里作为“ID 大于该值”的过滤条件。
小贴士:refresh token 和代理密码不要放在检查点里,放在环境变量或权限受限的独立密钥文件里。检查点你会复制、发给同事、附在错误报告里;密钥放在那里是多余的。
预期结果
验证:手动破坏 checkpoint 表里的游标,执行 UPDATE checkpoint SET cursor='broken';,然后启动下载器。控制台应该出现切换到 last_id 锚点的提示,导出继续而不崩溃。完成后的记录数应该跟不破坏游标的对照运行一致。
第 7 步:完整的稳健 Python 下载器骨架
本步目标:把前面所有内容整合到一个文件里,可以启动、中断、再启动,并得到没有重复的完整结果。
骨架结构
- 从环境变量读取配置:Proxeon 代理地址、API 地址、token、任务名。
- Store 类:SQLite,包含记录表和检查点表,每页一个事务。
- 用于幂等的记录键函数。
- 通过代理获取页的函数。
- 主循环:按检查点续传,断线后重建会话。
完整代码
# resumable_loader.py
import hashlib
import json
import os
import sqlite3
import time
from datetime import datetime, timezone
import requests
PROXY_URL = os.environ['PROXY_URL']# http://USER:PASS@HOST:PORT из кабинета Proxeon
API_BASE = os.environ.get('API_BASE', 'https://api.example.com')
API_TOKEN = os.environ.get('API_TOKEN', '')
DB_PATH = os.environ.get('DB_PATH', 'export.sqlite')
JOB = os.environ.get('JOB', 'orders-2026')
PAGE_SIZE = 500
MAX_ATTEMPTS = 8
FIELDS = ('cursor', 'last_id', 'pages', 'rows', 'started_at')
def now():
return datetime.now(timezone.utc).isoformat()
def record_key(rec):
canon = json.dumps({'id': rec['id']}, sort_keys=True, separators=(',', ':'))
return hashlib.blake2b(canon.encode('utf-8'), digest_size=16).digest()
class Store:
def __init__(self, path):
self.con = sqlite3.connect(path)
self.con.execute('PRAGMA journal_mode=WAL')
self.con.executescript('''
CREATE TABLE IF NOT EXISTS records(
key BLOB PRIMARY KEY,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
) WITHOUT ROWID;
CREATE TABLE IF NOT EXISTS checkpoint(
job TEXT PRIMARY KEY,
cursor TEXT,
last_id TEXT,
pages INTEGER NOT NULL,
rows INTEGER NOT NULL,
started_at TEXT,
updated_at TEXT
);
''')
def load(self, job):
row = self.con.execute(
'SELECT cursor, last_id, pages, rows, started_at FROM checkpoint WHERE job=?',
(job,)).fetchone()
if row is None:
return {'cursor': None, 'last_id': None, 'pages': 0, 'rows': 0, 'started_at': now()}
return dict(zip(FIELDS, row))
def commit_page(self, job, rows, state):
ts = now()
with self.con:
self.con.executemany(
'INSERT OR IGNORE INTO records(key, payload, fetched_at) VALUES (?, ?, ?)',
[(record_key(r), json.dumps(r, ensure_ascii=False), ts) for r in rows])
self.con.execute(
'INSERT INTO checkpoint(job, cursor, last_id, pages, rows, started_at, updated_at) '
'VALUES (?, ?, ?, ?, ?, ?, ?) '
'ON CONFLICT(job) DO UPDATE SET cursor=excluded.cursor, last_id=excluded.last_id, '
'pages=excluded.pages, rows=excluded.rows, updated_at=excluded.updated_at',
(job, state['cursor'], state['last_id'], state['pages'], state['rows'],
state['started_at'], ts))
def unique_count(self):
return self.con.execute('SELECT count(*) FROM records').fetchone()[0]
def make_session():
s = requests.Session()
s.proxies = {'http': PROXY_URL, 'https': PROXY_URL}
s.headers['User-Agent'] = 'resumable-loader/1.0'
if API_TOKEN:
s.headers['Authorization'] = 'Bearer ' + API_TOKEN
return s
def fetch_page(session, cursor):
params = {'limit': PAGE_SIZE}
if cursor:
params['cursor'] = cursor
r = session.get(API_BASE + '/orders', params=params, timeout=(15, 90))
r.raise_for_status()
body = r.json()
return body['items'], body.get('next_cursor')
def run():
store = Store(DB_PATH)
state = store.load(JOB)
print('старт: страниц %d, строк %d, уникальных в базе %d'
% (state['pages'], state['rows'], store.unique_count()))
session = make_session()
attempts = 0
while True:
try:
items, next_cursor = fetch_page(session, state['cursor'])
attempts = 0
except (requests.ConnectionError, requests.Timeout, requests.HTTPError) as e:
attempts += 1
if attempts > MAX_ATTEMPTS:
print('попытки исчерпаны, состояние сохранено, запустите снова позже')
raise
print('обрыв (%s), попытка %d из %d' % (type(e).__name__, attempts, MAX_ATTEMPTS))
time.sleep(min(60, 2 ** attempts)) # выбор пауз описан в статье про 429 и ретраи
session = make_session() # новая сессия: соединение через прокси пересоздаётся
continue
if not items:
break
state['pages'] += 1
state['rows'] += len(items)
state['last_id'] = str(items[-1]['id'])
state['cursor'] = next_cursor
store.commit_page(JOB, items, state)
if state['pages'] % 20 == 0:
print('страниц %d, строк %d' % (state['pages'], state['rows']))
if next_cursor is None:
break
print('готово: страниц %d, строк получено %d, уникальных в базе %d'
% (state['pages'], state['rows'], store.unique_count()))
if __name__ == '__main__':
run()如何适配自己的数据源
- 把路径
/orders以及字段名items、next_cursor、id换成你的 API 实际返回的。就在 fetch_page 和 record_key 里三个地方。 - 如果数据源按偏移量分页,把 cursor 参数换成 page,下一个值算成 state['pages'] + 1。检查点里保存页码而不是游标。
- 如果按键分页,传一个类似
after_id的参数,值取自 state['last_id'],去掉游标相关逻辑。 - 如果需要并发,把 fetch_page 作为第 5 步的 fetch_fn,commit_page 作为 write_fn。把导出拆成分片并填好任务队列。
- 在主循环前加上第 6 步的 preflight 函数。
结果验证:检查清单
在对实际跑几小时的数据量启动之前,先跑一遍这个清单。每项只花几分钟,加在一起能保证夜间导出不丢数据。
- 中断测试。启动下载器,30 秒后按 Ctrl+C。再启动。输出的第一行应该显示非零的页数,而不是“页数 0”。
- 重复测试。把 PAGE_SIZE 改成 10,连续五次在随机时刻中断。完成后比较唯一记录数和获取的总行数:唯一记录应该小于等于总行数,在干净的游标分页且源无插入时几乎相等。
- 一致性测试。任何一次中断后执行两个查询:
SELECT rows FROM checkpoint;和SELECT count(*) FROM records;。两者差值不应超过一页大小。超过说明检查点和数据不是在一个事务里写的。 - 代理测试。暂时停掉 PROXY_URL 环境变量或填错密码。下载器应该在第一个请求就以清晰错误退出,而不是卡住或直接绕过代理去访问源站。
- 游标过期测试。按第 6 步破坏数据库里的游标,确认切换到锚点的逻辑生效。
- 磁盘测试。一千页后看看 export.sqlite 的大小,乘以预计总页数。确认磁盘有 20% 余量的空间。
验证:成功标准是:经过三次故意中断和三次重启后,最终唯一记录数与同一数据源上不间断跑一次得到的结果一致,而且控制台里一次都没有出现从头开始的提示。
进阶能力与优化
- 进度与时间估计。如果知道总记录数,每二十页输出一次百分比和剩余时间估计。这既方便你自己看,也能区分卡住和跑得慢。
- payload 压缩。几千万条记录时,JSON 文本占空间很多。写入前用 zlib.compress 压缩 payload 字段并存为 BLOB。通常能省三到八倍。
- 迁移到服务端数据库。检查点跟数据同事务的方案几乎可以原样搬到 PostgreSQL。ON CONFLICT 构造支持,而“单写者”限制消失。
- 增量导出。保存每个任务的开始时间,全量完成后单独跑一个带“修改时间晚于”过滤的任务。这样不用全量重拉就能保持数据最新。
- 每个分片单独进程。不用线程,可以启动多个脚本实例,每个用不同的 JOB 和不同的 DB_PATH,最后合并结果。调试更简单,也完全避免了并发写入的问题。
- 断线指标。记录每次断线,包括异常类型和时间。过一天你就会发现断线集中在 Proxeon 通道 IP 轮换的间隔附近,可以据此调整轮换周期来匹配请求时长。
常见错误与解决方案
下面是几乎每个人头几次运行时都会踩的坑。格式:问题、原因、解决方案。
- 问题:重启后导出每次都从零开始。原因:检查点写在内存或无法在崩溃后保住的文件里,或者下载器启动时没读它。解决:确认 run 函数第一个动作是 store.load,而且每页之后在事务里更新状态。
- 问题:数据库里的记录比源站多一半。原因:去重键不稳定:里面包含了请求时间、页码,或者顺序随机的字段。解决:只用源 ID 或通过 fields 参数明确指定的稳定字段算键。
- 问题:数据库里的记录比源站少,但导出正常完成了。原因:检查点在写数据之前更新,崩溃后漏了一页。或者源站删除记录后,按偏移量分页把页码往回移了。解决:把顺序改成“先数据、后检查点”并在同一事务里;对有删除的源站改用按键分页。
- 问题:文件续传后大小对得上但压缩包损坏。原因:服务器某次返回了 200 而不是 206,文件被部分覆盖后再追加。解决:把 ETag 存在文件旁边,不一致就删文件;检查 Content-Range 是否与请求的偏移量匹配。
- 问题:并发时“database is locked”。原因:多个线程同时写 SQLite,或者连接被跨线程传递。解决:主线程单一写入点,worker 只下载;启用 WAL;连接在哪个线程用就在哪个线程建。
- 问题:跑了一小时后所有请求开始返回 401。原因:访问令牌过期。解决:到期前按计划刷新令牌;收到 401 时触发刷新并重发一次请求,不要算作断线。
- 问题:内存涨到几个 GB。原因:已见键集合或全部记录列表留在进程内存里。解决:用数据库主键或 seen 表去重;按页写数据,不累积。
- 问题:断线严格每隔几分钟发生一次。原因:跟代理通道 IP 轮换的间隔吻合。解决:这是正常情况,下载器必须能扛过去。如果请求较长,在 Proxeon 个人中心把轮换间隔调成明显大于单次请求的典型时长,或者用按请求轮换、在页与页之间换。
FAQ:稳健导出的常见问题
结果要 CSV,也必须用 SQLite 吗?
不必,但方便。SQLite 在这里充当可靠的状态存储和已见键集合。完成后一条命令就能从 records 表导出 CSV。如果想完全不用数据库,用第 3 步的文件方案,一页一个文件,配 JSON 检查点加原子替换。
检查点多频繁保存一次:每页还是更少?
每页。一次含几百行的 SQLite 事务只要几毫秒,跟通过代理的网络请求相比微不足道。省着不写检查点的收益抵不上丢几十页的风险。
如果 API 既不给游标也不给 ID,只有页码怎么办?
按页码做,把它存进检查点,并务必启用基于记录内容哈希的去重。接受一个事实:源站活跃变化时,页码偏移可能导致部分记录被跳过。对关键数据可以反向再跑一遍页:第一遍漏掉的,第二遍大概率能补上。
能用多线程按不同范围并发下载一个文件吗?
如果服务器支持 Range,可以。把文件切成 50-100 MB 的块,每块作为第 5 步队列里的一个任务,各自有临时文件,全部完成后按正确顺序拼接。每块检查大小,如果提供校验和就校验整个文件。
通过代理工作时设多少个 worker?
先在一个 Proxeon 通道上三到四个,观察 tasks 表里的错误率。错误率低于 1% 可以再加两个。错误率上升就减。单个通道超过十个线程很少带来收益:要么撞上通道带宽,要么惹毛数据源。
需要在两次运行之间保存会话 cookie 吗?
通常不用。启动时重新登录只要几秒,比恢复寿命未知的 cookie 更可靠。例外:数据源限制每天登录次数。那就保存 cookie,但一遇到 401 或跳转登录页就丢弃并重新登录。
怎么判断导出是完整结束而不是悄悄断了?
文件:大小等于 Content-Length 且校验和匹配。API:拿到没有 next_cursor 的页或空页,同时行数与源站声明的总数一致。在检查点里写上明确的完成标记,避免再次启动时又开一轮新的遍历。
怎么处理状态为 failed 的任务?
看 last_error 字段。网络错误的话,用 UPDATE 把任务改回 pending 再跑一遍。数据解析错误说明源里有格式不规范的记录:改代码后重跑。绝不要悄悄删掉 failed,它是导出缺什么的唯一证据。
不用 requests 而是用异步库可以吗?
可以,原则一样:原子工作单元、检查点与数据同写、幂等写入、磁盘队列。变的只是传输层。唯一注意:SQLite 写入保持同步和顺序,把并发放在网络请求层面。
结语
你从“文件不完整、心惊胆战地重启”一路走到了“断线无所谓”的下载器。总结一下做了什么。
- 搞定了 HTTP 断点续传:用 HEAD 检查 Accept-Ranges、用 Range 头带当前文件大小、区分 206 和 200、用 ETag 和 If-Range 防止文件被替换。
- 建立了分页导出的检查点:游标、页码或最后一条记录 ID 加辅助计数器,全部与数据在同一事务里。
- 通过主键和 INSERT OR IGNORE 或 ON CONFLICT DO UPDATE 实现幂等写入,文件方案用确定性文件名和原子替换。
- 把去重放在磁盘上,几百万个键不占内存,还能跨重启存活。
- 加了并发:SQLite 任务队列、失败任务自动回队、单一写入点。
- 考虑了长期中断后续传:刷新令牌、重新登录、游标失效后切换到 ID 锚点、启动前检查 Proxeon 代理。
- 把所有东西整合成一个可直接使用的骨架,改三四个字符串就能适配具体数据源。
下一步做什么
拿第 7 步的骨架,在一个小的真实数据量上跑,比如一万条记录。走一遍验证部分的清单。确认无误后,再启动整夜的完整导出。第二天早上你要么看到计数一致的“完成”,要么看到尝试用尽、状态已保存的提示,那时再跑一次脚本就行。
未来可以往哪发展
下一个层次是基于修改时间的增量导出而不是全量遍历,把状态搬到服务端数据库支持多机运行,以及根据响应状态码制定合理的重试策略——那部分有专门讲 429 和重试的文章。把重试那篇的策略和本文的稳健状态结合起来,你就得到了一个可以连续跑一周、不用打开终端的下载器。
最后一句。长时间导出断线不是事故,而是一种你如今已经能从容处理的正常情况。祝你导出顺利。