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

能否将异步函数定义为Celery任务?异步爬虫遇网络异常求助

回答:可以将协程包装为Celery任务,但你的实现有几个关键问题需要修复

首先明确:Celery本身不直接支持协程作为任务函数,但你可以在一个普通的同步Celery任务函数内部,启动asyncio事件循环来运行你的异步爬虫协程——这正是你代码里尝试的方向,这个思路完全可行,你的网络错误大概率是代码实现细节导致的,而非技术栈兼容性问题。

下面是你的代码里需要调整的几个核心问题:

1. 事件循环的创建方式不安全

你当前用asyncio.get_event_loop()获取事件循环,但在Celery的多进程worker环境中,这个方法可能会复用已经关闭或异常的循环,导致奇怪的网络错误。推荐使用Python 3.7+引入的asyncio.run(),它会自动创建、管理和关闭事件循环,更简洁可靠:

@celery_app.task
def run_crawler(url):
    asyncio.run(url_crawler(url))

如果你的Python版本低于3.7,应该显式创建新的事件循环:

@celery_app.task
def run_crawler(url):
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        loop.run_until_complete(url_crawler(url))
    finally:
        loop.close()

2. 共享状态的并发安全问题

你的代码里有两个共享状态变量存在问题:

  • urls字典:多个协程同时读写,虽然asyncio是单线程并发,但字典不是协程安全的(协程切换时可能导致数据损坏),建议用asyncio.Lock来保护对urls的修改:
    # 在url_crawler里创建锁
    url_lock = asyncio.Lock()
    
    # 在url_worker里修改urls时加锁
    async with url_lock:
        if link not in urls:
            urls[link] = depth + 1
            await queue.put(link)
    
  • request_count整数:因为整数是不可变类型,你在url_worker里的request_count +=1不会修改外层的变量(只是创建了一个新的局部变量),导致最终统计的请求数错误。可以用一个可变的容器(比如列表)来包装它:
    # 在url_crawler里初始化
    request_count = [1]
    
    # 在url_worker里修改
    request_count[0] += 1
    

3. 并发控制可能过度

Celery的worker默认是多进程模型,如果你同时设置了较多的Celery worker进程,再加上asyncio的max_concurrency,会导致总并发请求数过高,触发目标网站的反爬机制或者本地网络的连接限制。建议:

  • 调整Celery worker的进程数(比如celery worker --concurrency=4)
  • 降低asyncio的max_concurrency值,避免两者叠加导致并发量过大

4. 异常处理不够全面

你的url_worker只捕获了TimeoutError和OSError,但requests-html的AsyncHTMLSession底层基于aiohttp,还可能抛出aiohttp.ClientError等其他网络异常。建议扩展异常捕获范围,确保所有异常都被记录:

except (TimeoutError, OSError, aiohttp.ClientError) as e:
    logger.exception(e)
    queue.task_done()
    return

(记得先导入aiohttp)

修改后的核心代码片段

import asyncio
import aiohttp
from requests_html import AsyncHTMLSession
import celery
import time
import logging

logger = logging.getLogger(__name__)
config = ... # 你的配置对象
celery_app = celery.Celery(...) # 你的Celery实例

async def url_helper(resp):
    return resp.html.absolute_links

async def url_worker(session, queue, urls, request_count, url_lock):
    while True:
        try:
            current = await queue.get()
        except asyncio.QueueEmpty:
            return

        if not config.regex_url.match(current):
            queue.task_done()
            return

        depth = urls.get(current, 0)
        if depth == config.max_depth:
            logger.info("Out of depth: {}".format(current))
            queue.task_done()
            return

        try:
            logger.info("current link: {}".format(current))
            resp = await session.get(
                current, allow_redirects=True, timeout=15
            )
            new_urls = await url_helper(resp)
            # 修改request_count
            request_count[0] += 1
        except (TimeoutError, OSError, aiohttp.ClientError) as e:
            logger.exception(e)
            queue.task_done()
            return

        # 加锁修改urls
        async with url_lock:
            for link in new_urls:
                if config.regex_url.match(link) and link not in urls:
                    urls[link] = depth + 1
                    await queue.put(link)
        queue.task_done()

async def url_crawler(url):
    start = time.time()
    http_asession = AsyncHTMLSession()
    logger.debug("Sending GET request to start_url...")
    start_response = await http_asession.get(url=url)
    logger.debug("Received {} response.".format(start_response.status_code))
    request_count = [1]
    urls = {}
    start_urls = list(start_response.html.absolute_links)
    logger.info("start urls: %s", start_urls)
    queue = asyncio.Queue()
    max_workers = config.max_concurrency
    url_lock = asyncio.Lock()

    for url_item in start_urls:
        await queue.put(url_item)

    coros = asyncio.gather(
        *[
            url_worker(
                queue=queue,
                session=http_asession,
                urls=urls,
                request_count=request_count,
                url_lock=url_lock
            )
            for i in range(max_workers)
        ]
    )
    await coros
    await http_asession.close()
    logger.info("Crawled {} links.".format(len(urls)))
    logger.info(urls)
    logger.debug("Made {} HTTP requests.".format(request_count[0]))
    finish = time.time()
    logger.info("Execution time: {}".format(finish - start))

@celery_app.task
def run_crawler(url):
    asyncio.run(url_crawler(url))

总结

你的技术选型是合理的:用Celery管理分布式任务,用asyncio实现单进程内的异步爬取来提升性能。网络错误不是因为Celery和asyncio不兼容,而是代码里的事件循环管理、共享状态并发安全、并发控制和异常处理这些细节问题导致的。按照上面的调整修复后,应该能解决你遇到的问题。

内容的提问来源于stack exchange,提问作者moonesque

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 08:52:51