尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于SQLite与三层容错架构的内容获取工作流重构实践

基于SQLite与三层容错架构的内容获取工作流重构实践 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_modeWAL;一次即可之后所有连接都受益。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 fhttps://api.example.com/detail?item_id{content_id} try: resp requests.get(url, headersHEADERS, 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, size3): self.size size self.pool [] self.playwright sync_playwright().start() self.browser self.playwright.chromium.launch(headlessTrue) 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_untilnetworkidle等网络请求都静下来再提取内容成功率比固定 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_modeWAL;) 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, errorinvalid_url) return # 第一层 url fetch_layer1(content_id) if url: save_content(url, task[id]) update_status(task[id], success, layer1) return # 第二层 url fetch_layer2(content_id) if url: save_content(url, task[id]) update_status(task[id], success, layer2) return # 第三层 url fetch_layer3(content_id) if url: save_content(url, task[id]) update_status(task[id], success, layer3) return update_status(task[id], failed, errorall_layers_failed)每层成功后记录走到第几层这个数据后面做统计分析很有用。我跑完一批任务后会查一下各层的成功占比如果第三层占比超过 20%说明前两层需要优化了。4.4 断点续跑与状态恢复进程重启后先把所有processing状态的任务重置为pending因为上次处理到一半的肯定没完成。def recover_interrupted(): conn sqlite3.connect(db/tasks.db) conn.execute( UPDATE tasks SET statuspending WHERE statusprocessing ) conn.commit() conn.close()这个操作放在工作流启动时执行一次。有了这个机制我可以在任何时候 CtrlC 停掉进程改完代码重新跑之前处理过的不会重复没处理完的接着处理。4.5 内容落盘与命名规范内容文件名用content_id加时间戳避免重名。落盘前先检查文件是否已存在存在就跳过省一次下载。import os def save_content(url, task_id): filename foutput/{task_id}_{int(time.time())}.mp4 if os.path.exists(filename): return filename resp requests.get(url, streamTrue, timeout(5, 30)) with open(filename, wb) as f: for chunk in resp.iter_content(chunk_size8192): f.write(chunk) return filename下载这步也要设超时而且要用流式写入不能一次性读进内存大文件会把内存撑爆。5. 常见问题与排查技巧实录5.1 各层典型错误码与对应处理跑多了之后错误码基本能背下来。整理成表方便对照错误现象可能层级原因处理方式返回 403第一层请求头特征被识别降级到第二层换请求头返回空 JSON第一层参数名不对降级到第二层换参数组合连接超时第一层网络抖动直接降级不重试连续 429第二层频控触发加大退避间隔或暂停该批次页面加载超时第三层资源阻塞检查是否被重定向到验证页内容地址 404下载阶段地址过期重新走一遍解析流程这张表我贴在显示器边上排查的时候直接对号入座省得每次重新分析。5.2 SQLite 锁表与并发写入问题多进程跑的时候最容易遇到database is locked。根因是默认的 journal 模式不支持并发写。解决办法就一条开 WAL。PRAGMA journal_modeWAL; PRAGMA busy_timeout5000;busy_timeout设 5 秒遇到锁的时候会等而不是立刻报错。这两个设置加上之后我 4 个进程同时跑没再出现过锁表。5.3 浏览器兜底层的内存泄漏排查浏览器层跑久了内存会涨这是无头浏览器的通病。我的做法是每处理 50 条任务就重建一次浏览器实例不管有没有异常。这个阈值是试出来的50 条以内内存增长可控超过就开始明显。class BrowserPool: def __init__(self, size3, max_uses50): 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 statusprocessing 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 秒左右。最关键的是它让我从“盯着脚本跑”变成了“扔进去不用管”这才是工作流该有的样子。
返回列表