使用Dask时scheduler.address报错(ValueError:无法获取非运行服务器地址)
问题:Dask调度器启动后仍报错"cannot get address of non-running Server"
以下是调度器端代码:
# some imports import ... import dask.distributed import socket # some functions def some_function(): # etc. def handle_files(file, etc): # some code return some_output def find_free_port(): s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) s.bind(('', 0)) port = s.getsockname()[1] s.close() return port def main(): # some code list_of_files = ['file_path1', 'file_path2',...] scheduler = dask.distributed.Scheduler(protocol='tcp', host='xxx.xxx.xxx.xxx', port=find_free_port()) scheduler.start() client = dask.distributed.Client(scheduler.address) futures = [client.submit(handle_files, file, etc) for file in list_of_files] results = client.gather(futures) # rest of the code if __name__ == '__main__': main()
运行后出现错误:
File "C:\Users..\file.py",
line 401, in
main()File "C:\Users..\file.py",
line 382, in main
client = dask.distributed.Client(scheduler.address)File "C:\Anaconda\lib\site-packages\distributed\core.py", line 571,
in address
raise ValueError("cannot get address of non-running Server")ValueError: cannot get address of non-running Server
原因分析
调用scheduler.start()后,调度器是异步启动的,该方法不会阻塞等待调度器完全就绪。当你立刻访问scheduler.address时,调度器还未完成启动流程,所以触发"non-running Server"的错误。
解决方法
方法1:推荐使用Client自动创建调度器(简化流程)
不需要手动实例化和启动Scheduler,直接通过Client创建并连接调度器,这是Dask官方推荐的使用方式:
def main(): # some code list_of_files = ['file_path1', 'file_path2',...] # 直接创建Client,自动启动调度器并完成连接 client = dask.distributed.Client(protocol='tcp', host='xxx.xxx.xxx.xxx', port=find_free_port()) futures = [client.submit(handle_files, file, etc) for file in list_of_files] results = client.gather(futures) # rest of the code
方法2:手动启动调度器时等待就绪
如果必须手动管理Scheduler,调用start()后需要等待调度器完全启动,可使用scheduler.wait_for_start()方法:
def main(): # some code list_of_files = ['file_path1', 'file_path2',...] scheduler = dask.distributed.Scheduler(protocol='tcp', host='xxx.xxx.xxx.xxx', port=find_free_port()) scheduler.start() # 等待调度器启动完成,确保address可用 scheduler.wait_for_start() client = dask.distributed.Client(scheduler.address) futures = [client.submit(handle_files, file, etc) for file in list_of_files] results = client.gather(futures) # rest of the code
内容的提问来源于stack exchange,提问作者AxieKendy
相关产品推荐
相关产品推荐

