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

为Dask Worker配置事件循环以支持Actor与aiohttp资源

问题背景

我尝试在Dask Worker中通过Actor模式实例化一个遗留数据提取器,代码如下:

from dask.distributed import Client
client = Client()  
connector = Sharepoint(CONF.sources["sharepoint"])  
items = connector.enumerate_items()

# extraction
remote_extractor = client.submit(
    SharepointExtractor, CONF.sources["sharepoint"], connector, actor=True
)  # 在Worker上创建提取器
extractor = remote_extractor.result()  # 获取该对象的指针

futures = client.map(
    extractor.job,
    [i for i in items],
    retries=5,
    pure=False,
)
_ = await client.gather(futures)

SharepointExtractor的初始化方法会从连接器获取HTTP会话:

class SharepointExtractor:
    def __init__(
        self, conf: ConfigTree, connector: Sharepoint, *args, **kwargs
    ) -> None:
        self.conf = conf
        self.session = connector.session_factory()

.session_factory()本质上返回一个带有OAuth令牌的aiohttp.client.ClientSession(这也是选择Actor模式的原因)。

报错问题

ClientSession的构造函数会调用asyncio.get_event_loop(),但Worker中似乎没有可用的事件循环,报错如下:

...
 File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/eteel/connectors/rest.py", line 96, in session_factory
    connector=TCPConnector(limit=30),
  File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/aiohttp/connector.py", line 767, in __init__
    super().__init__(
  File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/aiohttp/connector.py", line 234, in __init__
    loop = get_running_loop(loop)
  File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/aiohttp/helpers.py", line 287, in get_running_loop
    loop = asyncio.get_event_loop()
  File "/usr/lib/python3.10/asyncio/events.py", line 656, in get_event_loop
    raise RuntimeError('There is no current event loop in thread %r.'
RuntimeError: There is no current event loop in thread 'Dask-Default-Threads-484036-0'.

我使用的是本地开发环境,集群为LocalCluster。

尝试过的解决方案
  • 切换异步模式
    原本以为改成异步模式会自动给Worker注入事件循环,代码如下:
client = await Client(asynchronous=True)  
connector = Sharepoint(CONF.sources["sharepoint"])
items = connector.enumerate_items()

# extraction
remote_extractor = await client.submit(
    SharepointExtractor, CONF.sources["sharepoint"], connector, actor=True
)  # 在Worker上创建提取器
extractor = await remote_extractor  # 获取该对象的指针

但仍报相同错误。

  • 显式设置事件循环
    尝试手动创建事件循环:
loop = asyncio.new_event_loop()
client = await Client(
    asynchronous=True, loop=loop
)

这次报错变更为:

....
  File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/distributed/client.py", line 923, in __init__
    self._loop_runner = LoopRunner(loop=loop, asynchronous=asynchronous)
  File "/home/zar3bski/.cache/pypoetry/virtualenvs/poc-dask-iG-N0GH5-py3.10/lib/python3.10/site-packages/distributed/utils.py", line 451, in __init__
    if not loop.asyncio_loop.is_running():
AttributeError: '_UnixSelectorEventLoop' object has no attribute 'asyncio_loop'

不清楚该构造函数要求的loop参数格式。

  • 单例模式获取提取器
    按照@mdurant的方法,从可导入模块中通过单例模式获取提取器:
def get_extractor(CONF):
    if extractor[0] is None:
        connector = Sharepoint(CONF.sources["sharepoint"])
        extractor[0] = SharepointBis(CONF.sources["sharepoint"], connector)
    return extractor[0]


def workload(CONF, item):
    extractor = get_extractor(CONF)
    return extractor.job(item)

def main(): 
    client = Client()
    connector = Sharepoint(CONF.sources["sharepoint"])
    items = connector.enumerate_items()
    futures = client.map(
        workload,
        [CONF for _ in range(len(items))],
        [i for i in items],
        retries=5,
        pure=False,
    )
    _ = client.gather(futures)

依旧报错:

2022-12-01 10:05:54,923 - distributed.worker - WARNING - Compute Failed
Key:       workload-ffcf0f1a-8aee-41d1-9ad2-f7eea91fa107-41
Function:  workload
args:      (<eteel.conf.ConfGenerator object at 0x7fae8040d4e0>, 'firex1.sharepoint.com,930e9ef8-6bdf-4484-9883-6aa9965c548f,aed0d0bd-a659-4dbf-bbaa-a56f4efa3b0c')
kwargs:    {}
Exception: 'RuntimeError("There is no current event loop in thread \'Dask-Default-Threads-166860-1\'.")'

使用Client(asynchronous=True)也会出现相同错误。

疑问

如何让Dask Worker的线程拥有可用的事件循环?有没有涉及aiohttp(或其他异步库)资源的Dask Actor示例?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 05:55:34