能否将异步函数定义为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

