1. 内容获取工作流的核心痛点与重构思路
做过内容批量采集的人都有一个共识:真正让人头疼的从来不是“能不能下载”,而是“下载过程稳不稳”。douyin-downloader 这类工具在圈子里流传已久,早期版本大多走的是单链路请求——解析一个视频 ID,发一次请求,拿到地址就落盘。这套逻辑在测试阶段跑十个八个链接看着挺美,一旦上量到几百上千条,失败率就会肉眼可见地往上飙。我自己最早搭的那套脚本,跑 500 条链接能成功 380 条就算烧高香了,剩下的要么超时、要么返回空数据、要么拿到的地址已经失效。
问题的根子在于:内容平台的接口行为不是静态的。同一个请求参数,在不同时间、不同网络出口、不同频次下,返回结果可能完全不一样。单链路架构把“一次请求成功”当成了默认前提,这本身就是个危险的假设。所以这次重构的核心目标很明确——把“一次成功”的赌注,改成“多层兜底”的确定性。
所谓 3 层容错架构,说白了就是给每一次内容获取准备三条命:第一层走主接口直连,追求速度和吞吐;第一层挂了自动降级到第二层,走备用解析通道;第二层也拿不到,第三层用浏览器环境兜底,模拟真实用户行为把内容捞回来。三层之间不是简单的重试关系,而是不同技术路径的接力,每一层解决的是前一层解决不了的那类失败。
这套思路的价值在于,它把“失败”从一个终点变成了一个中间状态。任何一层失败都不代表这条任务失败,只是意味着要换一条路走。配合 SQLite 做任务状态持久化,整个工作流就具备了断点续跑的能力——进程崩了、机器重启了,重新拉起来接着跑,不会从头再来。
适合谁来参考这套方案?如果你只是偶尔下几个视频,单链路脚本够用了,没必要上这套。但如果你在做内容归档、素材库建设、批量分析这类需要稳定吞吐的场景,或者你正在用 dify 工作流、扣子工作流、comfyui 工作流这类自动化编排工具做内容处理,那这套容错思路可以直接迁移过去。它本质上是一套任务可靠性工程的实践,跟具体下什么内容无关。
2. 三层容错架构的设计逻辑与选型考量
2.1 为什么是三层而不是两层或四层
层数不是拍脑袋定的。两层(主接口+浏览器兜底)的问题是中间缺少一个“轻量级补救”环节。主接口失败的原因有很多种:有的是参数问题,换个解析方式就能过;有的是频控问题,等一会儿换个通道就行;有的是内容本身受限,只能靠浏览器环境。如果只有两层,所有非主接口的失败都压到浏览器层,会导致浏览器资源被大量占用,吞吐直接崩掉。
四层呢?加一层“代理池轮换”听起来很美,但实际维护成本极高,而且代理质量参差不齐,反而引入新的不确定性。三层是一个平衡点:第一层保吞吐,第二层保成功率,第三层保底线。每层职责清晰,不会互相干扰。
2.2 各层的技术选型与职责边界
第一层我选的是直连接口解析。这层的目标是快,单条任务控制在 2 秒内完成。实现上用轻量 HTTP 客户端,设置合理的超时(连接 3 秒、读取 8 秒),拿到响应后直接提取内容地址。这层不做过多的重试,失败就立刻降级,避免在这里浪费时间。
第二层是备用解析通道。这层的核心是“换一种问法”。同样的内容,不同的接口路径、不同的参数组合、不同的请求头特征,返回结果可能就不一样。这层我会做 2 到 3 次带退避的重试,每次间隔递增(1 秒、3 秒、7 秒),给频控留出恢复窗口。这层的超时设置比第一层宽松,因为它的目标不是快,是稳。
第三层是浏览器兜底。这层用无头浏览器环境,完整加载页面,等内容渲染完成后再提取。这层最慢,单条可能 10 到 20 秒,但成功率最高。关键是这层不能滥用,只有前两层都失败才触发。为了控制资源,浏览器实例要复用,不能每条任务开一个新实例。
2.3 SQLite 在架构中的角色定位
很多人把 SQLite 当成一个简单的存储,在这套架构里它的角色要重得多。它承担的是任务状态机的职责。每条任务在 SQLite 里有一条记录,字段包括:任务 ID、原始链接、当前层级、重试次数、最后错误码、内容本地路径、更新时间。工作流每次处理任务前先查状态,处理完更新状态。
这样做的好处是,整个工作流变成了无状态的。进程随时可以重启,重启后从 SQLite 里读出“待处理”和“处理中”的任务接着跑。我用db browser for sqlite做可视化排查,哪条任务卡在哪一层、错误码是什么,一目了然。相比用 JSON 文件记状态,SQLite 的并发读写和事务能力让多进程协作也变得可行。
注意:SQLite 的写操作要开 WAL 模式,否则多进程同时写会锁表。执行
PRAGMA journal_mode=WAL;一次即可,之后所有连接都受益。
3. 核心细节解析与实操要点
3.1 任务状态表的设计与字段说明
表结构设计直接决定了后续排查的效率。我用的建表语句是这样的:
CREATE TABLE IF NOT EXISTS tasks ( id INTEGER PRIMARY KEY AUTOINCREMENT, source_url TEXT NOT NULL UNIQUE, content_id TEXT, status TEXT DEFAULT 'pending', layer INTEGER DEFAULT 1, retry_count INTEGER DEFAULT 0, last_error TEXT, local_path TEXT, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX idx_status ON tasks(status); CREATE INDEX idx_layer ON tasks(layer);status字段的取值我定义了五种:pending(待处理)、processing(处理中)、success(成功)、failed(彻底失败)、skipped(跳过)。layer记录当前走到第几层,retry_count记录在当前层的重试次数。last_error存最后一次的错误信息,排查时直接看这个字段就知道问题出在哪。
source_url加了 UNIQUE 约束,这样重复提交同一个链接不会产生脏数据,用INSERT OR IGNORE就能天然去重。这个细节在批量导入链接时特别有用,省掉了应用层的去重逻辑。
3.2 第一层直连接口的参数与超时控制
第一层的请求参数有几个关键点。请求头里的User-Agent不能太假,用主流浏览器的真实 UA 字符串。Referer要带上内容平台的域名,很多接口会校验这个。超时设置上,连接超时给 3 秒,读取超时给 8 秒,整体不超过 12 秒。
import requests HEADERS = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://www.douyin.com/", "Accept": "application/json, text/plain, */*", } def fetch_layer1(content_id): url = f"https://api.example.com/detail?item_id={content_id}" try: resp = requests.get(url, headers=HEADERS, timeout=(3, 8)) if resp.status_code == 200: data = resp.json() if data.get("url"): return data["url"] return None except requests.RequestException as e: return None这层不抛异常,统一返回 None 表示失败,由上层决定是否降级。这样做的好处是调用方逻辑简单,不用处理各种异常类型。
3.3 第二层备用通道的退避重试策略
第二层的重试不是简单的循环,而是带指数退避的。第一次失败等 1 秒,第二次等 3 秒,第三次等 7 秒。这个间隔是根据频控恢复时间估算的,太短了没用,太长了拖慢整体进度。
import time BACKOFF = [1, 3, 7] def fetch_layer2(content_id): for i, wait in enumerate(BACKOFF): if i > 0: time.sleep(wait) result = try_alternate_parse(content_id) if result: return result return Nonetry_alternate_parse里换的是接口路径和参数组合。比如主接口用item_id,备用通道可以试aweme_id,或者换一个返回格式。这层的核心思想是“东方不亮西方亮”,同一个内容,平台往往有多个入口能拿到。
3.4 第三层浏览器兜底的资源管理
浏览器兜底最怕的就是资源泄漏。我的做法是维护一个浏览器实例池,池子里最多 3 个实例,任务来了从池子里取,用完还回去。实例空闲超过 5 分钟就销毁重建,避免长时间运行后内存膨胀。
from playwright.sync_api import sync_playwright class BrowserPool: def __init__(self, size=3): self.size = size self.pool = [] self.playwright = sync_playwright().start() self.browser = self.playwright.chromium.launch(headless=True) def acquire(self): if self.pool: return self.pool.pop() return self.browser.new_context() def release(self, ctx): if len(self.pool) < self.size: self.pool.append(ctx) else: ctx.close()用 Playwright 而不是 Selenium,主要是因为它对无头模式的支持更干净,启动更快,而且自带等待机制。页面加载用wait_until="networkidle",等网络请求都静下来再提取内容,成功率比固定 sleep 高得多。
提示:浏览器兜底这层一定要设总超时,我设的是 25 秒。超过就强制关闭上下文,标记任务失败,不能让一条任务把整个工作流拖死。
4. 完整实操流程与关键环节实现
4.1 环境准备与依赖安装
先把基础环境搭起来。Python 3.10 以上,装这几个核心依赖:
pip install requests playwright aiohttp playwright install chromiumplaywright install chromium这步不能省,它会把浏览器内核下下来。如果网络环境下载慢,可以设置镜像源,但这里不展开。SQLite 是 Python 内置的,不用额外装。可视化工具db browser for sqlite单独下载安装,排查问题时用。
目录结构我习惯这样组织:
project/ main.py layers/ layer1.py layer2.py layer3.py db/ tasks.db output/ logs/分层放代码,后面哪层出问题改哪层,不会互相影响。
4.2 任务初始化与批量导入
批量导入链接的时候,用INSERT OR IGNORE一次性写入,避免逐条查询判断是否存在。
import sqlite3 def import_urls(urls): conn = sqlite3.connect("db/tasks.db") conn.execute("PRAGMA journal_mode=WAL;") cur = conn.cursor() for url in urls: cur.execute( "INSERT OR IGNORE INTO tasks (source_url, status) VALUES (?, 'pending')", (url,) ) conn.commit() conn.close()导入 1000 条链接实测不到 1 秒。导入完可以用db browser for sqlite打开看看,确认数据都进去了。
4.3 主工作流循环与层级调度
主循环的逻辑是:取一条 pending 任务,标记为 processing,然后按层级依次尝试,成功就更新状态为 success,三层都失败就标记 failed。
def process_task(task): content_id = extract_id(task["source_url"]) if not content_id: update_status(task["id"], "failed", error="invalid_url") return # 第一层 url = fetch_layer1(content_id) if url: save_content(url, task["id"]) update_status(task["id"], "success", layer=1) return # 第二层 url = fetch_layer2(content_id) if url: save_content(url, task["id"]) update_status(task["id"], "success", layer=2) return # 第三层 url = fetch_layer3(content_id) if url: save_content(url, task["id"]) update_status(task["id"], "success", layer=3) return update_status(task["id"], "failed", error="all_layers_failed")每层成功后记录走到第几层,这个数据后面做统计分析很有用。我跑完一批任务后会查一下各层的成功占比,如果第三层占比超过 20%,说明前两层需要优化了。
4.4 断点续跑与状态恢复
进程重启后,先把所有processing状态的任务重置为pending,因为上次处理到一半的肯定没完成。
def recover_interrupted(): conn = sqlite3.connect("db/tasks.db") conn.execute( "UPDATE tasks SET status='pending' WHERE status='processing'" ) conn.commit() conn.close()这个操作放在工作流启动时执行一次。有了这个机制,我可以在任何时候 Ctrl+C 停掉进程,改完代码重新跑,之前处理过的不会重复,没处理完的接着处理。
4.5 内容落盘与命名规范
内容文件名用content_id加时间戳,避免重名。落盘前先检查文件是否已存在,存在就跳过,省一次下载。
import os def save_content(url, task_id): filename = f"output/{task_id}_{int(time.time())}.mp4" if os.path.exists(filename): return filename resp = requests.get(url, stream=True, timeout=(5, 30)) with open(filename, "wb") as f: for chunk in resp.iter_content(chunk_size=8192): f.write(chunk) return filename下载这步也要设超时,而且要用流式写入,不能一次性读进内存,大文件会把内存撑爆。
5. 常见问题与排查技巧实录
5.1 各层典型错误码与对应处理
跑多了之后,错误码基本能背下来。整理成表方便对照:
| 错误现象 | 可能层级 | 原因 | 处理方式 |
|---|---|---|---|
| 返回 403 | 第一层 | 请求头特征被识别 | 降级到第二层,换请求头 |
| 返回空 JSON | 第一层 | 参数名不对 | 降级到第二层,换参数组合 |
| 连接超时 | 第一层 | 网络抖动 | 直接降级,不重试 |
| 连续 429 | 第二层 | 频控触发 | 加大退避间隔,或暂停该批次 |
| 页面加载超时 | 第三层 | 资源阻塞 | 检查是否被重定向到验证页 |
| 内容地址 404 | 下载阶段 | 地址过期 | 重新走一遍解析流程 |
这张表我贴在显示器边上,排查的时候直接对号入座,省得每次重新分析。
5.2 SQLite 锁表与并发写入问题
多进程跑的时候最容易遇到database is locked。根因是默认的 journal 模式不支持并发写。解决办法就一条:开 WAL。
PRAGMA journal_mode=WAL; PRAGMA busy_timeout=5000;busy_timeout设 5 秒,遇到锁的时候会等而不是立刻报错。这两个设置加上之后,我 4 个进程同时跑没再出现过锁表。
5.3 浏览器兜底层的内存泄漏排查
浏览器层跑久了内存会涨,这是无头浏览器的通病。我的做法是每处理 50 条任务就重建一次浏览器实例,不管有没有异常。这个阈值是试出来的,50 条以内内存增长可控,超过就开始明显。
class BrowserPool: def __init__(self, size=3, max_uses=50): self.max_uses = max_uses self.use_count = 0 # ... 其他初始化 def maybe_recycle(self): self.use_count += 1 if self.use_count >= self.max_uses: self.restart() self.use_count = 0重建的时候先把池子里的上下文全关掉,再关浏览器,最后重新 launch。顺序不能反,反了会残留进程。
5.4 任务卡在 processing 状态的定位方法
有时候任务会卡在 processing 不动,既不成功也不失败。这种情况用db browser for sqlite查updated_at字段,看最后更新时间。如果超过 10 分钟没更新,基本可以判定是卡死了。
SELECT * FROM tasks WHERE status='processing' AND updated_at < datetime('now', '-10 minutes');查出来的任务手动重置为 pending 重新跑。为了自动化,我在主循环里加了个看门狗,每 5 分钟扫一次,超时的自动重置。
5.5 提升整体吞吐的几条实操心得
第一,第一层的超时要设短。我一开始设的 15 秒,结果大量任务卡在第一层等超时,整体吞吐上不去。改成 8 秒后,降级更果断,整体反而快了。
第二,第二层的退避间隔不要设太长。我试过 5 秒、15 秒、30 秒的退避,结果一批任务跑了一晚上。改成 1、3、7 之后,成功率没降,速度翻倍。
第三,第三层要限流。浏览器层并发太高会互相抢资源,反而都慢。我限制同时最多 3 个浏览器上下文在跑,超出的排队等。
第四,SQLite 的写入要批量提交。每条任务都 commit 一次太慢,我改成每 20 条 commit 一次,性能提升明显。但要注意,批量提交意味着崩溃时可能丢最后一批的状态,所以批量大小不能太大,20 是个平衡点。
第五,日志要记全。每条任务的层级切换、错误码、耗时都记下来,后面分析瓶颈全靠这些数据。我用的是标准 logging 模块,按天切文件,跑一周下来能清楚看到哪层是瓶颈。
这套架构跑下来,500 条链接的成功率从原来的 76% 提到了 98% 以上,平均单条耗时从 12 秒降到了 4 秒左右。最关键的是,它让我从“盯着脚本跑”变成了“扔进去不用管”,这才是工作流该有的样子。