You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 07:18:14