为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
相关产品推荐
相关产品推荐

