Skip to content

Commit

Permalink
Enhanced thread-safety in zict.File (#7691)
Browse files Browse the repository at this point in the history
  • Loading branch information
crusaderky committed Mar 22, 2023
1 parent bfd2351 commit d338d65
Show file tree
Hide file tree
Showing 2 changed files with 9 additions and 8 deletions.
11 changes: 5 additions & 6 deletions distributed/tests/test_spill.py
Expand Up @@ -18,16 +18,15 @@
from distributed.utils import RateLimiterFilter
from distributed.utils_test import captured_logger

requires_zict_220 = pytest.mark.skipif(
not has_zict_220,
reason="requires zict version >= 2.2.0",
)


def psize(tmp_path: Path, **objs: object) -> tuple[int, int]:
# zict <= 2.2.0: tmp_path/key
# zict >= 2.3.0: tmp_path/key#0
fnames = tmp_path.glob("*")
key_to_fname = {fname.name.split("#")[0]: fname for fname in fnames}
return (
sum(sizeof(o) for o in objs.values()),
sum(os.stat(tmp_path / k).st_size for k in objs),
sum(os.stat(key_to_fname[k]).st_size for k in objs),
)


Expand Down
6 changes: 4 additions & 2 deletions distributed/tests/test_worker_memory.py
Expand Up @@ -911,9 +911,11 @@ def do_ignore_sigterm():
await wait(fut)
await c.run(lambda dask_worker: dask_worker.data.evict())
glob_out = await c.run(
lambda dask_worker: glob.glob(dask_worker.local_directory + "/**/myspill")
# zict <= 2.2.0: myspill
# zict >= 2.3.0: myspill#0
lambda dask_worker: glob.glob(dask_worker.local_directory + "/**/myspill*")
)
spill_fname = next(iter(glob_out.values()))[0]
spill_fname = glob_out[a.worker_address][0]
assert os.path.exists(spill_fname)

await leak_until_restart(c, s)
Expand Down

0 comments on commit d338d65

Please sign in to comment.