前言
2025年底公司有个老系统要下线,里面攒了几年的培训素材(视频、音频、文档、图片,几万条记录)都躺在 OSS 上,需要全部爬下来再转存到新的 OSS。量不小,单个视频还有好几个G的,直接写个 for 循环 requests 下载肯定歇菜,于是写了一个比较健壮的爬取脚本。完整代码在:
https://github.com/anTtutu/anttu.code.learn.python/tree/master/oss_download
这篇把思路和关键代码整理下。
1、难点在哪
先说清楚这个活的难点,也是写脚本时所有设计的出发点:
数据量大:
- 几万条素材记录,存在 MySQL 的
dim_resource_info 表里,每条记录带一个 resource_path(原 OSS 的 URL)
- 一条条串行下载根本跑不完,必须并发,但并发又不能把带宽和 OSS 打爆
- 跑到一半挂了不能从头来,几万条重跑一遍谁受得了,必须断点续传
- 哪些成功哪些失败得有账可查,不然没法向上交代
大文件难搞:
- 单个视频文件好几个G,默认超时时间(比如 aiohttp 总超时 5 分钟)跑到一半必断
- 大文件不能
response.read() 一把梭进内存,得流式分块写盘
- 还有 m3u8 格式的视频流,不是普通文件下载,要 ffmpeg 转封装
- 下载完还要转存到新 OSS,大文件上传同样是问题
另外还有些杂活:文件名没有后缀的要补、m3u8 要转成 mp4、文档要提取页数、图片要拿宽高,全部回写数据库。
2、整体思路
MySQL(素材清单) → 并发下载到本地 → 上传到新 OSS → 回写状态(DB + Excel)
按文件类型拆了多个脚本(视频、音频、pdf、ppt、xls、图片、zip 各一个),结构一样只差 SQL 条件和后缀处理。为什么拆?因为视频是大文件大户,单独跑方便控制节奏、单独重试,不用被小文件拖着。
核心依赖:
1
2
3
4
5
6
7
8
|
aiohttp + asyncio # 异步并发下载
tenacity # 失败重试
pymysql # 读清单、回写状态
alibabacloud_oss_v2 # 阿里云 OSS V2 SDK(上传新 OSS)
m3u8 + ffmpeg # m3u8 流媒体转 mp4
PyPDF2 / python-docx / python-pptx / PIL # 文档页数、图片宽高
openpyxl # Excel 进度账本
tqdm # 进度条
|
3、配置设计
所有参数收在一个 config 字典里,改脚本只动这一块:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
|
config = {
"db": {
"host": "your_db_host",
"user": "your_user",
"password": "your_password",
"database": "enterprise_lesson",
"port": 3306
},
"oss": {
"enable_upload": True, # 是否启用上传到 OSS
"access_key_id": "your_ak", # 建议走环境变量,别硬编码
"access_key_secret": "your_sk",
"bucket_name": "your-bucket",
"region": "cn-beijing",
"endpoint": "oss-cn-beijing.aliyuncs.com",
"base_url": "https://your-bucket.oss-cn-beijing.aliyuncs.com"
},
"download": {
"save_path": "./downloads",
"concurrent_limit": 100, # 最大并发数
"retry_attempts": 3, # 失败重试次数
"retry_interval": 5, # 重试间隔(秒)
"chunk_size": 1024 * 1024, # 分块 1MB,流式写盘
"time_out": 3600 * 8 # 总超时 8 小时!因为有几个G的文件
},
"m3u8": {
"ffmpeg_path": "ffmpeg",
"temp_path": "./temp_m3u8"
}
}
|
重点说 time_out: 3600 * 8:一开始用默认超时,大视频下到一半就断,排半天才发现是超时的事。几个G的文件按家宽上传带宽算,8 小时是留了充足余量的。
OSS 凭证走环境变量(SDK 的 EnvironmentVariableCredentialsProvider 自动读取),比硬编码在代码里安全,代码还要开源,AK 写进去就是事故:
1
2
3
4
5
6
7
|
os.environ["OSS_ACCESS_KEY_ID"] = config['oss']['access_key_id']
os.environ["OSS_ACCESS_KEY_SECRET"] = config['oss']['access_key_secret']
credentials_provider = oss.credentials.EnvironmentVariableCredentialsProvider()
cfg = oss.config.load_default()
cfg.credentials_provider = credentials_provider
cfg.region = config['oss']['region']
return oss.Client(cfg)
|
4、大文件的处理
4.1 流式分块下载,不进内存
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
|
@retry(stop=stop_after_attempt(config["download"]["retry_attempts"]),
wait=wait_fixed(config["download"]["retry_interval"]))
async def download_file(session, url, save_path):
os.makedirs(os.path.dirname(save_path), exist_ok=True)
# 断点续传第一层:本地已有就跳过
if os.path.exists(save_path):
logger.info(f"文件已存在,跳过下载:{save_path}")
return save_path
async with session.get(url, timeout=aiohttp.ClientTimeout(
total=config['download']['time_out'])) as response:
if response.status != 200:
raise Exception(f"下载失败,状态码:{response.status}, URL: {url}")
with open(save_path, "wb") as f:
while True:
chunk = await response.content.read(config["download"]["chunk_size"])
if not chunk:
break
f.write(chunk)
return save_path
|
三个要点:
- 1MB 一块循环写盘,内存占用恒定,几个G的文件也吃得消
- tenacity 装饰器:网络抖动失败自动重试 3 次、间隔 5 秒,重试还带着断点判断(文件存在就跳过),不会重复下载
- 超时给到 8 小时,这最坑,默认超时大文件必挂
4.2 断点续传:两层保险
- 第一层:SQL 里带
where is_downloaded = 0,下载成功回写 is_downloaded = 1,重跑脚本自动跳过已完成的
- 第二层:本地文件存在直接跳过下载(防止 DB 状态没写成功但文件其实下完了的情况)
失败的任务会把错误信息写进 error_message 字段:
1
2
|
sql = "UPDATE dim_resource_info SET is_downloaded=%s, download_time=%s, \
is_upload_oss=%s, upload_time=%s, oss_path=%s, error_message=%s WHERE resource_id=%s"
|
重跑的时候把报错的 resource_id 挑出来单独处理,或者按 error_message 分类批量重试。
4.3 m3u8 视频流转 mp4
m3u8 不是完整文件,是一堆 ts 分片的索引,普通下载下来没法播,用 ffmpeg 转封装:
1
2
3
4
5
6
7
8
9
10
11
|
@retry(stop=stop_after_attempt(config["download"]["retry_attempts"]),
wait=wait_fixed(config["download"]["retry_interval"]))
async def download_m3u8(session, url, save_path):
temp_file = os.path.join(config["m3u8"]["temp_path"],
f"{Path(save_path).stem}.m3u8")
async with session.get(url, ...) as response:
...
f.write(await response.read())
# ffmpeg 拉流合并,-c copy 不转码,速度飞快
command = f"{config['m3u8']['ffmpeg_path']} -i {url} -c copy {save_path}"
subprocess.run(command, shell=True, check=True)
|
-c copy 是关键:只做转封装不转码,几个G的视频几分钟搞定;如果转码,那就是小时级了。
4.4 大文件上传
下载完了要转存到新 OSS,用 V2 SDK 的文件上传接口:
1
2
3
4
|
with open(local_file_path, 'rb') as file:
result = client.put_object_from_file(
PutObjectRequest(bucket=bucket, key=object_name, body=file),
local_file_path)
|
代码里预留了 upload_big_file_to_oss 的分片上传接口——这次的单文件几个G简单上传扛住了,如果以后有几十G的,就要换分片上传(resumable upload),支持并发分片和断点续传。
5、数据量大的处理
5.1 asyncio 并发 + Semaphore 限流
1
2
3
4
5
6
7
8
9
10
11
12
|
async def main():
materials = fetch_material_data() # 一次性拉出几万条清单
semaphore = asyncio.Semaphore(config["download"]["concurrent_limit"])
with tqdm(total=len(materials), desc="下载进度", unit="文件") as progress_bar:
async with aiohttp.ClientSession() as session:
tasks = [process_material_with_semaphore(m, session, semaphore, progress_bar)
for m in materials]
await asyncio.gather(*tasks) # 几万个任务一次挂上去
async def process_material_with_semaphore(material, session, semaphore, progress_bar):
async with semaphore: # 同时最多 100 个在跑
await process_material(material, session, progress_bar)
|
任务全部挂上去,但信号量保证同时只有 100 个在实际跑。为什么是 100?小文件可以更高,但这批里混着几个G的视频,并发太高带宽打满、每个都慢,还容易触发 OSS 限流,100 是试出来的平衡点。
5.2 状态账本:DB + Excel 双保险
- DB 回写:成功回写
is_downloaded=1、oss_path;失败写 error_message,重试有依据
- Excel 进度文件:每条任务(成功或失败)都往
progress.xlsx 追加一行,给不懂数据库的同事看的账本,出问题了打开表格就能对
1
2
3
4
5
6
7
8
9
10
11
12
13
|
def update_excel_progress(data):
progress_file = config["excel"]["progress_file"]
if not os.path.exists(progress_file):
wb = Workbook()
ws = wb.active
ws.append(["ID", "资源ID", "资源名称", "文件名", "资源路径",
"资源类型", "创建时间", "状态", "OSS完整路径", "备注"])
wb.save(progress_file)
wb = load_workbook(progress_file)
ws = wb.active
for row in data:
ws.append(row)
wb.save(progress_file)
|
5.3 分批分脚本跑
按类型拆脚本还有个好处:先跑文档图片这种小文件(几小时清完一大批,快速验证脚本没问题),最后跑视频大文件(挂机跑几天)。重试也是同理,SQL 条件改一改就能只跑失败的那批:
1
2
3
4
|
-- 正常批次
select ... from dim_resource_info where resource_type_tag = '视频' and is_downloaded = 0;
-- 失败重试批次:按 resource_id 挑着跑
select ... from dim_resource_info where resource_id in (..., ...) and is_downloaded = 0;
|
6、细节处理
这些不起眼但少一个都会出问题:
URL 合法性校验:表里混着第三方外链(不在我们 OSS 上),不校验直接下会报一堆没用的错:
1
2
3
|
def is_valid_resourcepath(resourcepath):
parsed = urlparse(resourcepath)
return parsed.scheme in ('http', 'https') and parsed.netloc != ''
|
后缀补全规则:resource_name 有后缀而 resource_path 没有的,用 resource_name 的补;m3u8 的转成 .mp4 或按 resource_name 的音视频后缀;下载成功后把文档页数(PyPDF2/docx/pptx 按类型解析)、图片宽高(PIL)也回写到库里,一次搬家顺便把数据治理也做了。
日志双写:文件 + 控制台,挂机跑几天回来只看 download.log 就能复盘全程。
7、踩坑记录
- 默认超时下大文件必挂:aiohttp 默认总超时 5 分钟,几个G的视频下到 4 分多钟准时断,一开始还以为是网络问题,白查半天
- 全量 read 进内存:一开始图省事
response.read(),大文件直接把进程内存撑爆,改成 1MB 分块循环写
- 几万个任务直接 gather:一次性创建几万个协程任务,内存和句柄都有压力,用 Semaphore 限流后平稳
- DB 状态和实际文件不一致:下载成功但回写 DB 失败的情况存在,所以加了本地文件存在即跳过的第二层判断
- AK 千万别硬编码进开源仓库:凭证走环境变量,代码里只留 your_ak 占位
总结
整个活下来最大的体会:爬虫脚本的健壮性不在于代码多花哨,而在于把「必然会失败的」都提前想好——超时按最坏情况给、失败自动重试、状态随时可查、断了能接着跑。几万条记录、TB 级的数据就是靠这些朴素手段平稳搬完的。脚本已开源,遇到类似的 OSS 搬家、素材库迁移的活可以直接参考改改。
相关阅读