前言

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

三个要点:

  1. 1MB 一块循环写盘,内存占用恒定,几个G的文件也吃得消
  2. tenacity 装饰器:网络抖动失败自动重试 3 次、间隔 5 秒,重试还带着断点判断(文件存在就跳过),不会重复下载
  3. 超时给到 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、踩坑记录

  1. 默认超时下大文件必挂:aiohttp 默认总超时 5 分钟,几个G的视频下到 4 分多钟准时断,一开始还以为是网络问题,白查半天
  2. 全量 read 进内存:一开始图省事 response.read(),大文件直接把进程内存撑爆,改成 1MB 分块循环写
  3. 几万个任务直接 gather:一次性创建几万个协程任务,内存和句柄都有压力,用 Semaphore 限流后平稳
  4. DB 状态和实际文件不一致:下载成功但回写 DB 失败的情况存在,所以加了本地文件存在即跳过的第二层判断
  5. AK 千万别硬编码进开源仓库:凭证走环境变量,代码里只留 your_ak 占位

总结

整个活下来最大的体会:爬虫脚本的健壮性不在于代码多花哨,而在于把「必然会失败的」都提前想好——超时按最坏情况给、失败自动重试、状态随时可查、断了能接着跑。几万条记录、TB 级的数据就是靠这些朴素手段平稳搬完的。脚本已开源,遇到类似的 OSS 搬家、素材库迁移的活可以直接参考改改。

相关阅读