Scrapy Pipeline集成aiomysql协程存数据遇pending Task报错求助
问题原因分析
偶现的「Task was destroyed but it is pending」错误主要有几个触发点:
- 手动转换asyncio Future到Twisted Deferred时,未对创建的Task持有强引用,导致GC提前回收pending状态的任务
- 爬虫关闭时,aiomysql连接池未等待所有正在执行的数据库操作任务完成就被关闭
- 使用的aiomysql版本(0.1.1)较旧,存在协程任务管理的bug
- 自定义的
as_deferred函数未正确适配Scrapy与asyncio的事件循环交互
解决方案及修正后的代码
1. 利用Scrapy原生的异步Pipeline支持
Scrapy 2.2+已经原生支持async/await风格的Pipeline方法,无需手动转换Deferred,这能避免大部分事件循环适配问题。
2. 跟踪所有数据库操作任务,确保爬虫关闭时全部完成
通过维护一个任务集合,跟踪所有process_item中启动的数据库操作任务,在爬虫关闭前等待所有任务完成,再关闭连接池。
3. 升级aiomysql版本
旧版本aiomysql存在协程相关的bug,建议升级到最新稳定版(如0.2.7+)。
修正后的完整示例代码:
# 确保已启用TWISTED_REACTOR = "twisted.internet.asyncioreactor.AsyncioSelectorReactor" import aiomysql import asyncio from scrapy.exceptions import DropItem class AsyncMysqlPipeline: def __init__(self): self.pool = None self._pending_tasks = set() async def open_spider(self, spider): self.pool = await aiomysql.create_pool( host="localhost", port=3306, user="root", password="pwd", db="db", minsize=5, maxsize=20, loop=asyncio.get_event_loop() ) spider.logger.info("MySQL连接池初始化完成") async def process_item(self, item, spider): # 创建任务并加入跟踪集合 task = asyncio.create_task(self._save_item(item, spider)) self._pending_tasks.add(task) # 任务完成后从集合移除 task.add_done_callback(self._pending_tasks.discard) try: await task return item except Exception as e: spider.logger.error(f"存储Item失败: {str(e)}") raise DropItem(f"存储失败: {str(e)}") async def _save_item(self, item, spider): async with self.pool.acquire() as conn: async with conn.cursor() as cursor: # 替换为实际的SQL语句 sql = "INSERT INTO your_table (col1, col2) VALUES (%s, %s)" values = (item["col1"], item["col2"]) await cursor.execute(sql, values) await conn.commit() spider.logger.debug(f"成功存储Item: {item['id']}") async def close_spider(self, spider): if self._pending_tasks: spider.logger.info(f"等待{len(self._pending_tasks)}个未完成的存储任务...") # 等待所有pending任务完成 await asyncio.gather(*self._pending_tasks) if self.pool: self.pool.close() await self.pool.wait_closed() spider.logger.info("MySQL连接池已关闭")
关键改进点说明
- 移除了自定义的
as_deferred函数,直接使用Scrapy原生支持的async Pipeline方法 - 新增
_pending_tasks集合跟踪所有数据库存储任务,确保爬虫关闭时等待全部任务完成 - 优化了连接池的初始化参数(设置合理的min/max连接数)
- 增加了日志输出,方便调试问题
- 在
process_item中捕获异常,避免单个Item存储失败导致整个爬虫崩溃
其他可选方案:使用Scrapy官方推荐的数据库集成方式
如果不想手动处理asyncio与Scrapy的交互,也可以使用以下方案:
- 使用
scrapy-pymysql等成熟的第三方Pipeline库,直接同步操作数据库(Scrapy会自动在线程池中执行,不阻塞事件循环) - 使用Twisted的
adbapi连接MySQL,这是Scrapy传统的异步数据库操作方式,稳定性更高:
from twisted.enterprise import adbapi import pymysql class TwistedMysqlPipeline: def __init__(self, db_args): self.db_pool = adbapi.ConnectionPool( "pymysql", **db_args ) @classmethod def from_crawler(cls, crawler): db_args = { "host": crawler.settings.get("MYSQL_HOST"), "port": crawler.settings.get("MYSQL_PORT"), "user": crawler.settings.get("MYSQL_USER"), "password": crawler.settings.get("MYSQL_PASSWORD"), "db": crawler.settings.get("MYSQL_DB"), "charset": "utf8mb4" } return cls(db_args) def process_item(self, item, spider): query = self.db_pool.runInteraction(self._save_item, item, spider) query.addErrback(self._handle_error, item, spider) return item def _save_item(self, cursor, item, spider): sql = "INSERT INTO your_table (col1, col2) VALUES (%s, %s)" cursor.execute(sql, (item["col1"], item["col2"])) def _handle_error(self, failure, item, spider): spider.logger.error(f"存储失败: {failure}")
内容的提问来源于stack exchange,提问作者ayuge
相关产品推荐
相关产品推荐

