dask / dask/distributed

`get_worker` can return the wrong worker with an async local cluster

Open
#4,959 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
Python
Stars
1.7k
Forks
778
Avg merge
2h 50m
Merged PRs (30d)
3

Description

Split out from https://github.com/dask/distributed/pull/4937#issuecomment-866234111.

**What happened**:

When using a local cluster in async mode, `get_worker` always returns the same Worker instance, no matter which worker it's being called within.

**What you expected to happen**:

`get_worker`

**Minimal Complete Verifiable Example**:

```python
@gen_cluster(client=True)
async def test_get_worker_async_local(c, s, *workers):
assert len(workers) > 1
def whoami():
return get_worker().address

results = await c.run(whoami)
assert len(set(results.values())) == len(workers), results
E AssertionError: {'tcp://127.0.0.1:59774': 'tcp://127.0.0.1:59776', 'tcp://127.0.0.1:59776': 'tcp://127.0.0.1:59776'}
E assert 1 == 2
E +1
E -2
```

**Anything else we need to know?**:

This probably almost never affects users directly. But since most tests use an async local cluster with `@gen_cluster`, I'm concerned what edge-case behavior we might testing incorrectly.

Also, note that the same issue affects `get_client`. This feels a tiny bit less bad (at least it's always the _right_ client, unlike`get_worker`), but still can have some strange effects. In multiple places, worker code updates the default Client instance assuming it's in a separate process. With multiple workers trampling the default client, I wonder if this affects tests around advanced `secede`/client-within-task workloads.

I feel like the proper solution here would be to set a [contextvar](https://docs.python.org/3/library/contextvars.html) for the current worker that's updated as we context-switch in and out of that worker. Identifying the points where those switches have to happen seems tricky though.

I also think it would be reasonable for `get_worker` to error if `len(Worker._instances) > 1`.

**Environment**:

- Dask version: ac35e0f8104baaccfb44f98b67a48ee359069316
- Python version: 3.9.5
- Operating System: macOS
- Install method (conda, pip, source): source

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.