TaskSpider:基于任务生命周期的轻量级Python爬虫骨架
简介这是一份面向Python初学者与中级开发者设计的轻量级网络爬虫框架源码包聚焦于降低爬虫开发门槛、提升代码可维护性与任务复用性。资源以TaskSpider框架为核心通过模块化Task节点、封装网络请求与HTML解析逻辑、支持BigTask并发处理及TaskMessage参数传递机制显著简化复杂爬取任务的编码工作。压缩包共60个文件含29个核心Python源码如NetworkTask、BigTask、Anlyst等、17个编译缓存文件、11个CSV数据样例及1个Excel模板辅以LICENSE、README和完整测试用例Test目录结构清晰、开箱即用整体体积仅760KB便于快速部署与学习调试。目前已有35人下载学习适合希望掌握面向对象爬虫架构、理解任务调度与数据流设计的实践者可直接运行Sample中的graduateMessageSample案例快速上手从URL配置、页面抓取到结构化存储的全流程。1. 这不是又一个 Requests 封装TaskSpider 是面向任务生命周期的轻量级爬虫骨架你写过这样的代码吗用requests.get()抓一页BeautifulSoup解析正则提取字段再pandas.to_csv()存文件——单页能跑通但加到 50 个 URL 就得手动改循环加个重试逻辑就得在每个请求前塞try/except想把解析结果传给下游清洗模块得自己拼字典、传参数、处理 None。这不是爬虫是胶水代码流水线。TaskSpider 的定位很明确它不替代requests或lxml也不试图做成 Scrapy 那样的全栈框架而是提供一套任务声明式定义 生命周期钩子 消息驱动传递的骨架。核心抽象是Task类——每个爬虫行为被建模为一个可配置、可复用、可组合的任务节点BigTask负责批量调度TaskMessage作为唯一数据载体在NetworkTask → Anlyst → Writer链路中流动。它解决的不是“怎么发请求”而是“怎么让 20 个不同网站的爬取逻辑共享同一套错误重试、超时控制、日志埋点和失败回溯机制”。适合需要快速交付多个中小型爬虫项目、又不愿重复造轮子的 Python 开发者尤其对刚脱离脚本阶段、开始接触工程化爬虫逻辑的 1–3 年经验者友好。2. 从Task类继承开始理解任务生命周期与消息契约2.1 为什么 TaskSpider 不直接封装 requests——解耦网络层与业务逻辑TaskSpider 的设计哲学是“职责分离”。它把网络访问NetWork.py、URL 管理URLs.py、解析Anlyst.py、写入Writer.py全部拆成独立模块而Task类本身只定义四个标准钩子方法class Task: def before_run(self, message: TaskMessage) - TaskMessage: 任务执行前预处理可修改 message 或抛出异常终止流程 return message def run(self, message: TaskMessage) - TaskMessage: 核心执行逻辑必须返回 TaskMessage 实例 raise NotImplementedError def after_run(self, message: TaskMessage) - TaskMessage: 任务成功后处理如日志记录、状态标记 return message def on_error(self, message: TaskMessage, error: Exception) - TaskMessage: 异常捕获后处理可重试、降级或标记失败 return message提示TaskMessage是唯一的数据容器所有任务间通信都通过它完成。它本质是一个带类型提示的dict子类预置了url,content,status_code,parsed_data,error_info等字段。你不该在run()中直接print()或open()文件而应把结果存入message.parsed_data由下游Writer统一落盘。这种设计带来两个实际好处一是测试友好——你可以 mockmessage输入单独验证Anlyst.run()的解析逻辑无需真实 HTTP 请求二是可组合性强——比如EmailTask可以监听Writer成功后的message自动触发告警而不用侵入爬取主流程。2.2 实战三步写出第一个可运行的 NetworkTask我们以抓取豆瓣电影 Top250 第一页为例演示如何基于NetworkTask快速构建任务2.2.1 定义任务类并覆盖run()方法# my_first_task.py from Task import NetworkTask from TaskMessage import TaskMessage class DoubanTop250Task(NetworkTask): def run(self, message: TaskMessage) - TaskMessage: # 1. 设置目标 URL注意NetworkTask 自动调用 self._request() message.url https://movie.douban.com/top250 # 2. 执行请求NetworkTask 内部已封装 requests.Session 默认 headers response self._request(message) if response.status_code ! 200: raise RuntimeError(fHTTP {response.status_code} for {message.url}) # 3. 将响应体注入 message.content供下游解析 message.content response.text message.status_code response.status_code return message2.2.2 配置并执行任务实例# main.py from my_first_task import DoubanTop250Task from TaskMessage import TaskMessage if __name__ __main__: # 创建空消息对象 msg TaskMessage() # 初始化任务并执行 task DoubanTop250Task() result task.execute(msg) # execute() 自动调用 before_run → run → after_run print(fStatus: {result.status_code}) print(fContent length: {len(result.content)} chars)2.2.3 关键参数说明与调试技巧参数位置说明常见修改场景self.timeoutNetworkTask基类属性默认 10 秒单位秒抓取慢网站时设为 30self.retry_timesNetworkTask基类属性默认 3 次重试对高防站点设为 5并配合on_error自定义退避策略self.session.headersNetworkTask._request()内部已预置User-Agent可追加Referer防反爬需添加Cookie或X-Requested-Withmessage.url任务run()中赋值必须为字符串支持http://和https://动态 URL 可从message.params中读取如message.url fhttps://api.example.com/{message.params[id]}注意NetworkTask.execute()返回的是新TaskMessage实例原msg对象不会被修改。这是为了保证任务链路的不可变性immutability避免上游任务意外污染下游数据。3. 构建完整爬虫链路NetworkTask → Anlyst → Writer 的协同机制3.1 解析层Anlyst如何安全提取结构化数据Anlyst.py提供了基于lxml的 XPath 和 CSS Selector 两种解析器但关键在于它强制要求所有解析结果必须通过TaskMessage.parsed_data字段返回且类型必须为dict。这消除了下游对数据格式的猜测成本。3.1.1 编写豆瓣电影标题解析器# douban_anlyst.py from Anlyst import Anlyst from TaskMessage import TaskMessage class DoubanMovieAnlyst(Anlyst): def run(self, message: TaskMessage) - TaskMessage: # 1. 确保上游已提供 content if not message.content: raise ValueError(No content to parse. Check upstream NetworkTask.) # 2. 使用 lxml 解析 HTMLAnlyst 基类已初始化 etree.HTMLParser tree self._parse_html(message.content) # 3. 提取标题列表XPath 示例 titles tree.xpath(//div[classhd]/a/span[1]/text()) # 4. 提取评分CSS Selector 示例 ratings tree.cssselect(div.star span.rating_num) scores [r.text.strip() for r in ratings] # 5. 严格按契约返回 dict message.parsed_data { titles: titles[:10], # 取前 10 条 scores: scores[:10], total_count: len(titles) } return message3.1.2Anlyst的容错设计要点self._parse_html()内部已处理编码自动检测chardet但若页面声明meta charsetgbk而实际是 UTF-8仍可能乱码。此时应在before_run()中手动设置message.encoding utf-8。XPath 表达式失败时tree.xpath()返回空列表不会抛异常——这是刻意设计避免单条数据缺失导致整个任务中断。cssselect()方法来自lxml.cssselect比 BeautifulSoup 的select()性能高约 30%但不支持伪类如:nth-child复杂选择需用 XPath。3.2 写入层Writer的三种落地方式与性能权衡Writer.py提供CSVWriter,JSONWriter,SQLiteWriter三个子类它们共享同一套write()接口但内部实现差异显著写入方式适用场景单次写入耗时万行并发安全注意事项CSVWriter快速导出、Excel 查看~120ms❌多进程需加锁字段名必须提前定义self.fields [title, score]JSONWriter结构化数据存档、API 响应~85ms✅纯内存操作message.parsed_data必须是 JSON 序列化安全类型无 datetime、bytesSQLiteWriter需要索引查询、增量更新~210ms✅事务隔离首次运行自动建表表结构由message.parsed_data的 key 推断3.2.1 使用 SQLiteWriter 实现去重入库# db_writer.py from Writer import SQLiteWriter from TaskMessage import TaskMessage class MovieDBWriter(SQLiteWriter): def run(self, message: TaskMessage) - TaskMessage: # 1. 确保 parsed_data 是 list of dictSQLiteWriter 要求 if not isinstance(message.parsed_data, list): raise TypeError(parsed_data must be list for SQLiteWriter) # 2. 添加唯一标识字段用于 ON CONFLICT IGNORE for item in message.parsed_data: item[url_hash] hash(item.get(title, )) # 简单哈希生产环境建议用 sha256 # 3. 执行写入自动建表主键为 url_hash self.write( datamessage.parsed_data, table_namemovies, unique_keyurl_hash ) return message提示SQLiteWriter.write()的unique_key参数会生成INSERT OR IGNORE INTO ...语句避免重复插入。若需更新已有记录应改用INSERT OR REPLACE此时需在run()中手动构造 SQL而非调用self.write()。4. 并发加速与任务编排BigTask 的批量调度与资源控制4.1 BigTask 不是简单 for 循环它管理连接池与任务队列BigTask的核心价值在于统一管控并发粒度与资源配额。它不直接启动线程/进程而是通过concurrent.futures.ThreadPoolExecutor封装任务执行并提供max_workers、timeout_per_task、retry_policy三层控制# batch_crawl.py from Task import BigTask from my_first_task import DoubanTop250Task from douban_anlyst import DoubanMovieAnlyst from db_writer import MovieDBWriter from TaskMessage import TaskMessage # 构建 10 个不同分页的 URL 列表 urls [fhttps://movie.douban.com/top250?start{i*25}filter for i in range(10)] # 初始化 BigTask指定最大并发数为 5 big_task BigTask( tasks[ DoubanTop250Task(), DoubanMovieAnlyst(), MovieDBWriter() ], max_workers5, # 同时最多 5 个线程执行 NetworkTask timeout_per_task30, # 单个任务超时 30 秒 retry_policy{max_retries: 2} # 每个失败任务重试 2 次 ) # 批量执行每个 URL 生成一个独立 TaskMessage messages [TaskMessage(urlu) for u in urls] results big_task.execute_batch(messages) # 统计成功/失败数量 success_count sum(1 for r in results if r.is_success()) fail_count len(results) - success_count print(fSuccess: {success_count}, Fail: {fail_count})4.1.1execute_batch()的底层执行流程任务拆分将messages列表按max_workers分组每组进入一个ThreadPoolExecutor.submit()链路执行每个线程内依次调用DoubanTop250Task.run()→DoubanMovieAnlyst.run()→MovieDBWriter.run()形成串行链路错误聚合任一环节抛异常on_error()被触发message.error_info记录堆栈message.is_success()返回False结果归集所有线程完成后返回TaskMessage列表顺序与输入messages一致。注意BigTask的max_workers控制的是并发任务数不是并发请求数。若每个NetworkTask内部使用aiohttp异步请求则需替换NetworkTask基类当前版本仅支持同步阻塞式请求。4.2 任务依赖与条件跳过用TaskMessage.params实现动态流程TaskMessage.params是一个自由dict专用于携带非标准字段。它让任务链具备条件分支能力# conditional_task.py from Task import Task from TaskMessage import TaskMessage class ConditionalWriter(Task): def run(self, message: TaskMessage) - TaskMessage: # 根据上游解析结果决定是否写入 if message.parsed_data.get(total_count, 0) 5: message.params[skip_write] True return message # 正常写入逻辑 with open(output.txt, a) as f: f.write(str(message.parsed_data) \n) return message # 在链路中插入此任务 big_task BigTask(tasks[ DoubanTop250Task(), DoubanMovieAnlyst(), ConditionalWriter(), # 位于 Writer 之前 MovieDBWriter() ])这种模式避免了在Writer内部写if判断保持各任务职责单一。params字段在整条链路中透传任何任务都可读写是 TaskSpider 实现轻量级工作流的核心机制。5. 生产环境必备日志追踪、失败回溯与性能调优技巧5.1 用TaskMessage.trace_id实现端到端请求追踪每个TaskMessage实例在创建时自动生成trace_idUUID4该 ID 会贯穿整个任务链路。你可以在before_run()和after_run()中打印它快速定位某次失败请求的完整执行路径# logging_task.py import logging from Task import Task logger logging.getLogger(TaskSpider) class LoggingTask(Task): def before_run(self, message: TaskMessage) - TaskMessage: logger.info(f[{message.trace_id}] Start task: {self.__class__.__name__}) return message def after_run(self, message: TaskMessage) - TaskMessage: logger.info(f[{message.trace_id}] Finish task: {self.__class__.__name__}) return message配合logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s)日志输出形如2024-06-15 14:22:33,123 - TaskSpider - INFO - [a1b2c3d4-5678-90ef-ghij-klmnopqrstuv] Start task: DoubanTop250Task 2024-06-15 14:22:35,456 - TaskSpider - INFO - [a1b2c3d4-5678-90ef-ghij-klmnopqrstuv] Finish task: DoubanTop250Task提示trace_id也出现在message.error_info的异常信息中可直接 grep 日志文件定位问题链路。5.2 性能瓶颈诊断三类典型慢任务与优化方案瓶颈类型表现诊断命令优化方案DNS 解析慢NetworkTask首次请求耗时 5stime nslookup movie.douban.com在NetworkTask.__init__()中设置self.session.mount(http://, requests.adapters.HTTPAdapter(pool_connections10))复用连接池解析慢Anlyst.run()耗时占比 70%python -m cProfile -o profile.out your_script.py改用lxml.etree.iterparse()流式解析大 HTML或预编译 XPath 表达式self.title_xpath etree.XPath(//div[classhd]/a/span[1]/text())写入慢SQLiteWriter.write()单次 200mssqlite3 your.db EXPLAIN QUERY PLAN INSERT INTO movies...为url_hash字段添加索引CREATE INDEX idx_url_hash ON movies(url_hash);5.2.1 用setup.py验证依赖完整性项目根目录的setup.py不仅用于pip install更是运行时依赖检查入口。执行以下命令可验证环境是否完备# 检查必需依赖是否安装 python -c import requests, lxml, chardet, sqlite3; print(All core deps OK) # 检查可选依赖如需 EmailTask python -c import smtplib; print(Email support OK)若报ModuleNotFoundError按setup.py中install_requires字段安装pip install requests lxml chardetsetup.py中未声明pandas或numpy因为 TaskSpider 故意规避重量级依赖确保在树莓派等资源受限设备上也能运行。5.3 测试驱动开发用Test/目录快速验证任务逻辑项目自带Test/目录包含NetWorkTest.py等单元测试。运行测试前需先安装pytestpip install pytest cd Test pytest --tbshort -v测试用例采用unittest.mock.patch模拟网络请求例如NetWorkTest.py中patch(requests.Session.get) def test_network_task_success(self, mock_get): mock_get.return_value.status_code 200 mock_get.return_value.text htmlbodytest/body/html task NetworkTask() msg TaskMessage(urlhttp://example.com) result task.execute(msg) self.assertEqual(result.status_code, 200) self.assertIn(test, result.content)提示所有测试用例均不发起真实网络请求因此可在 CI 环境中稳定运行。新增任务类后务必在Test/下添加对应测试文件覆盖before_run,run,on_error三个关键路径。真正让 TaskSpider 在中小型爬虫项目中立住脚的不是它多快或多强而是它用TaskMessage这一根线把请求、解析、存储、通知这些原本散落在各处的胶水代码串成一条可追溯、可插拔、可替换的流水线。当你下次面对 15 个不同结构的网站抓取需求时不必重写 15 套requestsbs4csv只需继承NetworkTask、Anlyst、Writer填入各自的 XPath 和字段映射——剩下的重试、超时、日志、并发框架已为你托底。本文还有配套的精品资源点击获取