mirror of
https://github.com/wassname/ray.git
synced 2026-08-18 12:20:14 +08:00
[autoscaler] Create worker_file_mounts config (#9762)
This commit is contained in:
@@ -103,6 +103,7 @@ _hash_cache = {}
|
||||
|
||||
|
||||
def hash_runtime_conf(file_mounts,
|
||||
cluster_synced_files,
|
||||
extra_objs,
|
||||
generate_file_mounts_contents_hash=False):
|
||||
"""Returns two hashes, a runtime hash and file_mounts_content hash.
|
||||
@@ -111,20 +112,22 @@ def hash_runtime_conf(file_mounts,
|
||||
contents have changed. It is used at launch time (ray up) to determine if
|
||||
a restart is needed.
|
||||
|
||||
The file_mounts_content hash is used to determine if the file_mounts
|
||||
contents have changed. It is used at monitor time to determine if
|
||||
additional file syncing is needed.
|
||||
The file_mounts_content hash is used to determine if the file_mounts or
|
||||
cluster_synced_files contents have changed. It is used at monitor time to
|
||||
determine if additional file syncing is needed.
|
||||
"""
|
||||
runtime_hasher = hashlib.sha1()
|
||||
contents_hasher = hashlib.sha1()
|
||||
|
||||
def add_content_hashes(path):
|
||||
def add_content_hashes(path, allow_non_existing_paths: bool = False):
|
||||
def add_hash_of_file(fpath):
|
||||
with open(fpath, "rb") as f:
|
||||
for chunk in iter(lambda: f.read(2**20), b""):
|
||||
contents_hasher.update(chunk)
|
||||
|
||||
path = os.path.expanduser(path)
|
||||
if allow_non_existing_paths and not os.path.exists(path):
|
||||
return
|
||||
if os.path.isdir(path):
|
||||
dirs = []
|
||||
for dirpath, _, filenames in os.walk(path):
|
||||
@@ -146,15 +149,28 @@ def hash_runtime_conf(file_mounts,
|
||||
if conf_str not in _hash_cache or generate_file_mounts_contents_hash:
|
||||
for local_path in sorted(file_mounts.values()):
|
||||
add_content_hashes(local_path)
|
||||
contents_hash = contents_hasher.hexdigest()
|
||||
head_node_contents_hash = contents_hasher.hexdigest()
|
||||
|
||||
# Generate a new runtime_hash if its not cached
|
||||
# The runtime hash does not depend on the cluster_synced_files hash
|
||||
# because we do not want to restart nodes only if cluster_synced_files
|
||||
# contents have changed.
|
||||
if conf_str not in _hash_cache:
|
||||
runtime_hasher.update(conf_str)
|
||||
runtime_hasher.update(contents_hash.encode("utf-8"))
|
||||
runtime_hasher.update(head_node_contents_hash.encode("utf-8"))
|
||||
_hash_cache[conf_str] = runtime_hasher.hexdigest()
|
||||
|
||||
else:
|
||||
contents_hash = None
|
||||
# Add cluster_synced_files to the file_mounts_content hash
|
||||
if cluster_synced_files is not None:
|
||||
for local_path in sorted(cluster_synced_files):
|
||||
# For cluster_synced_files, we let the path be non-existant
|
||||
# because its possible that the source directory gets set up
|
||||
# anytime over the life of the head node.
|
||||
add_content_hashes(local_path, allow_non_existing_paths=True)
|
||||
|
||||
return (_hash_cache[conf_str], contents_hash)
|
||||
file_mounts_contents_hash = contents_hasher.hexdigest()
|
||||
|
||||
else:
|
||||
file_mounts_contents_hash = None
|
||||
|
||||
return (_hash_cache[conf_str], file_mounts_contents_hash)
|
||||
|
||||
Reference in New Issue
Block a user