From 9715658668f53cd25f373ec04cd4a0f955fa94a6 Mon Sep 17 00:00:00 2001 From: James Bourbeau Date: Fri, 10 Jan 2020 17:16:49 -0600 Subject: [PATCH 1/2] Close comm connection task before retrying --- distributed/comm/core.py | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) diff --git a/distributed/comm/core.py b/distributed/comm/core.py index 11f74a1aba8..bf8ee799360 100644 --- a/distributed/comm/core.py +++ b/distributed/comm/core.py @@ -1,14 +1,12 @@ from abc import ABC, abstractmethod, abstractproperty import asyncio -from datetime import timedelta import logging import weakref import dask -from tornado import gen from ..metrics import time -from ..utils import parse_timedelta, ignoring +from ..utils import parse_timedelta from . import registry from .addressing import parse_address @@ -211,13 +209,13 @@ def _raise(error): future = connector.connect( loc, deserialize=deserialize, **(connection_args or {}) ) - with ignoring(gen.TimeoutError): - comm = await gen.with_timeout( - timedelta(seconds=min(deadline - time(), 1)), - future, - quiet_exceptions=EnvironmentError, + try: + comm = await asyncio.wait_for( + future, timeout=min(deadline - time(), 1) ) break + except asyncio.TimeoutError: + pass if not comm: _raise(error) except FatalCommClosedError: From 4e76ff98a6a0f41bb86d514349cae9d89ba35d98 Mon Sep 17 00:00:00 2001 From: James Bourbeau Date: Fri, 10 Jan 2020 17:47:12 -0600 Subject: [PATCH 2/2] Switch back to ignoring context manager --- distributed/comm/core.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/distributed/comm/core.py b/distributed/comm/core.py index bf8ee799360..42c95e3579e 100644 --- a/distributed/comm/core.py +++ b/distributed/comm/core.py @@ -6,7 +6,7 @@ import dask from ..metrics import time -from ..utils import parse_timedelta +from ..utils import parse_timedelta, ignoring from . import registry from .addressing import parse_address @@ -209,13 +209,11 @@ def _raise(error): future = connector.connect( loc, deserialize=deserialize, **(connection_args or {}) ) - try: + with ignoring(asyncio.TimeoutError): comm = await asyncio.wait_for( future, timeout=min(deadline - time(), 1) ) break - except asyncio.TimeoutError: - pass if not comm: _raise(error) except FatalCommClosedError: