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

Prefect多线程流中使用async_playwright遇_channel为None错误求助

解决Prefect中Playwright并行爬取的_channel为None错误

我需要爬取约800个URL,计划分成200个URL的块并行执行,块内URL串行处理。使用async_playwright,且明确Playwright非线程安全,需为每个线程创建独立实例,但运行时出现错误,根源是Playwright的_browser_context.py中return from_channel(await self._channel.send("newPage"))行的self._channel为None。

原代码示例

from playwright.async_api import async_playwright
import asyncio
import re
from prefect import task, flow
from prefect.cache_policies import NONE

# Each thread uses its own Playwright instance as Playwright is not thread-safe
class SpiderThread:
    def __init__(self, p, browser, context):
        self.p = p
        self.browser = browser
        self.context = context

    @task(cache_policy=NONE) # Cache policy needed to disable pickling
    def run(self): # This will later run through hundreds of websites
        page = asyncio.run(self.context.new_page())
        response = asyncio.run(page.goto("https://www.ard.de/"))
        extracted_title = re.search('<title>(.*)</title>', response.text()).group(1)
        # Should print/return "ARD"
        print(extracted_title)
        return extracted_title

    @classmethod
    async def create(cls):
        p = await async_playwright().start()
        chromium = p.chromium
        browser = await chromium.launch()
        context = await browser.new_context()
        return SpiderThread(p, browser, context)

# This is the main flow that runs multiple SpiderThreads in parallel
@flow(name="Main")
def main_flow():
    threads= [asyncio.run(SpiderThread.create()), asyncio.run(SpiderThread.create())]
    print(f"Threads: {threads}")

    prefect_threads = []
    for t in threads:
        prefect_threads.append(t.run.submit())
    # Wait for all of them to finish
    print(f"Prefect threads: {prefect_threads}")
    results = [prefect_thread.result() for prefect_thread in prefect_threads]
    print(results)

if __name__ == "__main__":
    main_flow()

错误信息

12:20:29.653 | INFO    | Flow run 'grumpy-meerkat' - Beginning flow run 'grumpy-meerkat' for flow 'Main'
12:20:29.659 | INFO    | Flow run 'grumpy-meerkat' - View at http://127.0.0.1:4200/runs/flow-run/6e87d290-5ffa-43f2-9258-c6c9f6387ffb
Threads: [<__main__.SpiderThread object at 0x000002282AA90FD0>, <__main__.SpiderThread object at 0x000002282AA929B0>]
Prefect threads: [<prefect.futures.PrefectConcurrentFuture object at 0x000002282AE36D70>, <prefect.futures.PrefectConcurrentFuture object at 0x000002282AE37A90>]
12:20:30.795 | ERROR   | Task run 'run-f28' - Task run failed with exception: AttributeError("BrowserContext.new_page: 'NoneType' object has no attribute 'send'") - Retries are exhausted
Traceback (most recent call last):
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\prefect\task_engine.py", line 819, in run_context
    yield self
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\prefect\task_engine.py", line 1395, in run_task_sync
    engine.call_task_fn(txn)
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\prefect\task_engine.py", line 836, in call_task_fn
    result = call_with_parameters(self.task.fn, parameters)
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\prefect\utilities\callables.py", line 210, in call_with_parameters
    return fn(*args, **kwargs)
  File "C:\Users\mn\IdeaProjects\scraper\playwright-tests-prefect2.py", line 16, in run
    page = asyncio.run(self.context.new_page())
  File "C:\Users\mn\AppData\Local\Programs\Python\Python310\lib\asyncio\runners.py", line 44, in run
    return loop.run_until_complete(main)
  File "C:\Users\mn\AppData\Local\Programs\Python\Python310\lib\asyncio\base_events.py", line 641, in run_until_complete
    return future.result()
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\playwright\async_api\_generated.py", line 12787, in new_page
    return mapping.from_impl(await self._impl_obj.new_page())
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\playwright\_impl\_browser_context.py", line 325, in new_page
    return from_channel(await self._channel.send("newPage"))
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\playwright\_impl\_connection.py", line 61, in send
    return await self._connection.wrap_api_call(
  File "C:\Users\mn\IdeaProjects\scraper\venv\lib\site-packages\playwright\_impl\_connection.py", line 528, in wrap_api_call
    raise rewrite_error(error, f"{parsed_st['apiName']}: {error}") from None
AttributeError: BrowserContext.new_page: 'NoneType' object has no attribute 'send'

错误原因

  1. 同步任务嵌套异步调用:run方法是同步Prefect任务,但内部用asyncio.run调用Playwright的异步方法,导致Playwright的上下文在不同事件循环中被访问,连接被意外销毁,_channel变为None。
  2. 类实例绑定任务的问题:Prefect任务绑定在SpiderThread实例上,任务执行时可能在不同进程/线程中运行,实例的Playwright资源无法正确传递和维护。
  3. 手动管理Playwright实例:在Flow中手动创建SpiderThread实例,不符合Prefect的任务调度模型,导致资源生命周期管理混乱。

修正后的代码

from playwright.async_api import async_playwright
import re
from prefect import task, flow
from prefect.cache_policies import NONE

# 处理单个URL的异步函数
async def scrape_single_url(context, url):
    page = await context.new_page()
    response = await page.goto(url)
    content = await response.text()
    extracted_title = re.search('<title>(.*)</title>', content).group(1)
    await page.close()
    return extracted_title

# 处理一个URL块的异步任务
@task(cache_policy=NONE)
async def process_url_block(urls):
    # 每个块创建独立的Playwright实例
    async with async_playwright() as p:
        browser = await p.chromium.launch()
        context = await browser.new_context()
        
        # 块内串行处理URL
        results = []
        for url in urls:
            title = await scrape_single_url(context, url)
            results.append(title)
        
        await context.close()
        await browser.close()
        return results

# 主Flow:并行处理多个块
@flow(name="Main Scraper Flow")
async def main_flow():
    # 模拟800个URL,分成4个200的块
    all_urls = ["https://www.ard.de/"] * 800
    url_blocks = [all_urls[i:i+200] for i in range(0, len(all_urls), 200)]
    
    # 并行提交所有块任务
    block_tasks = [process_url_block.submit(block) for block in url_blocks]
    
    # 等待所有任务完成并收集结果
    all_results = []
    for task in block_tasks:
        block_results = await task.result()
        all_results.extend(block_results)
    
    print(f"Total titles scraped: {len(all_results)}")
    return all_results

if __name__ == "__main__":
    import asyncio
    asyncio.run(main_flow())

关键改进说明

  • 异步任务与Flow:将整个Flow和任务改为异步,适配Playwright的异步API,避免事件循环冲突。
  • 块级Playwright实例:每个process_url_block任务独立创建和销毁Playwright实例,确保线程安全,每个块拥有专属的浏览器上下文。
  • 块内串行处理:在块任务内部用循环串行处理URL,符合需求的同时简化资源管理。
  • Prefect原生调度:利用Prefect的submit方法实现块的并行执行,由Prefect统一管理任务生命周期和资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:07:34