[Tune] Parametrize Cloud Syncing Frequency (#8771)

This commit is contained in:
Amog Kamsetty
2020-06-04 18:55:50 -07:00
committed by GitHub
parent d452932740
commit 9410e5884d
5 changed files with 81 additions and 13 deletions
+41 -10
View File
@@ -5,6 +5,7 @@ import time
from shlex import quote
from ray import ray_constants
from ray import services
from ray.tune.cluster_info import get_ssh_key, get_ssh_user
from ray.tune.sync_client import (CommandBasedClient, get_sync_client,
@@ -12,7 +13,13 @@ from ray.tune.sync_client import (CommandBasedClient, get_sync_client,
logger = logging.getLogger(__name__)
SYNC_PERIOD = 300
# Syncing period for syncing local checkpoints to cloud.
# In env variable is not set, sync happens every 300 seconds.
CLOUD_SYNC_PERIOD = ray_constants.env_integer(
key="TUNE_CLOUD_SYNC_S", default=300)
# Syncing period for syncing worker logs to driver.
NODE_SYNC_PERIOD = 300
_log_sync_warned = False
_syncers = {}
@@ -70,12 +77,23 @@ class Syncer:
self.last_sync_down_time = float("-inf")
self.sync_client = sync_client
def sync_up_if_needed(self):
if time.time() - self.last_sync_up_time > SYNC_PERIOD:
def sync_up_if_needed(self, sync_period):
"""Syncs up if time since last sync up is greather than sync_period.
Arguments:
sync_period (int): Time period between subsequent syncs.
"""
if time.time() - self.last_sync_up_time > sync_period:
self.sync_up()
def sync_down_if_needed(self):
if time.time() - self.last_sync_down_time > SYNC_PERIOD:
def sync_down_if_needed(self, sync_period):
"""Syncs down if time since last sync down is greather than sync_period.
Arguments:
sync_period (int): Time period between subsequent syncs.
"""
if time.time() - self.last_sync_down_time > sync_period:
self.sync_down()
def sync_up(self):
@@ -131,6 +149,19 @@ class Syncer:
return self._remote_dir
class CloudSyncer(Syncer):
"""Syncer for syncing files to/from the cloud."""
def __init__(self, local_dir, remote_dir, sync_client):
super(CloudSyncer, self).__init__(local_dir, remote_dir, sync_client)
def sync_up_if_needed(self):
return super(CloudSyncer, self).sync_up_if_needed(CLOUD_SYNC_PERIOD)
def sync_down_if_needed(self):
return super(CloudSyncer, self).sync_down_if_needed(CLOUD_SYNC_PERIOD)
class NodeSyncer(Syncer):
"""Syncer for syncing files to/from a remote dir to a local dir."""
@@ -158,12 +189,12 @@ class NodeSyncer(Syncer):
def sync_up_if_needed(self):
if not self.has_remote_target():
return True
return super(NodeSyncer, self).sync_up_if_needed()
return super(NodeSyncer, self).sync_up_if_needed(NODE_SYNC_PERIOD)
def sync_down_if_needed(self):
if not self.has_remote_target():
return True
return super(NodeSyncer, self).sync_down_if_needed()
return super(NodeSyncer, self).sync_down_if_needed(NODE_SYNC_PERIOD)
def sync_up_to_new_location(self, worker_ip):
if worker_ip != self.worker_ip:
@@ -226,16 +257,16 @@ def get_cloud_syncer(local_dir, remote_dir=None, sync_function=None):
return _syncers[key]
if not remote_dir:
_syncers[key] = Syncer(local_dir, remote_dir, NOOP)
_syncers[key] = CloudSyncer(local_dir, remote_dir, NOOP)
return _syncers[key]
client = get_sync_client(sync_function)
if client:
_syncers[key] = Syncer(local_dir, remote_dir, client)
_syncers[key] = CloudSyncer(local_dir, remote_dir, client)
return _syncers[key]
sync_client = get_cloud_sync_client(remote_dir)
_syncers[key] = Syncer(local_dir, remote_dir, sync_client)
_syncers[key] = CloudSyncer(local_dir, remote_dir, sync_client)
return _syncers[key]