dask / dask/distributed

Outdated and/or unclear documentation about SchedulerPlugin

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

Description

**Describe the issue**:
[The current documentation regarding SchedulerPlugin with full task state access](https://distributed.dask.org/en/stable/plugins.html#accessing-full-task-state) gives the following code as example:
```python
from distributed.diagnostics.plugin import SchedulerPlugin

class MyPlugin(SchedulerPlugin):
def __init__(self, scheduler):
self.scheduler = scheduler

def transition(self, key, start, finish, *args, **kwargs):
# Get full TaskState
ts = self.scheduler.tasks[key]

@click.command()
def dask_setup(scheduler):
plugin = MyPlugin(scheduler)
scheduler.add_plugin(plugin)
```

The `dask_setup` function seems to appear out of no-where and `click` isn't imported. Moreover, this approach does not seem to correctly register the plugin. For example, the follow code runs successfully:
```python
import click
from distributed.diagnostics.plugin import SchedulerPlugin
from distributed import Client, LocalCluster

class MyPlugin(SchedulerPlugin):
def __init__(self, scheduler):
self.scheduler = scheduler

def transition(self, key, start, finish, *args, **kwargs):
ts = self.scheduler.tasks[key]
raise Exception

@click.command()
def dask_setup(scheduler):
plugin = MyPlugin(scheduler)
scheduler.add_plugin(plugin)

def job(i):
return i*2

if __name__ == "__main__":
cluster = LocalCluster()
client = Client(cluster)
ret = client.submit(job, 1).result()
print(ret)
```

The console log reads "2", while we would expect an Exception to be raised.

It is unclear how to register a scheduler plugin that can access the full task state. `client.scheduler` does not seem to be serializable, so the following code does not work and crashes with the error `TypeError: cannot pickle 'coroutine' object`:
```python
from distributed.diagnostics.plugin import SchedulerPlugin
from distributed import Client, LocalCluster

class MyPlugin(SchedulerPlugin):
def __init__(self, scheduler):
self.scheduler = scheduler

def transition(self, key, start, finish, *args, **kwargs):
ts = self.scheduler.tasks[key]
raise Exception

def job(i):
return i*2

if __name__ == "__main__":
cluster = LocalCluster()
client = Client(cluster)

plugin = MyPlugin(client.scheduler)
client.register_scheduler_plugin(plugin)

ret = client.submit(job, 1).result()
print(ret)
```

**Environment**:

- Dask version: 2024.4.2
- Python version: 3.10.12
- Operating System: Linux Mint 21.2
- Install method (conda, pip, source): pip

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.