Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 5 additions & 9 deletions packages/dask-task-models-library/requirements/_base.txt
Original file line number Diff line number Diff line change
Expand Up @@ -19,17 +19,17 @@ cloudpickle==3.1.2
# via
# dask
# distributed
dask==2026.3.0
dask==2026.7.1
# via
# -r requirements/_base.in
# distributed
distributed==2026.3.0
distributed==2026.7.1
# via dask
dnspython==2.8.0
# via email-validator
email-validator==2.3.0
# via pydantic
fsspec==2026.2.0
fsspec==2026.7.0
# via dask
idna==3.11
# via email-validator
Expand All @@ -53,7 +53,7 @@ markupsafe==3.0.3
# via jinja2
mdurl==0.1.2
# via markdown-it-py
msgpack==1.1.2
msgpack==1.2.1
# via distributed
orjson==3.11.9
# via
Expand Down Expand Up @@ -122,7 +122,7 @@ toolz==1.1.0
# dask
# distributed
# partd
tornado==6.5.5
tornado==6.5.8
# via distributed
typer==0.24.1
# via -r requirements/../../../packages/settings-library/requirements/_base.in
Expand All @@ -138,9 +138,5 @@ typing-inspection==0.4.2
# pydantic-settings
tzdata==2025.3
# via arrow
urllib3==2.7.0
# via
# -c requirements/../../../requirements/constraints.txt
# distributed
zict==3.0.0
# via distributed
14 changes: 6 additions & 8 deletions services/autoscaling/requirements/_base.txt
Original file line number Diff line number Diff line change
Expand Up @@ -90,12 +90,12 @@ cloudpickle==3.1.2
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
# distributed
dask==2025.11.0
dask==2026.7.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# -r requirements/_base.in
# distributed
distributed==2025.11.0
distributed==2026.7.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
Expand Down Expand Up @@ -128,7 +128,7 @@ frozenlist==1.8.0
# via
# aiohttp
# aiosignal
fsspec==2025.10.0
fsspec==2026.7.0
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
Expand Down Expand Up @@ -195,7 +195,7 @@ markupsafe==3.0.3
# jinja2
mdurl==0.1.2
# via markdown-it-py
msgpack==1.1.2
msgpack==1.2.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# distributed
Expand Down Expand Up @@ -333,7 +333,7 @@ protobuf==6.33.5
# -c requirements/../../../requirements/constraints.txt
# googleapis-common-protos
# opentelemetry-proto
psutil==7.1.3
psutil==7.2.2
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# -r requirements/../../../packages/aws-library/requirements/../../../packages/service-library/requirements/_base.in
Expand Down Expand Up @@ -449,7 +449,7 @@ toolz==1.1.0
# dask
# distributed
# partd
tornado==6.5.2
tornado==6.5.8
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# distributed
Expand Down Expand Up @@ -500,9 +500,7 @@ tzdata==2025.2
urllib3==2.7.0
# via
# -c requirements/../../../requirements/constraints.txt
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# botocore
# distributed
# requests
# sentry-sdk
uvicorn==0.38.0
Expand Down
2 changes: 1 addition & 1 deletion services/autoscaling/requirements/_test.txt
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ pluggy==1.6.0
# pytest-cov
pprintpp==0.4.0
# via pytest-icdiff
psutil==7.1.3
psutil==7.2.2
# via
# -c requirements/_base.txt
# -r requirements/_test.in
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,9 +61,12 @@ def _scheduler_identity_key_builder(func: Callable[..., Any], client: distribute
async def _get_scheduler_identity(client: distributed.Client) -> SchedulerInfo:
"""Returns all workers from the scheduler with a 2-second TTL cache.

client.scheduler_info() is a local cache capped at 5 workers for async clients
since https://github.com/dask/distributed/pull/9045.
client.scheduler.identity(n_workers=-1) is a live RPC but we cache it briefly
We use a live RPC instead of client.scheduler_info(), a local cache that
(a) capped results at 5 workers for async clients before
https://github.com/dask/distributed/pull/9308 (fixed in distributed 2026.7.0), and
(b) can still be empty/incomplete right after a client connects, since it is only
populated once the scheduler pushes its periodic info broadcast.
client.scheduler.identity(n_workers=-1) has neither limitation; we cache it briefly
to avoid redundant round-trips within a single autoscaling tick.
"""
assert client.scheduler # nosec
Expand Down Expand Up @@ -110,8 +113,7 @@ async def _dask_worker_from_ec2_instance(
) -> tuple[DaskWorkerUrl, DaskWorkerDetails]:
"""
Uses client.scheduler.identity() RPC to get all workers live from the scheduler,
bypassing client.scheduler_info() which is a local cache capped at 5 workers
for async clients regardless of the n_workers argument.
avoiding client.scheduler_info()'s local-cache limitations (see _get_scheduler_identity).

Raises:
Ec2InvalidDnsNameError
Expand Down
21 changes: 8 additions & 13 deletions services/autoscaling/tests/unit/test_modules_dask.py
Original file line number Diff line number Diff line change
Expand Up @@ -475,12 +475,15 @@ def dask_workers_config_with_more_than_5_workers() -> dict[str, Any]:
}


async def test_get_scheduler_identity_returns_all_workers_beyond_default_cap(
async def test_get_scheduler_identity_returns_all_workers_reliably(
dask_workers_config_with_more_than_5_workers: dict[str, Any],
dask_scheduler_config: dict[str, Any],
):
"""Regression test: scheduler_info() caps at 5 workers for async clients;
_get_scheduler_identity() must return all workers regardless."""
"""Regression test: client.scheduler_info() is a local cache that (a) used to cap results at
5 workers for async clients (fixed in distributed>=2026.7.0, see distributed#9308) and (b) can
still be empty/incomplete immediately after connecting, since it is only synced once the
scheduler pushes its periodic broadcast. _get_scheduler_identity() uses a live RPC instead and
must always return all connected workers, regardless of client-side cache state."""
_NUM_WORKERS: Final[int] = 6
assert len(dask_workers_config_with_more_than_5_workers) == _NUM_WORKERS

Expand All @@ -492,16 +495,8 @@ async def test_get_scheduler_identity_returns_all_workers_beyond_default_cap(
) as cluster,
distributed.Client(cluster.scheduler_address, asynchronous=True) as client,
):
# Demonstrate the bug: scheduler_info() with default n_workers=5 misses
# the 6th worker even when called with n_workers=-1 on an async client
info_default = client.scheduler_info()
assert len(info_default["workers"]) == 5, "scheduler_info() should return at most 5 workers (the default cap)"
info_minus1 = client.scheduler_info(n_workers=-1)
assert len(info_minus1["workers"]) == 5, (
"scheduler_info(n_workers=-1) is still broken for async clients - it ignores the argument"
)

# Our fix: _get_scheduler_identity() bypasses the cache and returns all workers
# _get_scheduler_identity() must return all workers right away, unlike scheduler_info()
# which relies on a local cache that may not be synced yet right after connecting.
identity = await _get_scheduler_identity(client)
assert len(identity["workers"]) == _NUM_WORKERS, (
f"_get_scheduler_identity() must return all {_NUM_WORKERS} workers"
Expand Down
14 changes: 6 additions & 8 deletions services/clusters-keeper/requirements/_base.txt
Original file line number Diff line number Diff line change
Expand Up @@ -87,12 +87,12 @@ cloudpickle==3.1.2
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
# distributed
dask==2025.11.0
dask==2026.7.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# -r requirements/_base.in
# distributed
distributed==2025.11.0
distributed==2026.7.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
Expand Down Expand Up @@ -125,7 +125,7 @@ frozenlist==1.8.0
# via
# aiohttp
# aiosignal
fsspec==2025.10.0
fsspec==2026.7.0
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# dask
Expand Down Expand Up @@ -192,7 +192,7 @@ markupsafe==3.0.3
# jinja2
mdurl==0.1.2
# via markdown-it-py
msgpack==1.1.2
msgpack==1.2.1
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# distributed
Expand Down Expand Up @@ -330,7 +330,7 @@ protobuf==6.33.5
# -c requirements/../../../requirements/constraints.txt
# googleapis-common-protos
# opentelemetry-proto
psutil==7.1.3
psutil==7.2.2
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# -r requirements/../../../packages/aws-library/requirements/../../../packages/service-library/requirements/_base.in
Expand Down Expand Up @@ -446,7 +446,7 @@ toolz==1.1.0
# dask
# distributed
# partd
tornado==6.5.2
tornado==6.5.8
# via
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# distributed
Expand Down Expand Up @@ -497,9 +497,7 @@ tzdata==2025.2
urllib3==2.7.0
# via
# -c requirements/../../../requirements/constraints.txt
# -c requirements/../../../services/dask-sidecar/requirements/_dask-distributed.txt
# botocore
# distributed
# requests
# sentry-sdk
uvicorn==0.38.0
Expand Down
2 changes: 1 addition & 1 deletion services/clusters-keeper/requirements/_test.txt
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,7 @@ propcache==0.4.1
# -c requirements/_base.txt
# aiohttp
# yarl
psutil==7.1.3
psutil==7.2.2
# via
# -c requirements/_base.txt
# -r requirements/_test.in
Expand Down
37 changes: 12 additions & 25 deletions services/dask-sidecar/requirements/_base.txt
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,9 @@ attrs==25.4.0
# aiohttp
# jsonschema
# referencing
blosc==1.11.3
blosc==1.11.4
# via -r requirements/_base.in
bokeh==3.8.1
bokeh==3.10.0
# via dask
boto3==1.40.61
# via aiobotocore
Expand Down Expand Up @@ -90,19 +90,17 @@ cloudpickle==3.1.2
# via
# dask
# distributed
contourpy==1.3.3
# via bokeh
cryptography==50.0.0
# via
# -c requirements/../../../requirements/constraints.txt
# -r requirements/_base.in
dask==2025.11.0
dask==2026.7.1
# via
# -c requirements/constraints.txt
# -r requirements/../../../packages/dask-task-models-library/requirements/_base.in
# -r requirements/_base.in
# distributed
distributed==2025.11.0
distributed==2026.7.1
# via dask
dnspython==2.8.0
# via email-validator
Expand All @@ -122,7 +120,7 @@ frozenlist==1.8.0
# via
# aiohttp
# aiosignal
fsspec==2025.10.0
fsspec==2026.7.0
# via
# -r requirements/_base.in
# dask
Expand Down Expand Up @@ -183,7 +181,7 @@ markupsafe==3.0.3
# via jinja2
mdurl==0.1.2
# via markdown-it-py
msgpack==1.1.2
msgpack==1.2.1
# via distributed
multidict==6.7.0
# via
Expand All @@ -192,11 +190,8 @@ multidict==6.7.0
# yarl
narwhals==2.11.0
# via bokeh
numpy==2.3.5
# via
# bokeh
# contourpy
# pandas
numpy==2.5.2
# via bokeh
opentelemetry-api==1.44.0
# via
# -r requirements/../../../packages/aws-library/requirements/../../../packages/service-library/requirements/_base.in
Expand Down Expand Up @@ -296,8 +291,6 @@ packaging==25.0
# opentelemetry-instrumentation-sqlalchemy
pamqp==3.3.0
# via aiormq
pandas==2.3.3
# via bokeh
partd==1.4.2
# via dask
pillow==12.2.0
Expand All @@ -315,7 +308,7 @@ protobuf==6.33.5
# -c requirements/../../../requirements/constraints.txt
# googleapis-common-protos
# opentelemetry-proto
psutil==7.1.3
psutil==7.2.2
# via
# -r requirements/../../../packages/aws-library/requirements/../../../packages/service-library/requirements/_base.in
# distributed
Expand Down Expand Up @@ -355,11 +348,8 @@ python-dateutil==2.9.0.post0
# aiobotocore
# arrow
# botocore
# pandas
python-dotenv==1.2.1
# via pydantic-settings
pytz==2025.2
# via pandas
pyyaml==6.0.3
# via
# -c requirements/../../../requirements/constraints.txt
Expand Down Expand Up @@ -387,7 +377,7 @@ rpds-py==0.29.0
# via
# jsonschema
# referencing
s3fs==2025.10.0
s3fs==2026.7.0
# via fsspec
s3transfer==0.14.0
# via boto3
Expand Down Expand Up @@ -417,7 +407,7 @@ toolz==1.1.0
# dask
# distributed
# partd
tornado==6.5.2
tornado==6.5.8
# via
# bokeh
# distributed
Expand Down Expand Up @@ -460,14 +450,11 @@ typing-inspection==0.4.2
# pydantic
# pydantic-settings
tzdata==2025.2
# via
# arrow
# pandas
# via arrow
urllib3==2.7.0
# via
# -c requirements/../../../requirements/constraints.txt
# botocore
# distributed
# requests
wrapt==1.17.3
# via
Expand Down
Loading
Loading