NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
PyPI · #661 most downloaded on PyPI
Parallel PyData with Task Scheduling
Last release 1 months ago
24 Aug 2026
Ships fairly regularly
a new release about every 5 weeks
Nearly every release is documented
notes for 60 of the last 60 stable releases
2 versions withdrawn
withdrawn after publishing
12 years old
239 releases · first in 2015
docs: fix duplicated-word typos in docstrings @Sreekant13
See the Changelog for more information.
One column per quarter.
See the Changelog for more information.
CI: Suppress numpy 'generic' timedelta deprecation from pandas @crusaderky
__dask_exprs__ protocol @mrocklin (#12457)test_concat_categorical (again) @crusaderky (#12463)See the Changelog for more information.
Support composite expressions with the dask_exprs protocol and avoid graph materialization when checking expr-wrapper Dask collections ( dask#12457 , dask#12476 ) Matthew Rocklin
Update compatibility for NumPy 2.5 type stubs, pandas 3 nightlies, pytest 9.1, and free-threaded msgpack builds ( dask#12483 , dask#12469 , dask#12465 , dask#12467 , distributed#9312 , distributed#9304 , distributed#9303 , distributed#9302 ) Guido Imperiale
Return all workers from scheduler_info() by default ( distributed#9308 ) Matthew Rocklin
Keep Client.scatter from unpacking custom containers ( distributed#9298 ) Adrin Jalali
Fix Distributed worker startup and aborted TCP communication races ( distributed#9313 , distributed#9315 ) Guido Imperiale
Additional changes
Upgrade pixi ( dask#12460 ) Matthew Rocklin
Fix release publisher warnings ( dask#12456 , distributed#9299 ) Matthew Rocklin
XFAIL test_concat_categorical (again) ( dask#12463 ) Guido Imperiale
CI: Suppress warning about slow disk ( dask#12464 ) Guido Imperiale
Support for pytest 9.1 ( dask#12465 , distributed#9303 ) Guido Imperiale
Free-threading: use upstream msgpack ( dask#12467 , distributed#9302 ) Guido Imperiale
CI: Suppress numpy generic timedelta deprecation from pandas ( dask#12468 ) Guido Imperiale
Pin pandas 3 nightlies ( dask#12469 , distributed#9304 ) Guido Imperiale
Bump actions/checkout from 6 to 7 ( dask#12473 , distributed#9306 )
Fix randomly failing s3fs checkouts in nightly CI ( dask#12474 ) Guido Imperiale
Avoid graph materialization in is_dask_collection for expr wrappers ( dask#12476 ) Matthew Rocklin
Bump prefix-dev/setup-pixi from 0.9.6 to 0.10.0 ( dask#12478 , distributed#9310 )
Bump actions/cache from 5 to 6 ( dask#12479 , distributed#9309 )
Avoid NumPy 2.5 type stubs ( dask#12483 , distributed#9312 ) Guido Imperiale
Bump actions/download-artifact from 7 to 8 ( dask#12461 , distributed#9301 )
Support composite expressions with dask_exprs protocol ( dask#12457 ) Matthew Rocklin
Return all workers from scheduler_info() by default ( distributed#9308 ) Matthew Rocklin
Client.scatter should not unpack custom containers ( distributed#9298 ) Adrin Jalali
Suppress mindeps warning when running a single test in CI crusaderky
AI agent skill to debug flaky tests ( distributed#9319 ) Guido Imperiale
Fix Nanny race conditions when worker fails to start ( distributed#9313 ) Guido Imperiale
Fix server hanging on aborted TCP comms ( distributed#9315 ) Guido Imperiale
Fix conda build CI ( distributed#9314 ) Guido Imperiale
Update deprecated CUDA system requirements in Pixi ( distributed#9311 ) Guido Imperiale
Support composite expressions with the __dask_exprs__ protocol and avoid graph materialization when checking expr-wrapper Dask collections (12457, 12476) `Matthew Rocklin`_
Update compatibility for NumPy 2.5 type stubs, pandas 3 nightlies, pytest 9.1, and free-threaded msgpack builds (12483, 12469, 12465, 12467, 9312, 9304, 9303, 9302) `Guido Imperiale`_
Return all workers from scheduler_info() by default (9308) `Matthew Rocklin`_
Keep Client.scatter from unpacking custom containers (9298) `Adrin Jalali`_
Fix Distributed worker startup and aborted TCP communication races (9313, 9315) `Guido Imperiale`_
Pandas 3.1: The inplace keyword in DataFrame.drop is deprecated @crusaderky
test_concat_categorical is still flaky on Pandas 2.0 @crusaderky (#12454)assert_eq helper @kyo5uke (#12441)test_interrupt @crusaderky (#12437)add_prefix/add_suffix in pandas nightly; support explicit axis @crusaderky (#12414)import_optional_dependency @crusaderky (#12412)test_describe vs. pandas nightly @crusaderky (#12409)test_shared_tasks in free-threading @crusaderky (#12398)test_store_locks @crusaderky (#12393)test_len @crusaderky (#12396)test_map_block_series @crusaderky (#12392)test_from_delayed_future @crusaderky (#12391)test_ffill_bfill @crusaderky (#12378)__package__ → __spec__.parent @DimitriPapadopoulos (#12333)See the Changelog for more information.
Improve pandas 3.1 compatibility, including add_prefix / add_suffix and DataFrame.drop changes ( dask#12414 , dask#12447 ) Guido Imperiale
Fix groupby compatibility with pandas 2.2 and 2.3 and test against intermediate NumPy and pandas versions ( dask#12372 ) Guido Imperiale
Store empty arrays with Zarr 3.2.0 ( dask#12366 ) Guido Imperiale
Fix quantile and nanquantile behavior ( dask#12380 ) Guido Imperiale
Add support for Cholesky decomposition with complex dtypes ( dask#12416 ) Niclas Rieger
Fix dataframe correctness issues in DataFrame.merge and Series.map ( dask#12430 , dask#12432 ) nn , Mike Evdokimov
Overhaul the publish_dataset extension ( distributed#9217 ) Guido Imperiale
Overhaul task stream handling ( distributed#9230 , distributed#9282 , distributed#9293 ) Guido Imperiale , MohammadYusif
Fix a race condition in SubprocessCluster startup ( distributed#9292 ) Guido Imperiale
Clean up deprecated Distributed APIs across scheduler, worker, CLI, plugin, and progress handling ( distributed#9222 , distributed#9225 , distributed#9238 , distributed#9240 , distributed#9250 ) Guido Imperiale
Migrate CI and local workflows to Pixi ( dask#12389 , distributed#9276 ) Guido Imperiale
Add AI contribution guidance with AGENTS.md and CLAUDE.md ( dask#12415 , distributed#9279 ) Guido Imperiale
Add PyPI release workflows with GitHub Actions Trusted Publishing ( dask#12452 , distributed#9297 ) Matthew Rocklin
Additional changes
Pandas 3.1: The inplace keyword in DataFrame.drop is deprecated ( dask#12447 ) Guido Imperiale
test_concat_categorical is still flaky on Pandas 2.0 ( dask#12454 ) Guido Imperiale
Add tests for the array assert_eq helper ( dask#12441 ) きょうすけ
Revert security autoescape change in HTML reprs ( dask#12451 ) Guido Imperiale
Fix Dataframe.merge losing data when there are 128 or 129 partitions ( dask#12430 ) nn
Security: enable Jinja2 autoescape to prevent XSS in HTML reprs ( dask#12423 ) dfgvaetyj3456356-hash
Test vs. NumPy 1.26 and PyArrow 18 ( dask#12445 ) Guido Imperiale
Do not build the dask-core pixi package twice ( dask#12443 ) Guido Imperiale
Fix broken condition in test_cpu_affinity_taskset ( dask#12399 ) Elliott Sales de Andrade
Fix flaky test_interrupt ( dask#12437 ) Guido Imperiale
Free-threading: enable msgpack C extension ( dask#12439 ) Guido Imperiale
Fix Series.map(Series) producing wrong results for non-co-aligned inputs ( dask#12432 ) Mike Evdokimov
Repair Upstream CI ( dask#12436 ) Guido Imperiale
Strictly expect test failures ( dask#12435 ) Guido Imperiale
Update pixi lockfile ( dask#12433 ) Guido Imperiale
Add more thorough tests for Series.map and fix pandas nightly CI ( dask#12425 ) Guido Imperiale
Add AGENTS.md / CLAUDE.md and clarify guidelines for AI pull requests ( dask#12415 ) Guido Imperiale
Tweak local coverage HTML report ( dask#12422 ) Guido Imperiale
Make tests less dependent on urllib3 and requests ( dask#12419 ) Guido Imperiale
Fix trivial Sphinx warning crusaderky
Do not run scheduled tests on forks ( dask#12418 ) Guido Imperiale
Only upload conda packages from main branch crusaderky
Repair conda upload ( dask#12417 ) Guido Imperiale
Fix add_prefix / add_suffix in pandas nightly and support explicit axis ( dask#12414 ) Guido Imperiale
Remove legacy pandas dtype testing strings vs. decimals ( dask#12413 ) Guido Imperiale
Better type annotations for import_optional_dependency ( dask#12412 ) Guido Imperiale
Resurrect PySpark tests ( dask#12410 ) Guido Imperiale
Test on Linux ARM ( dask#12408 ) Guido Imperiale
Fix failures in test_describe vs. pandas nightly ( dask#12409 ) Guido Imperiale
Migrate CI to Pixi ( dask#12389 ) Guido Imperiale
Fix typos ( dask#12405 ) Guido Imperiale
XFAIL regression in click 8.4.0 ( dask#12407 ) Guido Imperiale
Fix flaky test_shared_tasks in free-threading ( dask#12398 ) Guido Imperiale
Enforce Python 3.10 compatibility in black ( dask#12397 ) Guido Imperiale
Update Optuna+Dask example to use optuna-integration ( dask#12385 ) Mike Evdokimov
Fix flaky test_store_locks ( dask#12393 ) Guido Imperiale
XFAIL flaky test_len ( dask#12396 ) Guido Imperiale
Run some tests serially in pytest-xdist ( dask#12395 ) Guido Imperiale
test_orc segfaults with PyArrow 16 ( dask#12394 ) Guido Imperiale
Unskip test_map_block_series ( dask#12392 ) Guido Imperiale
Speed up test_from_delayed_future ( dask#12391 ) Guido Imperiale
Fix duplicated words in dask.array routines and slicing docstrings ( dask#12390 ) lphuc2250gma
Simplify test_dot and suppress spurious test outputs ( dask#12388 ) Guido Imperiale
Test vs. intermediate versions of numpy and pandas ( dask#12372 ) Guido Imperiale
Clean up obsolete version check in test_http ( dask#12387 ) Guido Imperiale
Fix LaTeX formula for split_every in docstrings ( dask#12379 ) stephenworsley
Remove numexpr ( dask#12384 ) Guido Imperiale
Unit-less timedelta64 is deprecated in NumPy nightly ( dask#12383 ) Guido Imperiale
Require pyyaml >=5.4.1 ( dask#12382 ) Guido Imperiale
Re-raise ModuleNotFoundError instead of ImportError ( dask#12381 ) Guido Imperiale
Fix flaky test_ffill_bfill ( dask#12378 ) Guido Imperiale
Raise NotImplementedError for da.quantile with weights on NumPy < 2.0 ( dask#12370 ) Dr Alex Mitre
Fix typo in array-chunks.rst ( dask#12375 ) Henry
Minor docs tweaks ( dask#12377 ) Guido Imperiale
Clean up minimum pyarrow version checks ( dask#12376 ) Guido Imperiale
Use Python 3.14 in additional and upstream envs and dev docs ( dask#12373 ) Guido Imperiale
More test fixes for pandas-nightly ( dask#12369 ) Guido Imperiale
Switch back to conda ( dask#12368 ) Guido Imperiale
Disable CI for Windows 3.14t ( dask#12367 ) Guido Imperiale
Use f-strings ( dask#12362 ) Dimitri Papadopoulos Orfanos
Store empty arrays with zarr 3.2.0 ( dask#12366 ) Guido Imperiale
Fix some pandas nightly deprecation warnings ( dask#12341 ) Trevin Chow
3.14t CI: tweak notes ( dask#12363 ) Guido Imperiale
Use f-strings ( dask#12354 ) Dimitri Papadopoulos Orfanos
Docs: order array arguments in docstring ( dask#12342 ) Peter A. Jonsson
numba is now available on conda-forge for 3.14t ( dask#12358 ) Guido Imperiale
package to spec.parent ( dask#12333 ) Dimitri Papadopoulos Orfanos
Update pre-commit black hook ( dask#12344 ) Dimitri Papadopoulos Orfanos
Update pre-commit ruff hook ( dask#12345 ) Dimitri Papadopoulos Orfanos
Fix typos ( dask#12339 ) Dimitri Papadopoulos Orfanos
Add numba, sparse, and h5py to 3.14t CI ( dask#12338 ) Guido Imperiale
Improve Contribution Policy ( dask#12320 ) Jacob Tomlinson
Fix test failures caused by port 8787 already in use ( distributed#9296 ) Guido Imperiale
Migrate from black to ruff format ( distributed#9295 ) Guido Imperiale
Fix race condition in SubprocessCluster startup ( distributed#9292 ) Guido Imperiale
Refactor get_task_stream ( distributed#9293 ) Guido Imperiale
Do not build the distributed pixi package twice ( distributed#9289 ) Guido Imperiale
Fix flaky test_steal_twice and test_steal_when_more_tasks ( distributed#9288 ) Guido Imperiale
Clarify scheduler unpickling in protocol docs ( distributed#9286 ) Peter Chen J.
Reinstate no-queue tests ( distributed#9287 ) Guido Imperiale
Clarify unsatisfied resource restrictions in resource docs ( distributed#9269 ) Peter Chen J.
Fix import-time coverage ( distributed#9284 ) Guido Imperiale
Improve CLI error message for unknown options to dask worker ( distributed#9281 ) Zeus Almightee
Run most CUDA tests in pixi ( distributed#9285 ) Guido Imperiale
Migrate CI to pixi ( distributed#9276 ) Guido Imperiale
Add AGENTS.md / CLAUDE.md and guidelines for AI pull requests ( distributed#9279 ) Guido Imperiale
Fix tests vs. NumPy nightly builds ( distributed#9280 ) Guido Imperiale
uvloop tests ( distributed#9278 ) Guido Imperiale
Fix flaky test_call_stack_future ( distributed#9277 ) Guido Imperiale
Drop dependency on urllib3 ( distributed#9273 ) James Lamb
Remove avoid_ci mark ( distributed#9275 ) Guido Imperiale
Print host_info in CI without needing NumPy ( distributed#9274 ) Guido Imperiale
Fix typos ( distributed#9268 ) Guido Imperiale
Support pixi in dask/dask CI ( distributed#9264 ) Guido Imperiale
Truncate large coro repr in retry log output ( distributed#9197 ) Ernest Provo
Add explicit check for name not None ( distributed#9265 ) Maneesh Sutar
XFAIL test_scheduler_bokeh.py::test_simple in Windows ( distributed#9266 ) Guido Imperiale
Clean up deprecated loop properties ( distributed#9231 ) Guido Imperiale
Clean up deprecations in Scheduler, Worker, Nanny ( distributed#9238 ) Guido Imperiale
Clean up deprecated client context in different tasks/threads ( distributed#9233 ) Guido Imperiale
Clean up deprecated Prometheus metrics ( distributed#9249 ) Guido Imperiale
Clean up deprecations in distributed.deploy ( distributed#9244 ) Guido Imperiale
Clean up deprecated stream RPC handler argument ( distributed#9242 ) Guido Imperiale
Deprecations in CLI ( distributed#9240 ) Guido Imperiale
Clean up deprecations in security ( distributed#9239 ) Guido Imperiale
Clean up deprecated rpc synchronous context manager ( distributed#9235 ) Guido Imperiale
Clean up deprecated async listener.stop() ( distributed#9234 ) Guido Imperiale
Clean up deprecations in register_plugin ( distributed#9225 ) Guido Imperiale
Clean up deprecations in remove_worker ( distributed#9222 ) Guido Imperiale
Clean up minimum pyarrow version checks ( distributed#9260 ) Guido Imperiale
Partial review of task streams ( distributed#9230 ) Guido Imperiale
Update outdated work stealing docs regarding worker restrictions ( distributed#9214 ) Kevin Ziroldi
Fix memray CI; move back from mamba to conda ( distributed#9258 ) Guido Imperiale
Install keras with conda ( distributed#9259 ) Guido Imperiale
Clean up Cython ( distributed#9257 ) Guido Imperiale
Use @dataclass(slots=True) with constraints from Python <3.14 ( distributed#9256 ) Guido Imperiale
Type annotations for client.Future ( distributed#9255 ) Guido Imperiale
Clean up deprecated Lock client parameter ( distributed#9237 ) Guido Imperiale
Clean up deprecated nested_deserialize ( distributed#9243 ) Guido Imperiale
Deprecations in distributed.compatibility ( distributed#9232 ) Guido Imperiale
Clean up deprecations in utils_test ( distributed#9226 ) Guido Imperiale
Clean up deprecations in distributed.utils ( distributed#9236 ) Guido Imperiale
Clean up deprecations in ProgressBar ( distributed#9250 ) Guido Imperiale
Run more tests when requests is not installed ( distributed#9251 ) Guido Imperiale
Standardize and increase async_poll_for timeout ( distributed#9248 ) Guido Imperiale
Relax unreasonably short test timeouts ( distributed#9247 ) Guido Imperiale
Use f-strings ( distributed#9245 ) Dimitri Papadopoulos Orfanos
Update link to bokeh sources in comment ( distributed#9241 ) Guido Imperiale
Apply ruff/pyupgrade rule UP031 ( distributed#9204 ) Dimitri Papadopoulos Orfanos
Homogeneous environment names ( distributed#9209 ) Dimitri Papadopoulos Orfanos
Overhaul publish_dataset extension ( distributed#9217 ) Guido Imperiale
Fix flaky test_mixing_clients_same_scheduler ( distributed#9229 ) Guido Imperiale
Type hints: @contextmanager should return Generator[T] ( distributed#9228 ) Guido Imperiale
Fix flaky test_actor.py::test_failed_worker ( distributed#9227 ) Guido Imperiale
Clean up info parameters in transition to memory ( distributed#9221 ) Guido Imperiale
Move state validation from Scheduler to SchedulerState ( distributed#9224 ) Guido Imperiale
Bump to black 26.3.1 ( distributed#9223 ) Guido Imperiale
Tweak scheduler annotations ( distributed#9220 ) Guido Imperiale
Fix typos ( distributed#9203 ) Dimitri Papadopoulos Orfanos
Bug: dashboard would not show no-worker tasks ( distributed#9215 ) Guido Imperiale
Improve pandas 3.1 compatibility, including add_prefix/add_suffix and DataFrame.drop changes (12414, 12447) `Guido Imperiale`_
Fix groupby compatibility with pandas 2.2 and 2.3 and test against intermediate NumPy and pandas versions (12372) `Guido Imperiale`_
Store empty arrays with Zarr 3.2.0 (12366) `Guido Imperiale`_
Fix quantile and nanquantile behavior (12380) `Guido Imperiale`_
Add support for Cholesky decomposition with complex dtypes (12416) `Niclas Rieger`_
Fix dataframe correctness issues in DataFrame.merge and Series.map (12430, 12432) `nn`_, `Mike Evdokimov`_
Overhaul the publish_dataset extension (9217) `Guido Imperiale`_
Overhaul task stream handling (9230, 9282, 9293) `Guido Imperiale`_, `MohammadYusif`_
Fix a race condition in SubprocessCluster startup (9292) `Guido Imperiale`_
Clean up deprecated Distributed APIs across scheduler, worker, CLI, plugin, and progress handling (9222, 9225, 9238, 9240, 9250) `Guido Imperiale`_
Migrate CI and local workflows to Pixi (12389, 9276) `Guido Imperiale`_
Add AI contribution guidance with AGENTS.md and CLAUDE.md (12415, 9279) `Guido Imperiale`_
Add PyPI release workflows with GitHub Actions Trusted Publishing (12452, 9297) `Matthew Rocklin`_
Silence deprecations in global config if local config overrides them @crusaderky
dask.array.linalg.svd @ayanbag (#12292)test_tokenize_range_index fails if cityhash is not installed @crusaderky (#12286)See the Changelog for more information.
Preliminary Python 3.14t support ( dask#12223 ) Guido Imperiale
Bokeh 3.9.0 compatibility ( distributed#9205 ) Dimitri Papadopoulos Orfanos
Additional changes
docs: document approximate algorithm and Dask-specific params in describe() ( dask#12300 ) Maxime Grenu_
docs: clarify coarsen reduction function contract ( dask#12314 ) monkeyjack123_
Fix misleading TypeError for scalar overflow in dask.array elemwise ( dask#12301 ) Maxime Grenu_
Stricter warnings filter ( dask#12274 ) Guido Imperiale
Clean up obsolete PANDAS_GE markers ( dask#12279 ) Guido Imperiale
Remove mention of obsolete default value for ‘boundary’ parameter. ( dask#12304 ) Marianne Corvellec_
Pandas in 3.14t CI ( dask#12284 ) Guido Imperiale
Quadratic definition time in xarray.DataArray.to_zarr(compute=False) ( dask#12299 ) Guido Imperiale
test_tokenize_range_index fails if cityhash is not installed ( dask#12286 ) Guido Imperiale
Bump minimum version of scipy ( dask#12271 ) Guido Imperiale
Fix flaky categorical concat test ( dask#12276 ) Harshith J_
Doc: document Zarr compression options for to_zarr ( dask#12269 ) Harshith J_
Disable the GIL on 3.14t Windows CI ( dask#12280 ) Guido Imperiale
Update obsolete pandas URLs ( dask#12278 ) Guido Imperiale
Suppress warning: Consolidated metadata is not part of Zarr 3 ( dask#12273 ) Guido Imperiale
Pandas4Warning: Copy-on-Write is always enabled with pandas >= 3.0 ( dask#12272 ) Guido Imperiale
Disable the GIL in 3.14t CI ( dask#12270 ) Guido Imperiale
Propagate contextvars to worker threads; catch warnings in 3.14t ( dask#12224 ) Guido Imperiale
Fix bugs in env.yaml / pytest.xml upload ( dask#12266 ) Guido Imperiale
Added full_matrices parameter to dask.array.linalg.svd ( dask#12292 ) Ayan Bag_
fix: zarr.create_array for better backward compatibility ( dask#12291 ) Wouter-Michiel Vierdag
Silence deprecations in global config if local config overrides them ( dask#12315 ) Guido Imperiale
Fix Total CPU % on /workers tab to normalize by total nthreads ( distributed#9195 ) Ernest Provo
setproctitle: avoid being caught by dask.config; add to test envs ( distributed#9202 ) Guido Imperiale
Add return type annotation for Client._register_plugin ( distributed#9201 ) Simon-Martin Schröder
docs: fix Scheduler.close docstring ( distributed#9198 ) Chase Naples_
Fix Total CPU % on /workers tab to normalize by total nthreads ( distributed#9195 ) Ernest Provo
XFAIL test_handle_null_partitions_2 ( distributed#9191 ) Guido Imperiale
Type hints for Future.status ( distributed#9188 ) Navid_
Pin sphinx=8 ( distributed#9190 ) Guido Imperiale
Preliminary Python 3.14t support (12223) `Guido Imperiale`_
Bokeh 3.9.0 compatibility (9205) `Dimitri Papadopoulos Orfanos`_
See the Changelog for more information.
Ignore deprecationwarning on np.fix @TomAugspurger
zarr_read_kwargs to mode argument @melonora (#12205)test_merge_groupby_to_frame @TomAugspurger (#12244)See the Changelog for more information.
dask.dataframe now requires PyArrow 16 or greater (was 14)
Have **kwargs in to_zarr follow zarr-python API and add mode argument ( dask#12205 ) Wouter-Michiel Vierdag
Note
Passing on io-related arguments in **kwargs in to_zarr will be deprecated and read_kwargs argument as well as zarr_array_kwargs (dict) introduced in 2025.12.0 has been removed. If you passed on either mode or read_only as **kwargs or read_kwargs in to_zarr , please use the new mode argument. The read_only argument can still be passed on, but it will give a warning and have no effect (given that to_zarr is meant to write this should not be an issue). For now no error will be thrown. **kwargs in to_zarr has been renamed as **zarr_array_kwargs to indicate that this directly follows the zarr-python API of Group.create_array when zarr>v3.0.0 and zarr.create for zarr<v3.0.0 . Please see dask.array.to_zarr() for more.
Additional changes
Minimum version of optional dependency h5py bumped to 3.7.0 (was 3.4.0)
Minimum version of optional dependency python-snappy bumped to 0.7.1 (was 0.6.0)
Minimum version of optional dependency tiledb bumped to 0.27.0 (was 0.12.0)
See the Changelog for more information.
Fix XSS vulnerability CVE-2026-23528 Jacob Tomlinson
Support duck-typed Futures in task graph processing ( dask#12213 ) Matthew Rocklin
Additional changes
Remove the Python 2 Comment ( dask#12229 ) Vipin Kataria
Fix changelog: distributed-pr -> pr-distributed ( dask#12227 ) Matthew Plough
Support duck-typed Futures in task graph processing ( dask#12213 ) Matthew Rocklin
Relax test_serialization ( dask#12226 ) Guido Imperiale
[cosmetic] Reorganise dependency groups in CI environment files ( dask#12222 ) Guido Imperiale
Review _array_expr_enabled() ( dask#12217 ) Guido Imperiale
Increase coverage; lower codecov threshold to pass ( dask#12214 ) Guido Imperiale
Test array expr on mindeps ( dask#12216 ) Guido Imperiale
Disable some Mac builds ( dask#12218 ) Guido Imperiale
Typing tweaks ( dask#12215 ) Guido Imperiale
[CI] unbreak codecov ( dask#12211 ) Guido Imperiale
Test array expr on Python 3.14 ( dask#12212 ) Guido Imperiale
Fix pickle compatibility for Python 3.14 ( dask#12206 ) Matthew Rocklin
Remove deprecated dask._compatibility.entry_points ( dask#12202 ) Guido Imperiale
Tweak MacOS CI ( dask#12200 ) Guido Imperiale
Remove obsolete CI pins ( dask#12199 ) Guido Imperiale
Fix XSS vulnerability CVE-2026-23528 Jacob Tomlinson
Clean up obsolete pins in CI ( distributed#9172 ) Guido Imperiale
Fix incompatibility of pyparsing vs. packaging in mindeps CI ( distributed#9170 ) Guido Imperiale
Bump mypy; fix mypy failure ( distributed#9171 ) Guido Imperiale
Fix XSS vulnerability CVE-2026-23528 `Jacob Tomlinson`_
Support duck-typed Futures in task graph processing (12213) `Matthew Rocklin`_
Remove deprecated dask._compatibility.entry_points @crusaderky
test_serialization @crusaderky (#12226)_array_expr_enabled() @crusaderky (#12217)dask._compatibility.entry_points @crusaderky (#12202)See the Changelog for more information.
Stable sort in Series.value\_counts for pandas 3.x @TomAugspurger
test_ufunc_meta for upstream-dev failure @TomAugspurger (#12170)See the Changelog for more information.
More improvements for pandas 3.x Tom Augspurger
Support zarr sharding through create_array ( dask#12153 ) Wouter-Michiel Vierdag
Various improvements for project linting and type hinting Dimitri Papadopoulos Orfanos
Add new “optimization.tune.active” configuration option to disable partition fusion ( dask#12194 ) Richard (Rick) Zamora
Additional changes
Stable sort in Series.value_counts for pandas 3.x ( dask#12191 ) Tom Augspurger
Add new “optimization.tune.active” configuration option to disable partition fusion ( dask#12194 ) Richard (Rick) Zamora
Build llms.txt files in Sphinx documentation ( dask#12192 ) Jacob Tomlinson
Support zarr sharding through create_array ( dask#12153 ) Wouter-Michiel Vierdag
Support min/max of datetime ( dask#12183 ) Julia Signell
pandas 3.x compatibility ( dask#12180 ) Tom Augspurger
Minimal version of setuptools-scm ( dask#12184 ) Dimitri Papadopoulos Orfanos
Update test_ufunc_meta for upstream-dev failure ( dask#12170 ) Tom Augspurger
Upstream compat ( dask#12165 ) Tom Augspurger
Enforce a few more ruff rules ( dask#12157 ) Dimitri Papadopoulos Orfanos
Enforce ruff/refurb rules (FURB) ( dask#12144 ) Dimitri Papadopoulos Orfanos
DEP: bump minimal requirement on toolz (0.10.0 -> 0.12.0) ( dask#12163 ) Clément Robert
Fix execution stop in da.to_zarr due to (misleading) PerformanceWarning raised as exception ( dask#12161 ) Marvin Albert
Use f-string interpolation where possible ( dask#12140 ) Dimitri Papadopoulos Orfanos
pre-commit black hook: use implicit defaults ( dask#12156 ) Dimitri Papadopoulos Orfanos
Enforce ruff/pygrep-hooks rules (PGH) ( dask#12143 ) Dimitri Papadopoulos Orfanos
Apply Repo-Review rules ( dask#12148 ) Dimitri Papadopoulos Orfanos
Document groupby: split_every, split_out ( dask#12135 ) Jayesh Manani
isort → ruff ( dask#12149 ) Dimitri Papadopoulos Orfanos
Enforce ruff/pyupgrade rule UP031 ( dask#12137 ) Dimitri Papadopoulos Orfanos
Replace pre-commit hook with ruff rule ( dask#12142 ) Dimitri Papadopoulos Orfanos
Fix reify to handle sparse arrays and other objects without len ( dask#12103 ) Gautham Hullikunte
Ruff supersedes absolufy-imports ( dask#12141 ) Dimitri Papadopoulos Orfanos
Enforce ruff/pyupgrade rule UP032 ( dask#12136 ) Dimitri Papadopoulos Orfanos
Typing fixes ( distributed#9159 ) Jacob Tomlinson
Explicit setuptools-scm minimum version ( distributed#9160 ) Jacob Tomlinson
Enforce ruff rules (RUF) ( distributed#9153 ) Dimitri Papadopoulos Orfanos
Clean up MANIFEST.in ( distributed#9149 ) Dimitri Papadopoulos Orfanos
isort → ruff ( distributed#9152 ) Dimitri Papadopoulos Orfanos
Ruff supersedes absolufy-imports ( distributed#9154 ) Dimitri Papadopoulos Orfanos
Bump minimum supported toolz to 0.12.0 ( distributed#9151 ) James Bourbeau
flake8, bugbear, pyupgrade → ruff ( distributed#9147 ) Dimitri Papadopoulos Orfanos
Fix typos found by codespell ( distributed#9145 ) Dimitri Papadopoulos Orfanos
Clean up setuptools-specific configuration ( distributed#9150 ) Dimitri Papadopoulos Orfanos
PEP 639 compliance ( distributed#9146 ) Dimitri Papadopoulos Orfanos
Update black ( distributed#9148 ) Dimitri Papadopoulos Orfanos
Fix empty progress bar ( distributed#9144 ) Jacob Tomlinson
Exclude broken tblib versions in CI ( distributed#9141 ) Jacob Tomlinson
More improvements for pandas 3.x `Tom Augspurger`_
Support zarr sharding through create_array (12153) `Wouter-Michiel Vierdag`_
Various improvements for project linting and type hinting `Dimitri Papadopoulos Orfanos`_
Add new "optimization.tune.active" configuration option to disable partition fusion (12194) `Richard (Rick) Zamora`_
Fix deprecated quantile 'interpolation' being passed to numpy @djhoese
MANIFEST.in @DimitriPapadopoulos (#12041)np.accumulate workaround comment @jacobtomlinson (#12129)test\_parquet for pyarrow==22.0 @TomAugspurger (#12116)pip pin for docs @jrbourbeau (#12102)meta arguments in GroupByApply @rjzamora (#12099)See the Changelog for more information.
Use shard shape when available in to_zarr ( dask#12105 ) Davis Bennett
Improve worker and nanny support for ipv6 ( distributed#9133 ) Jianyu Sun
Linting and type hinting improvements across the codebase
Additional changes
Replace versioneer with setuptools-scm ( dask#12133 ) Jacob Tomlinson
Apply ruff/Pylint Refactor rules (PLR) ( dask#12010 ) Dimitri Papadopoulos Orfanos
Remove files from MANIFEST.in ( dask#12041 ) Dimitri Papadopoulos Orfanos
Stabilize test_filter_nonpartition_columns ( dask#12131 ) DongWon
Enforce ruff/pyupgrade rules UP007 and UP033 ( dask#12125 ) Dimitri Papadopoulos Orfanos
Update np.accumulate workaround comment ( dask#12129 ) Jacob Tomlinson
flake8 , bugbear , pyupgrade → ruff ( dask#12002 ) Dimitri Papadopoulos Orfanos
Adjust pyarrow version skip in test_parquet ( dask#12124 ) Tom Augspurger
Fix ufunc in dask.array.cumreduction ( dask#12119 ) Tony Ding
Fix docs footer ( dask#12120 ) Jacob Tomlinson
Use integer multiple of shard shape when rechunking in to_zarr ( dask#12106 ) Davis Bennett
Ensure that the shard shape is used as the default chunk shape for sharded Zarr arrays ( dask#12104 ) Davis Bennett
Skip test_parquet for pyarrow==22.0 ( dask#12116 ) Tom Augspurger
Clean up setuptools-specific configuration ( dask#12040 ) Dimitri Papadopoulos Orfanos
PEP 639 compliance ( dask#12024 ) Dimitri Papadopoulos Orfanos
Fix deprecated quantile interpolation being passed to numpy ( dask#12108 ) David Hoese
Add uv.lock to .gitignore ( dask#12110 ) Jacob Tomlinson
Use shard shape when available in to_zarr ( dask#12105 ) Davis Bennett
Add more optional dependencies to Python 3.13 CI builds ( dask#12100 ) James Bourbeau
Remove pip pin for docs ( dask#12102 ) James Bourbeau
Address collection-based meta arguments in GroupByApply ( dask#12099 ) Richard (Rick) Zamora
Replace versioneer with setuptools-scm ( distributed#9137 ) Jacob Tomlinson
Improve worker and nanny support for ipv6 ( distributed#9133 ) Jianyu Sun
Fix CI Multiple aliased keys in file /Users/runner/.condarc ( distributed#9136 ) Jacob Tomlinson
Remove pip pin for docs ( distributed#9132 ) James Bourbeau
Remove UCX configuration schema ( distributed#9127 ) Peter Andreas Entschev
Add generic type support to Future and Client methods ( distributed#9123 ) Simon-Martin Schröder
Use shard shape when available in to_zarr (12105) `Davis Bennett`_
Improve worker and nanny support for ipv6 (9133) `Jianyu Sun`_
Linting and type hinting improvements across the codebase
Use updated docs theme @jacobtomlinson
dask.array.cumprod does not deal with dtype @tonyyuyiding (#12097)dask.dataframe.read_sql_query() @jacobtomlinson (#12091)\_ExprSequence.\_simplify\_down @rjzamora (#12081)_array_like_safe @ilan-gold (#12078)See the Changelog for more information.
Several Dask Array bug fixes including dask#12097 , dask#12089 , dask#12088 , and dask#12090 .
Additional changes
Use updated docs theme ( dask#12093 ) Jacob Tomlinson
Fix: dask.array.cumprod does not deal with dtype ( dask#12097 ) Tony Ding
CuPy compatibility for percentile ( dask#12098 ) Tom Augspurger
Avoid using methods.concat on empty lists ( dask#12096 ) Tony Ding
Add distribution check for optional dependencies ( dask#12087 ) James Bourbeau
Fix percentile inconsistencies ( dask#12088 ) Oisin-M
Fix warning in test_ufunc_where_no_out ( dask#12094 ) Tom Augspurger
Fix/choose trivial case ( dask#12090 ) Oisin-M
Add input validation on dask.dataframe.read_sql_query() ( dask#12091 ) Jacob Tomlinson
Numpy 2.2 updates for cov function with tests ( dask#12079 ) Mike McCarty
Fix nanvar ( dask#12089 ) Oisin-M
Document manually triggering the conda-forge bots ( dask#12083 ) Jacob Tomlinson
Fix mixed HLG/Expr handling in _ExprSequence._simplify_down ( dask#12081 ) Richard (Rick) Zamora
Add dask.tokenize to API docs ( dask#12080 ) Username46786
CreateOverlappingPartitions : Add before and after to prepend name ( dask#11965 ) Fabien Aulaire
Fix scipy.sparce.csc_matrix scalar declaration in _array_like_safe ( dask#12078 ) Ilan Gold
Update docs theme and remove docs env pins ( distributed#9125 ) Jacob Tomlinson
Add worker name as prefix to ThreadPoolExecutor name ( distributed#9120 ) Maneesh Sutar
Skip hanging SSH tests on Windows ( distributed#9115 ) Jacob Tomlinson
Fix macOS CI failure during job startup ( distributed#9113 ) Jacob Tomlinson
Prevent task stream dashboard showing 1970 date ( distributed#9109 ) Guillaume Eynard-Bontemps
Several Dask Array bug fixes including 12097, 12089, 12088, and 12090.
See the Changelog for more information.
See the Changelog for more information.
This is a backport security release only.
See CVE-2026-23528 for more details.
This is a backport security release only.
See CVE-2026-23528 for more details.
Bump scientific-python/issue-from-pytest-log-action from 1.3.0 to 1.4.0 @[dependabot[bot]]
.groups @TomAugspurger (#12071)See the Changelog for more information.
Avoid unconditional pyarrow dependency in dataframe.backends ( dask#12075 ) Tom Augspurger
pandas 3.x compatibility for .groups ( dask#12071 ) Tom Augspurger
Additional changes
Avoid unconditional pyarrow dependency in dataframe.backends ( dask#12075 ) Tom Augspurger
pandas 3.x compatibility for .groups ( dask#12071 ) Tom Augspurger
Expose details about worker start timeout in the exception message ( distributed#9092 ) Taylor Braun-Jones
pynvml => nvidia-ml-py in CI ( distributed#9111 ) Jacob Tomlinson
Avoid unconditional pyarrow dependency in dataframe.backends (12075) `Tom Augspurger`_
pandas 3.x compatibility for .groups (12071) `Tom Augspurger`_
MAINT: address NumPy deprecation in np.minimum @MarcoGorelli
0 scalar setting for scipy.sparse @ilan-gold (#12027)take @keewis (#11998)np.minimum @MarcoGorelli (#12059)pyarrow chunked array conversion @jrbourbeau (#12034)xfail condition for pyarrow large\_string issue @jrbourbeau (#12032)name not passed to blockwise in map_blocks @ilan-gold (#11952)See the Changelog for more information.
pandas 3.x compatibility ( dask#12025 ) Tom Augspurger
Remove protocol=”ucx” support in favor of distributed-ucxx ( distributed#9105 ) Peter Andreas Entschev
Additional changes
Fix 0 scalar setting for scipy.sparse ( dask#12027 ) Ilan Gold
Workaround failing upstream-dev tests ( dask#12061 ) Tom Augspurger
avoid instantiating a potentially very large arange in take ( dask#11998 ) Justus Magin
MAINT: address NumPy deprecation in np.minimum ( dask#12059 ) Marco Edward Gorelli
CI fixes ( dask#12058 ) Tom Augspurger
MAINT: Address NumPy DeprecationWarning ( dask#12056 ) Marco Edward Gorelli
Fix test_enforce_columns on Python 3.14 ( dask#12047 ) Elliott Sales de Andrade
Fix “th” –> “the” typo in DataFrame SQL docs ( dask#12038 ) Peter A. Jonsson
Advance rng state in permutation ( dask#12031 ) James Bourbeau
Fix pyarrow chunked array conversion ( dask#12034 ) James Bourbeau
Fix xfail condition for pyarrow large_string issue ( dask#12032 ) James Bourbeau
pandas 3.x compatibility ( dask#12025 ) Tom Augspurger
Fix name not propagated correctly in map_blocks ( dask#11952 ) Ilan Gold
Clean tuples dict keys from workers_info in /api/v1/retire_workers. ( distributed#8996 ) Florian Courtial
Remove protocol=”ucx” support in favor of distributed-ucxx ( distributed#9105 ) Peter Andreas Entschev
pandas 3.x compatibility (12025) `Tom Augspurger`_
Remove protocol="ucx" support in favor of distributed-ucxx (9105) `Peter Andreas Entschev`_
CI: update actions location @bsipocz
MapPartitions @rjzamora (#11875)jinja2 is not installed @lukasbindreiter (#11987)builtins.any from being shadowed in dask.array.reductions @m-albert (#11988)\_\_main\_\_ in pickle normalization @jrbourbeau (#11970)upstream CI installation @jrbourbeau (#11976)Dispatch @jrbourbeau (#11974)See the Changelog for more information.
Account for main in pickle normalization ( dask#11970 ) James Bourbeau
Enable column projection in MapPartitions ( dask#11875 ) Richard (Rick) Zamora
Add config option for direct-to-workers ( distributed#9097 ) James Bourbeau
Additional changes
CI: update actions location ( dask#12019 ) Brigitta Sipőcz
Apply ruff/flake8-comprehensions rules (C4) ( dask#12004 ) Dimitri Papadopoulos Orfanos
Apply ruff/flake8-pie rules (PIE) ( dask#12006 ) Dimitri Papadopoulos Orfanos
Apply ruff/Pylint Error rules (PLE) ( dask#12013 ) Dimitri Papadopoulos Orfanos
Apply ruff/Pylint Convention rules (PLC) ( dask#12012 ) Dimitri Papadopoulos Orfanos
Apply ruff/flake8-pyi rules (PYI) ( dask#12007 ) Dimitri Papadopoulos Orfanos
Apply ruff/flake8-simplify rules (SIM) ( dask#12008 ) Dimitri Papadopoulos Orfanos
Apply ruff/Pylint Warning rules (PLW) ( dask#12011 ) Dimitri Papadopoulos Orfanos
Apply ruff/flake8-implicit-str-concat rules (ISC) ( dask#12005 ) Dimitri Papadopoulos Orfanos
Apply ruff/pycodestyle rule E714 ( dask#12000 ) Dimitri Papadopoulos Orfanos
Fix typos found by codespell ( dask#12001 ) Dimitri Papadopoulos Orfanos
Update PyPI URL for official nightly pyarrow repository ( dask#11996 ) Raúl Cumplido
Fall-back to textual repr in case jinja2 is not installed ( dask#11987 ) Lukas Bindreiter
Prevent builtins.any from being shadowed in dask.array.reductions ( dask#11988 ) Marvin Albert
Bump conda-incubator/setup-miniconda from 3.1.1 to 3.2.0 ( dask#11982 )
Skip groupby cov test for pandas 3.x ( dask#11977 ) Tom Augspurger
Fix upstream CI installation ( dask#11976 ) James Bourbeau
Make module name logic more resilient in Dispatch ( dask#11974 ) James Bourbeau
Ensure memray profiler runs on all workers ( distributed#9095 ) James Bourbeau
Update def to class typo in actors docs ( distributed#9091 ) Peter Fackeldey
Bump conda-incubator/setup-miniconda from 3.1.1 to 3.2.0 ( distributed#9090 )
Update persist in tests for async clients ( distributed#9089 ) Tom Augspurger
Fix pyarrow FileInfo import ( distributed#9078 ) James Bourbeau
Make module name logic more resilient in _always_use_pickle_for ( distributed#9086 ) James Bourbeau
Temporarily pin pytest in CI to avoid coverage error ( distributed#9088 ) James Bourbeau
Remove s3fs from testing CI environment ( distributed#9087 ) James Bourbeau
Reuse Comm objects in Scheduler.broadcast ( distributed#9083 ) Tom Augspurger
Fix test_resubmit_nondeterministic_task_different_deps ( distributed#9085 ) James Bourbeau
Account for __main__ in pickle normalization (11970) `James Bourbeau`_
Enable column projection in MapPartitions (11875) `Richard (Rick) Zamora`_
Add config option for direct-to-workers (9097) `James Bourbeau`_
Revert "Dont handle tuple in task\_spec.parse\_input" @fjetter
See the Changelog for more information.
Fixed Dask Array slicing regression introduced in the 2025.5.0 release. See dask#11947 from Florian Jetter for more details. Additional changes
Speed up slicing graph generation ( dask#11945 ) Florian Jetter
Revert “Don’t handle tuple in task_spec.parse_input ” ( dask#11953 ) Florian Jetter
Optimize slicing graph generation ( dask#11946 ) Florian Jetter
Fix xarray slicing regression ( dask#11947 ) Florian Jetter
Don’t handle tuple in task_spec.parse_input ( dask#11948 ) Florian Jetter
Fixed Dask Array slicing regression introduced in the 2025.5.0 release. See 11947 from `Florian Jetter`_ for more details.
Speed up slicing graph generation @fjetter
to_dask_array for single partition @jrbourbeau (#11931)setup-miniconda step in CI @jrbourbeau (#11925)See the Changelog for more information.
Fixed Array setitem when both the array and the indexer have unknown shape. See dask#11753 from Tom Augspurger for more details.
Fixed several delayed graph handling issues introduced in the 2025.4.0 release. See dask#11917 , dask#11907 , and distributed#9071 from Florian Jetter for more details.
Additional changes
Speed up slicing graph generation ( dask#11945 ) Florian Jetter
Optimize dask order for worst case of get_target ( dask#11935 ) Florian Jetter
Raise on local executor if tasks are missing dependency ( dask#11944 ) Florian Jetter
Fix to_dask_array for single partition ( dask#11931 ) James Bourbeau
Ensure parquet plan is fully cached during optimization ( dask#11933 ) Florian Jetter
Better documentation for expression system ( dask#11915 ) Florian Jetter
Simplify (and speed up) culling ( dask#11899 ) Florian Jetter
Update pre-commit ( dask#11926 ) Florian Jetter
Don’t run post setup-miniconda step in CI ( dask#11925 ) James Bourbeau
Try to pin pip for readthedocs ( dask#11923 ) Florian Jetter
Fix windows CI ( dask#11919 ) Florian Jetter
Use stable crick for py310 ( distributed#9072 ) Florian Jetter
Remove internal dependencies mapping in update_graph ( distributed#9036 ) Florian Jetter
Partially forgotten dependencies ( distributed#9068 ) Florian Jetter
Replace filesystem-spec in CI environment with fsspec ( distributed#9069 ) James Bourbeau
Ensure actors set erred state properly in case of worker failure ( distributed#9067 ) Florian Jetter
Refactor timeouts in start cluster ( distributed#9062 ) Florian Jetter
Fix workers / threads / memory displayed in client repr ( distributed#9066 ) James Bourbeau
Pin pip for readthedocs ( distributed#9063 ) Florian Jetter
Skip TLS functional tests ( distributed#9061 ) Florian Jetter
Ensure client submit does not serialize unnecessarily ( distributed#9057 ) Florian Jetter
Fixed Array setitem when both the array and the indexer have unknown shape. See 11753 from `Tom Augspurger`_ for more details.
Fixed several delayed graph handling issues introduced in the 2025.4.0 release. See 11917, 11907, and 9071 from `Florian Jetter`_ for more details.
Ensure only HLGs are probited reuse @fjetter
See the Changelog for more information.
This release contains several graph optimization fixes for issues introduced in the 2025.4.0 release.
See dask#11906 , dask#11898 , dask#11903 , and dask#11904 by Florian Jetter for more details. Additional changes
Implement ufuncs and gufunc for array-expr ( dask#11818 ) Patrick Hoefler
Implement map_overlap for array-expr ( dask#11822 ) Patrick Hoefler
This release contains several graph optimization fixes for issues introduced in the 2025.4.0 release.
See 11906, 11898, 11903, and 11904 by `Florian Jetter`_ for more details.
Ensure Future value is in da.from\_delayed task graph @TomAugspurger
DataFrame.isin as object type numpy arrays @mroeschke (#11869)map_blocks in array.store to avoid materialization and dropping of annotations @fjetter (#11844)h5py from upstream CI job @jrbourbeau (#11847)Expr.__setattr__ for subclasses @TomAugspurger (#11845)See the Changelog for more information.
When computing multiple Dask-Expr backed collections like DataFrames, they are now optimized together instead of individually.
Graph materialization and low level optimization is now being performed on the scheduler of a distributed cluster (if available).
New kwarg force for DataFrame.shuffle which signals the optimizer to not drop the shuffle during optimization.
Collections that are passed to Dask methods as arguments are now properly optimized. If multiple collections are passed as arguments they will be optimized together. Collections passed this way are prohibited from being being reused, i.e. if the collection is used again in another function call it will be computed again. This pattern is used to avoid pipeline breakers which typically drive memory usage. Avoiding those should reduce memory pressure on the cluster but can cause runtime regressions.
(Special case of above point) Collections passed to Delayed objects are now optimized automatically.
Support for custom low level optimizers removed.
Top level dask.optimize will now always trigger graph materialization. Previously this was not always the case. This also causes any low level HLG annotations to be dropped.
DataFrame and Array compute results are now always concatenated on the cluster. Previously, the behavior was dependent on the API used to call compute ( dask.compute , DaskCollection.compute , or Client.compute ).
dask.base.collections_to_dsk has been renamed to collections_to_expr and no longer returns a HighLevelGraph or dict object but instead guarantees an dask._expr.Expr object. Further, it no longer performs low level optimization immediately but instead delays until the Expr instance is materialized, i.e. the returned object is no longer a mapping such that converting it to dict or iterating over it is not possible any more.
Additional changes
Ensure Future value is in da.from_delayed task graph ( dask#11896 ) Tom Augspurger
Fix annotations passed to delayed ( dask#11893 ) Florian Jetter
Migrate delayed unpack_collections ( dask#11881 ) Florian Jetter
Remove Pub / Sub references from docs ( dask#11891 ) James Bourbeau
Ensure only classes without custom init are singletons ( dask#11886 ) Florian Jetter
Remove custom initializers for delayed expressions ( dask#11888 ) Florian Jetter
Fix persisting multiple DFs at the same time ( dask#11887 ) Florian Jetter
Avoid always parsing list inputs to DataFrame.isin as object type numpy arrays ( dask#11869 ) Matthew Roeschke
Unskip pandas-dev cov / corr tests ( dask#11873 ) Tom Augspurger
HLG blockwise fix ( dask#11871 ) Florian Jetter
Ensure annotations for HLG objects are properly generated ( dask#11866 ) Florian Jetter
Factor out singleton logic from base Expr class ( dask#11868 ) Florian Jetter
Ensure HLGs are using dependencies properly in optimization ( dask#11859 ) Florian Jetter
Ensure dictionaries tokenize deterministically ( dask#11867 ) Florian Jetter
Ensure default dask scheduler only compute what’s needed ( dask#11861 ) Florian Jetter
Faster tokenization of pd.RangeIndex ( dask#11863 ) Florian Jetter
Update link to Quansight in community doc ( dask#11860 ) Pavithra Eswaramoorthy
Relax tolerance in autocorr test ( dask#11857 ) Tom Augspurger
Use map_blocks in array.store to avoid materialization and dropping of annotations ( dask#11844 ) Florian Jetter
Ensure repartition does not trigger memory size computation during lowering (i.e. on the scheduler) ( dask#11855 ) Florian Jetter
Support args and kwargs for rolling aggregations ( dask#11856 ) Florian Jetter
Remove nightly h5py from upstream CI job ( dask#11847 ) James Bourbeau
Ensure HLGExpr tokenize uniquely ( dask#11849 ) Florian Jetter
Do not inject median in describe for pandas 3 ( dask#11846 ) Florian Jetter
Fixed Expr.setattr for subclasses ( dask#11845 ) Tom Augspurger
Wrap HLGs in an Expr to avoid Client side materialization ( dask#11736 ) Florian Jetter
Improve error when submitting work from a closed client ( distributed#9049 ) James Bourbeau
Return a default value if address resolution fails ( distributed#9051 ) Sandro
Avoid deepcopy when submitting graph ( distributed#8633 ) Florian Jetter
Dynamically scale heartbeat and scheduler_info intervals ( distributed#9046 ) Florian Jetter
Speed up process startup time by avoiding importing packages on version check ( distributed#9048 ) Florian Jetter
Reduce size of scheduler_info ( distributed#9045 ) Florian Jetter
Cache WorkerState host property ( distributed#9044 ) Florian Jetter
Clear ci env cache ( distributed#9047 ) Florian Jetter
Remove deprecated Pub / Sub ( distributed#9039 ) Florian Jetter
Perform explicit culling step only if LLG is submitted ( distributed#9040 ) Florian Jetter
Do not fully materialize global annotations by type ( distributed#9035 ) Florian Jetter
Allow nested worker_client calls ( distributed#9038 ) George Sakkis
Dump ci cache ( distributed#9037 ) Florian Jetter
Scheduler type annotations ( distributed#9030 ) Florian Jetter
Reduce dask.order overhead by removing stripped_dep computation ( distributed#9031 ) Florian Jetter
Use Expr instead of HLG ( distributed#9008 ) Florian Jetter
When computing multiple Dask-Expr backed collections like DataFrames, they are now optimized together instead of individually.
Graph materialization and low level optimization is now being performed on the scheduler of a distributed cluster (if available).
New kwarg force for DataFrame.shuffle which signals the optimizer to not drop the shuffle during optimization.
Collections that are passed to Dask methods as arguments are now properly optimized. If multiple collections are passed as arguments they will be optimized together. Collections passed this way are prohibited from being being reused, i.e. if the collection is used again in another function call it will be computed again. This pattern is used to avoid pipeline breakers which typically drive memory usage. Avoiding those should reduce memory pressure on the cluster but can cause runtime regressions.
(Special case of above point) Collections passed to Delayed objects are now optimized automatically.
Support for custom low level optimizers removed.
Top level dask.optimize will now always trigger graph materialization. Previously this was not always the case. This also causes any low level HLG annotations to be dropped.
DataFrame and Array compute results are now always concatenated on the cluster. Previously, the behavior was dependent on the API used to call compute (dask.compute, DaskCollection.compute, or Client.compute).
dask.base.collections_to_dsk has been renamed to collections_to_expr and no longer returns a HighLevelGraph or dict object but instead guarantees an dask._expr.Expr object. Further, it no longer performs low level optimization immediately but instead delays until the Expr instance is materialized, i.e. the returned object is no longer a mapping such that converting it to dict or iterating over it is not possible any more.
See the Changelog for more information.
Fix dataset info cache assignment @fjetter
to\_orc to DataFrame API @TomAugspurger (#11807)to_pandas_dispatch registration for cudf @rjzamora (#11799)apply\_gufunc @phofl (#11683)dask.dataframe.__all__ @flying-sheep (#11782)dask.bag @flying-sheep (#11781)dask.array.__all__ @flying-sheep (#11780)asarray(..., like=...) vs. scipy.sparse objects @crusaderky (#11755)flaky to tests extra @TomAugspurger (#11770)arange(..., like=x) embeds the graph of x @crusaderky (#11754)map_partitions returns Series object if function returns scalar @fjetter (#11756)See the Changelog for more information.
apply_ufunc requires the core dimension to have chunksize=-1 . The underlying rechunking operation will automatically adjust the chunksize of the core dimension but keep the other dimensions the same. This can cause exploding chunksizes under the hood.
This release adds an intermediate step that resizes the non-core dimensions by the same factor that the core dimension will increase to keep the maximum chunksize under control. This behavior is automatically enabled when allow_rechunk=True is set.
import xarray as xr import dask.array as da arr = xr . DataArray ( da . random . random (( 1 , 750 , 45910 ), chunks = ( 1 , "auto" , - 1 )), dims = [ "band" , "y" , "x" ], ) result = arr . interp ( y = arr . coords [ "y" ], method = "linear" , )
Previously
Individual chunks are exploding to 25 GiB, likely causing out of memory errors.
Now
Dask will now automatically split individual chunks into chunks that will have the same chunksize minus a small tolerance.
Additional changes
Fix dataset info cache assignment ( dask#11840 ) Florian Jetter
Expr setattr ( dask#11836 ) Florian Jetter
Follow up to expression tokenization caching ( dask#11837 ) Florian Jetter
Consolidate getattr for expr classes ( dask#11835 ) Florian Jetter
Reduce pickle size of ReadParquet expression ( dask#11797 ) Florian Jetter
arange loses precision on ~2**63 ( dask#11801 ) Guido Imperiale
Remove numbagg from upstream build ( dask#11821 ) Patrick Hoefler
Dispatch to numbagg for nanmedian and nanquantile ( dask#11817 ) Patrick Hoefler
Make missing meta warning more ergonomic ( dask#11814 ) Patrick Hoefler
Remove name doc from from_pandas ( dask#11812 ) Patrick Hoefler
Implement an Array Scalar ( dask#11810 ) Patrick Hoefler
Added to_orc to DataFrame API ( dask#11807 ) Tom Augspurger
Implement reverse indexing for DataFrames ( dask#11803 ) Patrick Hoefler
Add lazy to_pandas_dispatch registration for cudf ( dask#11799 ) Richard (Rick) Zamora
Fix missing imports in array-expr ( dask#11796 ) Florian Jetter
Cache tokens on expressions and restore after pickle roundtrip ( dask#11791 ) Florian Jetter
Use random dashboard ports for LocalCluster in distributed tests ( dask#11795 ) Florian Jetter
Implement slicing for array-expr ( dask#11783 ) Patrick Hoefler
Never use an asynchronous Client when calling top level compute function ( dask#11790 ) Florian Jetter
Refactor import tests ( dask#11794 ) Florian Jetter
Migrate base.unpack_collections to Task class ( dask#11793 ) Florian Jetter
Ensure map_blocks generates unique tokens ( dask#11792 ) Florian Jetter
Speed up normalize_pickle by 50 percent ( dask#11788 ) Florian Jetter
Fix divisions calculation with duplicates ( dask#11787 ) Patrick Hoefler
Fix assign align for duplicated divisions ( dask#11786 ) Patrick Hoefler
Ensure concat optimize project does not raise ( dask#11784 ) Florian Jetter
Add array-expr from_array ( dask#11772 ) Patrick Hoefler
Keep chunksizes consistent in apply_gufunc ( dask#11683 ) Patrick Hoefler
Test dask.dataframe.all ( dask#11782 ) Philipp A.
Add all to dask.bag ( dask#11781 ) Philipp A.
Add test for dask.array.all ( dask#11780 ) Philipp A.
Bump JamesIves/github-pages-deploy-action from 4.7.2 to 4.7.3 ( dask#11777 )
Export dask.array members ( dask#11779 ) Philipp A.
Fix sorted_divisions_locations with duplicates ( dask#11773 ) Tom Augspurger
Fix small typo in best-practices.rst ( dask#11775 ) Sergey Kolesnikov
Allow unknown chunks in blockwise adjust_chunks ( dask#11769 ) Lindsey Gray
Fix crash in asarray(..., like=...) vs. scipy.sparse objects ( dask#11755 ) Guido Imperiale
Remove flaky optional dependency ( dask#11771 ) Tom Augspurger
Add support for scipy sparray ( dask#11750 ) Philipp A.
Added flaky to tests extra ( dask#11770 ) Tom Augspurger
Ensure divisions are plain scalars ( dask#11767 ) Tom Augspurger
Remove divisions code duplication ( dask#11764 ) Florian Jetter
Ensure divisions not diverging from npartitions in Merge ( dask#11762 ) Florian Jetter
Skip test_visualize_int_overflow on windows ( dask#11761 ) Florian Jetter
Reduce pickle size for tasks ( dask#11687 ) Florian Jetter
Implement unify_chunks and Rechunk ( dask#11692 ) Patrick Hoefler
Fix expression getitem to avoid alignment ( dask#11760 ) Patrick Hoefler
arange(..., like=x) embeds the graph of x ( dask#11754 ) Guido Imperiale
Simplify assert_divisions ( dask#11745 ) Florian Jetter
Fix Projection logic for Series objects ( dask#11747 ) Patrick Hoefler
Remove bytes as keys ( dask#11757 ) Florian Jetter
Ensure map_partitions returns Series object if function returns scalar ( dask#11756 ) Florian Jetter
Don’t upload env twice ( dask#11748 ) Patrick Hoefler
Fix badges in readme ( distributed#9029 ) Florian Jetter
Properly forward cancellation reason ( distributed#9028 ) Florian Jetter
Fix bokeh circle ( distributed#9026 ) Florian Jetter
Ensure FileInfo can be serialized ( distributed#9025 ) Florian Jetter
Add ipykernel to skipped modules in code sampling ( distributed#9022 ) Matthew Rocklin
SpecCluster: add option to not shut down the scheduler when the cluster is closed ( distributed#9021 ) Taylor Braun-Jones
Fix CI by using client.persist(collection) instead of collection.persist() ( distributed#9020 ) Hendrik Makait
Add redirect from prefix root to status ( distributed#9015 ) Isaac
Bump JamesIves/github-pages-deploy-action from 4.7.2 to 4.7.3 ( distributed#9018 )
Remove bytes keys from tests ( distributed#9017 ) Jacob Tomlinson
apply_ufunc requires the core dimension to have chunksize=-1. The underlying rechunking operation will automatically adjust the chunksize of the core dimension but keep the other dimensions the same. This can cause exploding chunksizes under the hood.
This release adds an intermediate step that resizes the non-core dimensions by the same factor that the core dimension will increase to keep the maximum chunksize under control. This behavior is automatically enabled when allow_rechunk=True is set.
import xarray as xr
import dask.array as da
arr = xr.DataArray(
da.random.random((1, 750, 45910), chunks=(1, "auto", -1)),
dims=["band", "y", "x"],
)
result = arr.interp(
y=arr.coords["y"],
method="linear",
)
Add big array example @jrbourbeau
__setitem__ with dask bool mask @crusaderky (#11728)arange: fix extreme values @crusaderky (#11707)arange: support kwargs @crusaderky (#11710)normalize_token is threadsafe @fjetter (#11709)scikit-image nightly back to upstream CI"" @phofl (#11667)\_\_all\_\_ to init @phofl (#11664)_execute_subgraph @hendrikmakait (#11655)normalize\_chunks @phofl (#11650)concatenate3 in overlap and rechunking graphs @hendrikmakait (#11621)concrete in task graph @hendrikmakait (#11620)dask.array.rechunk with chunks='auto' @schlunma (#11622)scikit-image nightly back to upstream CI" @phofl (#11616)See the Changelog for more information.
This release includes a critical fix that fixes a deadlock that can arise when seceded task are rescheduled, or cancelled and resubmitted, e.g. due to a worker being lost.
See distributed#8991 by Hendrik Makait for more details. Additional changes
Add big array example ( dask#11744 ) James Bourbeau
Fix exploding chunksizes in pad for constant padding ( dask#11743 ) Patrick Hoefler
Move optimize method to base class ( dask#11742 ) Florian Jetter
Add changelog entry for fixed deadlock ( dask#11741 ) Hendrik Makait
Fix graph creation in dask-expr to_delayed ( dask#11739 ) Patrick Hoefler
Remove culling from delayed optimisation ( dask#11737 ) Patrick Hoefler
Compute meta for from_map on the cluster ( dask#11738 ) Patrick Hoefler
Bugs in setitem with dask bool mask ( dask#11728 ) Guido Imperiale
Implement infrastructure, random, blockwise and Elemwise ( dask#11689 ) Patrick Hoefler
array / asarray with both like= and dtype= ( dask#11733 ) Guido Imperiale
Fix annotations warnings test ( dask#11734 ) Patrick Hoefler
Catch warnings when writing to remote storage with to_parquet ( dask#11731 ) Patrick Hoefler
Remove LocalCluster from tests ( dask#11729 ) Patrick Hoefler
Fix partition pruning when using from_array ( dask#11725 ) Patrick Hoefler
Fix concatentation with mixed dtype columns ( dask#11727 ) Patrick Hoefler
arange : fix extreme values ( dask#11707 ) Guido Imperiale
Graph corruption on scalar getitem -> setitem ( dask#11723 ) Guido Imperiale
Never share buffers after compute() ( dask#11697 ) Guido Imperiale
Extract Dask Array from xarray DataArray in from_array ( dask#11712 ) Patrick Hoefler
arange : support kwargs ( dask#11710 ) Guido Imperiale
Ensure normalize_token is threadsafe ( dask#11709 ) Florian Jetter
Expand advise for instance types and processes ( dask#11705 ) Florian Jetter
Drop legacy timeseries implementation ( dask#11704 ) Florian Jetter
Update Dask Cloud Provider documentation to include Nebius as a supported cloud option ( dask#11703 ) Alexander
Fix normalize_chunks when squashing into a single chunk ( dask#11702 ) Patrick Hoefler
Fix positional indexing with newaxis ( dask#11699 ) Patrick Hoefler
Set array backend in scipy-sparse-indexing ( dask#11700 ) Tom Augspurger
Fix value_counts shuffling strategy ( dask#11698 ) Patrick Hoefler
Disentangle core expression class from dataframe specific code ( dask#11688 ) Patrick Hoefler
Bump conda-incubator/setup-miniconda from 3.1.0 to 3.1.1 ( dask#11685 )
Fixup dataframe conversion from array methods ( dask#11684 ) Patrick Hoefler
Remove remaining artifacts of fastparquet ( dask#11682 ) Patrick Hoefler
Remove traceback from sizeof failure warning ( distributed#9006 ) Jacob Tomlinson
Hotfix: Ignore negative occupancy ( distributed#9012 ) Hendrik Makait
Remove expensive tokenization for key uniqueness check ( distributed#9009 ) Patrick Hoefler
Fix CI for changes in from_map ( distributed#9011 ) Patrick Hoefler
Avoid handling stale long-running messages on scheduler ( distributed#8991 ) Hendrik Makait
Bump test_stress timeout ( distributed#9002 ) Tom Augspurger
Poll in test_rmm_metrics test ( distributed#9004 ) Tom Augspurger
Cache occupancy in WorkStealing.balance() ( distributed#9005 ) Hendrik Makait
Homogeneous balancing by accounting for in-flight requests ( distributed#9003 ) Hendrik Makait
Consistent estimation of task duration between stealing, adaptive and occupancy calculation ( distributed#9000 ) Hendrik Makait
Increase default work-stealing interval by 10x ( distributed#8997 ) Hendrik Makait
Remove occupancy plot from status dashboard ( distributed#8995 ) Hendrik Makait
Bump conda-incubator/setup-miniconda from 3.1.0 to 3.1.1 ( distributed#8990 )
This release includes a critical fix that fixes a deadlock that can arise when seceded task are rescheduled, or cancelled and resubmitted, e.g. due to a worker being lost.
See 8991 by `Hendrik Makait`_ for more details.
This enforces the deprecation of the configuration:
This release drops the legacy Dask DataFrame implementation. The API with query planning is now the only available Dask DataFrame implementation.
This enforces the deprecation of the configuration:
dask . config . set ({ "dataframe.query-planning" : False })
Dask-Expr was merged into the dask package as well as the dask/dask repository. It is no longer necessary to install dask-expr separately.
Dask introduced a mechanism that is called root task queuing in 2022. This mechanism allows Dask to detect tasks that are reading data from storage and schedule them defensively to avoid memory pressure on the cluster through overproduction of these tasks. The underlying mechanism was very fragile and failed for specific types of computations like opening multiple zarr stores or loading a large number of netcdf files.
The recent changes in Dask’s task graph representation allow for more robust detection of root tasks. This change makes the detection mechanism independent of the workload running and is especially beneficial for Xarray workloads.
This results in significantly more memory stability and a reduced memory footprint for workloads where root task detection was previously failing and makes the expected memory profile deterministic and independent of the topology of the task graph.
Fix map\_overlap bug where rechunking and trim=False caused inconsistent chunkings @phofl
NestedContainers in case of trivial inputs @hendrikmakait (#11600)TaskRef in Alias @hendrikmakait (#11597)array.push @fjetter (#11576)See the Changelog for more information.
This release reduces the number of Python object references related to tracking tasks by the Dask scheduler. This increases scheduler responsiveness by reducing the time needed to run garbage collection on the scheduler.
See dask#8958 , dask#11608 , dask#11600 , dask#11598 , dask#11597 , and distributed#8963 from Hendrik Makait for more details. Additional changes
Fix map_overlap bug where rechunking and trim=False caused inconsistent chunkings ( dask#11605 ) Patrick Hoefler
Avoid legacy implementation in read-csv ( dask#11603 ) Patrick Hoefler
Remove legacy DataFrame import ( dask#11604 ) Patrick Hoefler
asarray ignores dtype for array inputs ( dask#11586 ) crusaderky
Add back LLM chatbot to Dask docs ( dask#11594 ) dchudz
Bump JamesIves/github-pages-deploy-action from 4.6.9 to 4.7.2 ( dask#11593 )
Migrate dask array creation routines to task spec ( dask#11582 ) James Bourbeau
Migrate most of dask array random to task spec ( dask#11581 ) James Bourbeau
Do not use local function in array.push ( dask#11576 ) Florian Jetter
Bump conda-incubator/setup-miniconda from 3.0.3 to 3.1.0 ( distributed#8922 )
Pick random dashboard port in tests ( distributed#8965 ) Hendrik Makait
Fix formatting for NoValidWorkerException message ( distributed#8967 ) Hendrik Makait
Support pynvml>=11.5 in WSL ( distributed#8962 ) Richard (Rick) Zamora
Bump JamesIves/github-pages-deploy-action from 4.6.9 to 4.7.2 ( distributed#8960 )
This release reduces the number of Python object references related to tracking tasks by the Dask scheduler. This increases scheduler responsiveness by reducing the time needed to run garbage collection on the scheduler.
See 8958, 11608, 11600, 11598, 11597, and 8963 from `Hendrik Makait`_ for more details.
Revert "Add LLM chatbot to Dask docs (#11556)" @dchudz
Task class @fjetter (#11568)Bag graphs to TaskSpec graphs during optimization @fjetter (#11569)scikit-image nightly back to upstream CI @jrbourbeau (#11530)from\_dask\_dataframe import @phofl (#11528)read\_only kwarg in zarr=3 @phofl (#11516)See the Changelog for more information.
This release adds support for Python 3.13. Dask now supports Python 3.10-3.13.
See dask#11456 and distributed#8904 from Patrick Hoefler and James Bourbeau for more details. Additional changes
Revert “Add LLM chatbot to Dask docs ( dask#11556 )” ( dask#11577 ) dchudz
Automatically rechunk if array in to_zarr has irregular chunks ( dask#11553 ) Patrick Hoefler
Blockwise uses Task class ( dask#11568 ) Florian Jetter
Migrate rechunk and reshape to task spec ( dask#11555 ) Patrick Hoefler
Cache svg-representation for arrays ( dask#11560 ) Deepak Cherian
Fix empty input for containers ( dask#11571 ) Florian Jetter
Convert Bag graphs to TaskSpec graphs during optimization ( dask#11569 ) Florian Jetter
Add LLM chatbot to Dask docs ( dask#11556 ) dchudz
Fuse data nodes in linear fusion too ( dask#11549 ) Patrick Hoefler
Migrate slicing code to task spec ( dask#11548 ) Patrick Hoefler
Speed up ArraySliceDep tokenization ( dask#11551 ) Patrick Hoefler
Fix fusing of p2p barrier tasks ( dask#11543 ) Patrick Hoefler
Remove infra/mentions of GPU CI ( dask#11546 ) Charles Blackmon-Luca
Temporarily disable gpuCI update CI job ( dask#11545 ) James Bourbeau
Use BlockwiseDep to implement map_blocks keywords ( dask#11542 ) Patrick Hoefler
Remove optimize_slices ( dask#11538 ) Patrick Hoefler
Make reshape_blockwise a noop if shape is the same ( dask#11541 ) Patrick Hoefler
Remove read-only flag from open_arry in open_zarr ( dask#11539 ) Patrick Hoefler
Implement linear_fusion for task spec class ( dask#11525 ) Patrick Hoefler
Remove recursion from TaskSpec ( dask#11477 ) Florian Jetter
Fixup test after dask-expr change ( dask#11536 ) Patrick Hoefler
Bump codecov/codecov-action from 3 to 5 ( dask#11532 )
Create dask-expr frame directly without roundtripping ( dask#11529 ) Patrick Hoefler
Add scikit-image nightly back to upstream CI ( dask#11530 ) James Bourbeau
Remove from_dask_dataframe import ( dask#11528 ) Patrick Hoefler
Ensure that from_array creates a copy ( dask#11524 ) Patrick Hoefler
Simplify and improve performance of normalize chunks ( dask#11521 ) Patrick Hoefler
Fix flaky nanquantile test ( dask#11518 ) Patrick Hoefler
Fix tests for new read_only kwarg in zarr=3 ( dask#11516 ) Patrick Hoefler
Fix test_jupyter.py::test_shutsdown_cleanly ( distributed#8954 ) Hendrik Makait
Install tornado from conda-forge in Python 3.13 CI ( distributed#8951 ) James Bourbeau
Restore retire workers API ( distributed#8939 ) Florian Jetter
Properly convert finalize dependencies to references ( distributed#8949 ) Hendrik Makait
Block fusion for barrier tasks ( distributed#8944 ) Patrick Hoefler
Remove infra/mentions of GPUCI ( distributed#8946 ) Charles Blackmon-Luca
Temporarily disable gpuCI update CI job ( distributed#8945 ) James Bourbeau
Remove recursion in task spec ( distributed#8920 ) Florian Jetter
Less verbose log messages for remove and register worker ( distributed#8938 ) Florian Jetter
Do not log full worker info in retire_workers ( distributed#8935 ) Florian Jetter
This release adds support for Python 3.13. Dask now supports Python 3.10-3.13.
See 11456 and 8904 from `Patrick Hoefler`_ and `James Bourbeau`_ for more details.
Remove only\_refs parsing option for TaskSpec @fjetter
nanpercentile for dask arrays @phofl (#11505)See the Changelog for more information.
Note
Versions 2024.11.0 and 2024.11.1 included a critical performance regression and should be skipped by every user.
This release deprecates the legacy Dask DataFrame implementation. The old implementation will be removed completely in a future release. Users are encourage to switch to the new implementation now and to report any issues they are facing.
Users are also encourage to check that they are only importing functions from dask.dataframe and not any of the submodules.
Dask Array added new quantile and nanquantile methods. Previously, Dask dispatched to the NumPy implementation, which blocked the GIL a lot. This caused large slowdowns on workers with more than one tread and could lead to runtimes over 200s per chunk.
The new quantile implementation avoids many of these problems and reduces runtime to around 1s per chunk independently of the number of threads.
Using Xarrays rolling(...).construct(...) with Dask Arrays led to very large chunksizes that rarely fit into memory on a single worker.
The underlying operations is a view on the smaller NumPy array, but triggering a copy of the data will lead to very large memory usage.
import xarray as xr import dask.array as da arr = xr . DataArray ( da . ones (( 93504 , 721 , 1440 ), chunks = ( "auto" , - 1 , - 1 )), dims = [ "time" , "lat" , "longitude" ], ) # Initial chunks are ~128 MiB arr . rolling ( time = 30 ) . construct ( "window_dim" )
Previously
Individual chunks are exploding to 10 GiB, likely causing out of memory errors.
Now
Dask will now automatically split individual chunks into chunks that will have the same chunksize minus a small tolerance.
map_overlap now creates smaller and more efficient graphs to keep task graphs generally a lot smaller.
The previous version injected a lot of tasks that weren’t necessary, increasing the number of tasks by a factor of 2-10x of what actually necessary. This caused a lot of stress on the scheduler.
Einstein summation historically led to very large chunksizes if applied to more than one Dask Array. This behavior is inherited from NumPy but led to out of memory errors on workers:
import dask.array as da arr = da . random . random (( 1024 , 64 , 64 , 64 , 64 ), chunks = ( 256 , 16 , 16 , 16 , 16 )) # Initial chunks are 128 MiB result = da . einsum ( "aijkl,amnop->ijklmnop" , arr , arr )
Previously
Individual chunks are exploding to 32 GiB, very likely causing out of memory errors.
Now
The operation keeps individual chunksizes the same.
Additional changes
Add changelog for Dask release ( dask#11502 ) Patrick Hoefler
Minor updates to optional dependencies table ( dask#11503 ) James Bourbeau
Add push for ffill like operations ( dask#11501 ) Patrick Hoefler
Remove func packing for TaskSpec ( dask#11496 ) Florian Jetter
Make tokenization for vindex more efficient ( dask#11493 ) Patrick Hoefler
Cut down runtime of einstein summation test ( dask#11499 ) Patrick Hoefler
Improve test runtime for test_rot90 ( dask#11498 ) Florian Jetter
Disable low level optimization for TaskSpec in Bags ( dask#11495 ) Florian Jetter
Add automatic rechunking to sliding-window-view ( dask#11479 ) Patrick Hoefler
Add load_stored kwarg to dask.array.store ( dask#11465 ) Deepak Cherian
Fix quantile error in two dimensions ( dask#11489 ) Patrick Hoefler
Bump conda-incubator/setup-miniconda from 3.0.4 to 3.1.0 ( dask#11490 )
Update map_blocks docstring ( dask#11491 ) Patrick Hoefler
Fix einsum with empty arrays ( dask#11488 ) Patrick Hoefler
Implement non gil-blocking quantile method ( dask#11473 ) Patrick Hoefler
Use internal keyword for trimming in map_overlap to reduce graph size ( dask#11486 ) Patrick Hoefler
Minor dask order refactor ( dask#11467 ) Florian Jetter
Remove empty tasks from map_overlap ( dask#11483 ) Patrick Hoefler
Fixup auto chunks calculation if single chunk goes below 1 ( dask#11485 ) Patrick Hoefler
Fix CI after pandas upstream changes ( dask#11482 ) Patrick Hoefler
Make sure that block_id and block_info don’t create extra tasks ( dask#11484 ) Patrick Hoefler
Use repeat to build nearest boundary ( dask#9666 ) Jean-Baptiste Bayle
Remove dead code from make_blockwise ( dask#11478 ) Florian Jetter
Patch auto-chunks calculation for rioxarray ( dask#11480 ) Patrick Hoefler
Skip legacy test because of flaky warning ( dask#11475 ) Patrick Hoefler
Unskip a few dask-expr tests ( dask#11474 ) Patrick Hoefler
Keep chunk sizes consistent in einsum ( dask#11464 ) Patrick Hoefler
Improve how normalize_chunks squashes together chunks when “auto” is set ( dask#11468 ) Patrick Hoefler
Fix resolve_aliases when multiple aliases are in graph ( dask#11469 ) Patrick Hoefler
Avoid cyclic import in dask.array ( dask#11472 ) Hendrik Makait
Unskip dataframe test ( dask#11471 ) Patrick Hoefler
Improve dask.order performance for large graphs ( dask#11466 ) Florian Jetter
Ensure that slice(None) just maps the keys ( dask#11450 ) Patrick Hoefler
Fix Task.repr() of unpickled object ( dask#11463 ) Peter Andreas Entschev
Use TaskSpec in local dask execution ( dask#11378 ) Florian Jetter
Adjust accuracy in test_solve_triangular_vector ( dask#11461 ) Florian Jetter
Update Aggregation docstring ( dask#11459 ) Guillaume Eynard-Bontemps
Implement fuse option for delayed objects ( dask#11441 ) Patrick Hoefler
Deprecate legacy dask dataframe implementation ( dask#11437 ) Patrick Hoefler
Fix na casting behavior for groupby.agg with arrow dtypes ( dask#11118 ) Patrick Hoefler
Fix behavior of keys_in_tasks for TaskSpec nodes ( dask#11445 ) Florian Jetter
Convert dtype to int instead of np.uint8 for visualizing large task graphs ( dask#11440 ) Patrick Hoefler
Ensure dependencies are not mutated ( dask#11438 ) Florian Jetter
Full support for task spec in dask.order ( dask#11347 ) Florian Jetter
Remove redundant methods in P2PBarrierTask ( distributed#8924 ) Florian Jetter
Fix skipif condition for test_tell_workers_when_peers_have_left ( distributed#8929 ) Florian Jetter
Ensure ConnectionPool is closed even if network stack swallows CancelledErrors ( distributed#8928 ) Florian Jetter
Fix flaky test_server_comms_mark_active_handlers ( distributed#8927 ) Florian Jetter
Make assumption in P2P’s barrier mechanism explicit ( distributed#8926 ) Hendrik Makait
Adjust timeouts in Jupyter cli test ( distributed#8925 ) Florian Jetter
Add stimulus_id to update_graph plugin hook ( distributed#8923 ) Hendrik Makait
Reduce P2P transfer task overhead ( distributed#8912 ) Hendrik Makait
Disable profiler on Python 3.11 ( distributed#8916 ) Florian Jetter
Fix test_restarting_does_not_deadlock ( distributed#8849 ) Florian Jetter
Adjust popen timeouts for testing ( distributed#8848 ) Florian Jetter
Add retry to shuffle broadcast ( distributed#8900 ) Florian Jetter
Fix test_shuffle_with_array_conversion ( distributed#8909 ) Florian Jetter
Refactor some tests ( distributed#8908 ) Florian Jetter
Graduate dask-expr from contrib to core project ( distributed#8911 ) Hendrik Makait
Skip test_tell_workers_when_peers_have_left on py10 ( distributed#8910 ) Florian Jetter
Internal cleanup of P2P code ( distributed#8907 ) Hendrik Makait
Use Task class instead of tuple ( distributed#8797 ) Florian Jetter
Increase connect timeout for test_tell_workers_when_peers_have_left ( distributed#8906 ) Florian Jetter
Remove dispatching in TaskCollection ( distributed#8903 ) Florian Jetter
Deduplicate requests to scheduler in P2P ( distributed#8899 ) Hendrik Makait
Add configurations for rootish taskgroup threshold ( distributed#8898 ) Patrick Hoefler
Note
Versions 2024.11.0 and 2024.11.1 included a critical performance regression and should be skipped by every user.
This release deprecates the legacy Dask DataFrame implementation. The old implementation will be removed completely in a future release. Users are encourage to switch to the new implementation now and to report any issues they are facing.
Users are also encourage to check that they are only importing functions from dask.dataframe and not any of the submodules.
Dask Array added new quantile and nanquantile methods. Previously, Dask dispatched to the NumPy implementation, which blocked the GIL a lot. This caused large slowdowns on workers with more than one tread and could lead to runtimes over 200s per chunk.
The new quantile implementation avoids many of these problems and reduces runtime to around 1s per chunk independently of the number of threads.
Using Xarrays rolling(...).construct(...) with Dask Arrays led to very large chunksizes that rarely fit into memory on a single worker.
The underlying operations is a view on the smaller NumPy array, but triggering a copy of the data will lead to very large memory usage.
import xarray as xr
import dask.array as da
arr = xr.DataArray(
da.ones((93504, 721, 1440), chunks=("auto", -1, -1)),
dims=["time", "lat", "longitude"],
) # Initial chunks are ~128 MiB
arr.rolling(time=30).construct("window_dim")
map_overlap now creates smaller and more efficient graphs to keep task graphs generally a lot smaller.
The previous version injected a lot of tasks that weren't necessary, increasing the number of tasks by a factor of 2-10x of what actually necessary. This caused a lot of stress on the scheduler.
Einstein summation historically led to very large chunksizes if applied to more than one Dask Array. This behavior is inherited from NumPy but led to out of memory errors on workers:
import dask.array as da
arr = da.random.random((1024, 64, 64, 64, 64), chunks=(256, 16, 16, 16, 16)) # Initial chunks are 128 MiB
result = da.einsum("aijkl,amnop->ijklmnop", arr, arr)
Ensure subgraphs release intermediate results @fjetter
See the Changelog for more information.
Deprecate legacy dask dataframe implementation @phofl
load_stored kwarg to dask.array.store. @dcherian (#11465)slice(None) just maps the keys @phofl (#11450)Task.__repr__() of unpickled object @pentschev (#11463)See the Changelog for more information.
Ensure broadcast_shapes() returns integers, not NumPy scalars. @trexfeathers
broadcast_shapes() returns integers, not NumPy scalars. @trexfeathers (#11434)See the Changelog for more information.
Zarr-Python 3 compatibility ( dask#11388 )
Avoid exponentially increasing taskgraph in overlap ( dask#11423 )
Ensure numba tokenization does not use slow pickle path ( dask#11419 )
Additional changes
Ensure broadcast_shapes() returns integers, not NumPy scalars. ( dask#11434 ) Martin Yeo
(fix): sparse indexing ( dask#11430 ) Ilan Gold
Ensure that recursively calling tokenize respects ensure_deterministic ( dask#11431 ) Florian Jetter
Make P2P more configurable ( distributed#8469 ) Hendrik Makait
Fit Dashboard worker table to page width ( distributed#8897 ) Jacob Tomlinson
Raise helpful error when using the wrong plugin base classes ( distributed#8893 ) Jacob Tomlinson
Fix url escaping on exceptions dashboard for non-string keys ( distributed#8891 ) Patrick Hoefler
Add meaningful error for out of disk exception during write ( distributed#8886 ) Hendrik Makait
Fix binary operations with scalar on the left ( dask-expr#1150 ) Patrick Hoefler
Raise exception when calculating divisons ( dask-expr#1149 ) Patrick Hoefler
Fix merge_asof for single partition ( dask-expr#1145 ) Patrick Hoefler
Improve handling of optional dependencies in analyze and explain ( dask-expr#1146 ) Hendrik Makait
Fix alignment issue with groupby index accessors ( dask-expr#1142 ) Patrick Hoefler
Fix displaying timestamp scalar ( dask-expr#1141 ) Patrick Hoefler
Zarr-Python 3 compatibility (11388)
Avoid exponentially increasing taskgraph in overlap (11423)
Ensure numba tokenization does not use slow pickle path (11419)
Improve error message for incorrect columns order in meta information @dbalabka
RAPIDS_VER to 24.12 @github-actions (#11407)zarr.open\_array instead of using the zarr.Array constructor @jhamman (#11387)See the Changelog for more information.
Adaptive scaling clusters now recover from spurious errors during scaling.
See distributed#8871 by Hendrik Makait for more details. Additional changes
Improve error message for incorrect columns order in meta information ( dask#11393 ) Dmitry Balabka
Update gpuCI RAPIDS_VER to 24.12 ( dask#11407 )
Bump jacobtomlinson/gha-anaconda-package-version from 0.1.3 to 0.1.4 ( dask#11405 )
Switch to using zarr.open_array instead of using the zarr.Array constructor ( dask#11387 ) Joe Hamman
Update gpuCI RAPIDS_VER to 24.12 ( distributed#8879 )
Don’t consider scheduler idle while executing Scheduler.update_graph ( distributed#8877 ) Hendrik Makait
Bump jacobtomlinson/gha-anaconda-package-version from 0.1.3 to 0.1.4 ( distributed#8878 )
Support P2P rechunking datetime arrays ( distributed#8875 ) James Bourbeau
Adaptive scaling clusters now recover from spurious errors during scaling.
See 8871 by `Hendrik Makait`_ for more details.
Revert "Improve normalize\_chunks calculation for "auto" setting" @jrbourbeau
bokeh minimum version to 3.1.0 @jrbourbeau (#11375)tokenize to dedicated submodule @fjetter (#11371)np.min\_scalar\_type in shuffle @jrbourbeau (#11369)See the Changelog for more information.
bokeh>=3.1.0 is now required for diagnostics and the distributed cluster dashboard.
See dask#11375 and distributed#8861 by James Bourbeau for more details.
Add a Task class to replace tuples for task specification.
See dask#11248 by Florian Jetter for more details. Additional changes
Bump peter-evans/create-pull-request from 6 to 7 ( dask#11380 )
Reduce overhead in tokenize ( dask#11373 ) Florian Jetter
Move tokenize to dedicated submodule ( dask#11371 ) Florian Jetter
Ensure process_runnables is not too eager in the presence of multiple splits ( dask#11367 ) Florian Jetter
Use np.min_scalar_type in shuffle ( dask#11369 ) James Bourbeau
Write indexing arrays into dask graph to reduce size for multiple xarray variables ( dask#11362 ) Patrick Hoefler
Cast indexer to minimal dtype in shuffle ( dask#11364 ) Patrick Hoefler
Reduce memory usage of dask.order ( dask#11361 ) Florian Jetter
Bump JamesIves/github-pages-deploy-action from 4.6.3 to 4.6.4 ( dask#11366 )
precommit autoupdate ( dask#11360 ) Florian Jetter
Homogeneously schedule P2P’s unpack tasks ( distributed#8873 ) Hendrik Makait
Work/fix firewall for localhost ( distributed#8868 ) Mario Linker
Use new tokenize module ( distributed#8858 ) James Bourbeau
Point to user code with idempotent plugin warning ( distributed#8856 ) James Bourbeau
Fix test nanny timeout ( distributed#8847 ) Florian Jetter
Bump JamesIves/github-pages-deploy-action from 4.5.0 to 4.6.4 ( distributed#8853 )
Speed up Client.map by computing token only once for func and kwargs ( distributed#8855 ) Florian Jetter
Update pre-commit ( distributed#8852 ) Florian Jetter
bokeh>=3.1.0 is now required for diagnostics and the distributed cluster dashboard.
See 11375 and 8861 by `James Bourbeau`_ for more details.
Add a Task class to replace tuples for task specification.
See 11248 by `Florian Jetter`_ for more details.
Add changelor entries for shuffle, vindex and blockwise\_reshape @phofl
normalize\_chunks @Illviljan (#11271)numpy and pyarrow versions in install docs @jrbourbeau (#11340)numpy>=1.24 and pyarrow>=14.0.1 minimum versions @jrbourbeau (#11331)crick back to Python 3.11+ CI builds @jrbourbeau (#11335)dask.array.fft mismatch with Numpy's interface (add support for norm argument) @joanrue (#10665)rechunk_p2p @hendrikmakait (#11319)axes are positive / add tests for negative axes @joanrue (#10812)See the Changelog for more information.
To enable users to rechunk data at larger scales than before, Dask now automatically chooses an appropriate rechunking method when rechunking on a cluster. This requires no additional configuration and is enabled by default.
Specifically, Dask chooses between task-based and P2P rechunking. While task-based rechunking has been the previous default, P2P rechunking is beneficial when rechunking requires almost all-to-all communication between the old and new chunks, e.g., when changing between spacial and temporal chunking. In these cases, P2P rechunking offers constant memory usage and creates smaller task graphs. As a result, it works for cases where tasks-based rechunking would have previously failed.
To disable automatic selection, users can select their preferred method via the configuration
import dask.config # Choose either "tasks" or "p2p" dask . config . set ({ "array.rechunk.method" : "tasks" })
or when rechunking
import dask.array as da arr = da . random . random ( size = ( 1000 , 1000 , 365 ), chunks = ( - 1 , - 1 , "auto" )) # Choose either "tasks" or "p2p" arr = arr . rechunk (( "auto" , "auto" , - 1 ), method = "tasks" )
See dask#11337 by Hendrik Makait for more details.
Dask added a shuffle-API to Dask Arrays. This API allows for shuffling the data along a single dimension. It will ensure that every group of elements along this dimension are in exactly one chunk. This is a very useful operation for GroupBy-Map patterns in Xarray. See shuffle() for more information and API signature.
See dask#11267 , dask#11311 and dask#11326 by Patrick Hoefler for more details.
The new blockwise_reshape() enables an embarassingly parallel reshaping operation for cases where you don’t care about the order of the underlying array. It is embarassingly parallel and doesn’t trigger a rechunking operation under the hood anymore. This is useful when you don’t care about the order of the resulting Array, i.e. if a reduction is applied to the array or if the reshaping is only temporary.
arr = da . random . random ( size = ( 100 , 100 , 48_000 ), chunks = ( 1000 , 100 , 83 ) result = reshape_blockwise ( arr , ( 10_000 , 48_000 )) result . sum () # or: do something that preserves the shape of each chunk result = reshape_blockwise ( result , ( 100 , 100 , 48_000 ), chunks = arr . chunks )
Dask will automatically calculate the resulting chunks if the number of dimensions is reduced, but you have to specify the resulting chunks if the number of dimensions is increased.
Reshaping a Dask Array oftentimes creates a very complicated computations with rechunk operations in between because Dask respect the C ordering of the Array by default. This ensures that the resulting Dask Array is returned in the same order as the corresponding NumPy Array. However, this can lead to very inefficient computations. The blockwise_reshape is a lot more efficient than the default implemenation if you don’t care about the order.
Warning
Blockwise reshape operations are more efficient as the default, but they will return an Array that is ordered differently. Use with care!
See dask#11328 by Patrick Hoefler for more details.
Indexing a Dask Array with vindex() previously created a single output chunk along the dimensions that were indexed. vindex is commonly used in Xarray when indexing multiple dimensions in a single step, i.e.:
arr = xr . DataArray ( da . random . random (( 100 , 100 , 100 ), chunks = ( 5 , 5 , 50 )), dims = [ 'a' , "b" , "c" ], )
Previously, this put the indexed dimensions into a single chunk:
Dask now uses an improved algorithm that ensures that the chunksizes are kept consistent:
See dask#11330 by Patrick Hoefler for more details. Additional changes
Add changelog entries for shuffle, vindex and blockwise_reshape ( dask#11350 ) Patrick Hoefler
Ensure persisted collections are released without GC ( dask#11348 ) Florian Jetter
Update zoom link for dask meeting ( dask#11357 ) Sarah Charlotte Johnson
Add more docstring examples for normalize_chunks ( dask#11271 ) Illviljan
Choose automatically between tasks-based and p2p rechunking ( dask#11337 ) Hendrik Makait
Implement blockwise reshaping API for arrays ( dask#11328 ) Patrick Hoefler
Make rechunking in shuffle more intelligent to distribute unevenly if necessary ( dask#11326 ) Patrick Hoefler
Increase visibility of GPU CI updates ( dask#11345 ) Charles Blackmon-Luca
Update numpy and pyarrow versions in install docs ( dask#11340 ) James Bourbeau
Fixup dask and distributed dependencies ( dask#11338 ) Patrick Hoefler
Bump numpy>=1.24 and pyarrow>=14.0.1 minimum versions ( dask#11331 ) James Bourbeau
Add crick back to Python 3.11+ CI builds ( dask#11335 ) James Bourbeau
Preserve chunksizes in vindex ( dask#11330 ) Patrick Hoefler
Fix dask.array.fft mismatch with Numpy’s interface (add support for norm argument) ( dask#10665 ) joanrue
Pass additional parameters to rechunk_p2p ( dask#11319 ) Hendrik Makait
Fix docstring formatting for map_overlap ( dask#11332 ) Tao Xin
Fix NumPy overflowing for prod on 2.0 ( dask#11327 ) Patrick Hoefler
Ensure axes are positive / add tests for negative axes ( dask#10812 ) joanrue
Fix map_overlap with new_axis ( dask#11128 ) David Stansby
Avoid capturing code of xdist ( distributed#8846 ) Florian Jetter
Reduce memory footprint of culling P2P rechunking ( distributed#8845 ) Hendrik Makait
Add tests for choosing default rechunking method ( distributed#8843 ) Hendrik Makait
Increase visibility of GPU CI updates ( distributed#8841 ) Charles Blackmon-Luca
Bump test_pause_while_idle timeout ( distributed#8844 ) Florian Jetter
Concatenate small input chunks before P2P rechunking ( distributed#8832 ) Hendrik Makait
Remove dump cluster from gen_cluster ( distributed#8823 ) Florian Jetter
Bump numpy>=1.24 and pyarrow>=14.0.1 minimum versions ( distributed#8837 ) James Bourbeau
Fix PipInstall plugin on Worker ( distributed#8839 ) Hendrik Makait
Remove more Python 3.10 compatibility code ( distributed#8824 ) James Bourbeau
Use task-based rechunking to prechunk along partial boundaries ( distributed#8831 ) Hendrik Makait
Ensure client_desires_keys does not corrupt Scheduler state ( distributed#8827 ) Florian Jetter
Bump minimum cloudpickle to 3 ( distributed#8836 ) James Bourbeau
To enable users to rechunk data at larger scales than before, Dask now automatically chooses an appropriate rechunking method when rechunking on a cluster. This requires no additional configuration and is enabled by default.
Specifically, Dask chooses between task-based and P2P rechunking. While task-based rechunking has been the previous default, P2P rechunking is beneficial when rechunking requires almost all-to-all communication between the old and new chunks, e.g., when changing between spacial and temporal chunking. In these cases, P2P rechunking offers constant memory usage and creates smaller task graphs. As a result, it works for cases where tasks-based rechunking would have previously failed.
To disable automatic selection, users can select their preferred method via the configuration
import dask.config
# Choose either "tasks" or "p2p"
dask.config.set({"array.rechunk.method": "tasks"})
or when rechunking
import dask.array as da
arr = da.random.random(size=(1000, 1000, 365), chunks=(-1, -1, "auto"))
# Choose either "tasks" or "p2p"
arr = arr.rechunk(("auto", "auto", -1), method="tasks")
See 11337 by `Hendrik Makait`_ for more details.
Dask added a shuffle-API to Dask Arrays. This API allows for shuffling the data along a single dimension. It will ensure that every group of elements along this dimension are in exactly one chunk. This is a very useful operation for GroupBy-Map patterns in Xarray. See ~dask.array.Array.shuffle for more information and API signature.
See 11267, 11311 and 11326 by `Patrick Hoefler`_ for more details.
The new ~dask.array.blockwise_reshape enables an embarassingly parallel reshaping operation for cases where you don't care about the order of the underlying array. It is embarassingly parallel and doesn't trigger a rechunking operation under the hood anymore. This is useful when you don't care about the order of the resulting Array, i.e. if a reduction is applied to the array or if the reshaping is only temporary.
arr = da.random.random(size=(100, 100, 48_000), chunks=(1000, 100, 83)
result = reshape_blockwise(arr, (10_000, 48_000))
result.sum()
# or: do something that preserves the shape of each chunk
result = reshape_blockwise(result, (100, 100, 48_000), chunks=arr.chunks)
Dask will automatically calculate the resulting chunks if the number of dimensions is reduced, but you have to specify the resulting chunks if the number of dimensions is increased.
Reshaping a Dask Array oftentimes creates a very complicated computations with rechunk operations in between because Dask respect the C ordering of the Array by default. This ensures that the resulting Dask Array is returned in the same order as the corresponding NumPy Array. However, this can lead to very inefficient computations. The blockwise_reshape is a lot more efficient than the default implemenation if you don't care about the order.
Warning
Blockwise reshape operations are more efficient as the default, but they will return an Array that is ordered differently. Use with care!
See 11328 by `Patrick Hoefler`_ for more details.
Indexing a Dask Array with ~dask.array.vindex previously created a single output chunk along the dimensions that were indexed. vindex is commonly used in Xarray when indexing multiple dimensions in a single step, i.e.:
arr = xr.DataArray(
da.random.random((100, 100, 100), chunks=(5, 5, 50)),
dims=['a', "b", "c"],
)
Previously, this put the indexed dimensions into a single chunk:
Dask now uses an improved algorithm that ensures that the chunksizes are kept consistent:
See 11330 by `Patrick Hoefler`_ for more details.
Ensure pickle does not change tokens @fjetter
asarray for array input with dtype @lucascolley (#11288)np dtypes in dask.array namespace @lucascolley (#11178)See the Changelog for more information.
Reshaping a Dask Array oftentimes squashed the dimensions to reshape into a single chunk. This caused very large output chunks and subsequently a lot of out of memory errors and performance issues.
arr = da . ones ( shape = ( 1000 , 100 , 48_000 ), chunks = ( 1000 , 100 , 83 )) arr . reshape ( 1000 , 100 , 4 , 12_000 )
Previously, this put the last dimension into a single chunk of size 12_000.
The new algorithm will ensure that the chunk-size between in- and output is kept the same. This will avoid large increases in chunk-size and fragmentation of chunks.
The scheduler previously created an inefficient execution graph for Xarray GroupBy-Reduction patterns that use the cohorts strategy:
import xarray as xr arr = xr . open_zarr ( ... ) arr . chunk ( time = TimeResampler ( "ME" )) . groupby ( "time.month" ) . mean ()
An issue in the algorithm that creates the execution order of the task graph lead to an inefficient execution strategy that accumulates a lot of unnecessary memory on the cluster. The improvement is very similar to the previous ordering improvement in 2024.08.0 .
This release drops support for Python 3.9 in accordance with NEP 29. Python 3.10 is now the required minimum version to run Dask.
See dask#11245 and distributed#8793 by Patrick Hoefler for more details. Additional changes
Ensure pickle does not change tokens ( dask#11320 ) Florian Jetter
Add changelog entry for reshape and ordering improvements ( dask#11324 ) Patrick Hoefler
Rename chunksize-tolerance option ( dask#11317 ) Patrick Hoefler
Upgrade gpuCI and fix Dask Array failures with “cupy” backend ( dask#11309 ) Richard (Rick) Zamora
Implement automatic rechunking for shuffle ( dask#11311 ) Patrick Hoefler
Ensure we test against numpy 2 in CI ( dask#11182 ) James Bourbeau
Revert “Test ordering on distributed scheduler ( dask#11310 )” ( dask#11321 ) Florian Jetter
Test ordering on distributed scheduler ( dask#11310 ) Florian Jetter
Add tests to cover more cases of new reshape implementation ( dask#11313 ) Patrick Hoefler
Order: Choose better target for branches with multiple leaf nodes ( dask#11303 ) Patrick Hoefler
Order: Ensure runnable tasks are certainly runnable ( dask#11305 ) Florian Jetter
Fix upstream numpy build ( dask#11304 ) Patrick Hoefler
Make shuffle a no-op if possible ( dask#11291 ) Patrick Hoefler
Keep chunksize consistent in reshape ( dask#11273 ) Patrick Hoefler
Enable slicing with only one unknown chunk ( dask#11301 ) Patrick Hoefler
Link to dask vs spark benchmarks on Dask docs ( dask#11289 ) Sarah Charlotte Johnson
Fix slicing for masked arrays ( dask#11300 ) Patrick Hoefler
Array: fix asarray for array input with dtype ( dask#11288 ) Lucas Colley
Add numpy constants to array api ( dask#11287 ) Lucas Colley
Ignore typing of return value ( dask#11286 ) Patrick Hoefler
Remove automatic resizing in reshape ( dask#11269 ) Patrick Hoefler
API: expose np dtypes in dask.array namespace ( dask#11178 ) Lucas Colley
Reduce frequency of unmanaged memory use warning ( distributed#8834 ) Patrick Hoefler
Update gpuCI RAPIDS_VER to 24.10 ( distributed#8786 )
Avoid RuntimeError: dictionary changed size during iteration in Server._shift_counters() ( distributed#8828 ) Hendrik Makait
Improve concurrent close for scheduler ( distributed#8829 ) Hendrik Makait
MINOR: Extract truncation logic out of partial concatenation in P2P rechunking ( distributed#8826 ) Hendrik Makait
avoid excessive attribute access overhead for remove_from_task_prefix_count ( distributed#8821 ) Florian Jetter
Avoid key validation if validation is disabled ( distributed#8822 ) Florian Jetter
Log worker_client event ( distributed#8819 ) James Bourbeau
Reshaping a Dask Array oftentimes squashed the dimensions to reshape into a single chunk. This caused very large output chunks and subsequently a lot of out of memory errors and performance issues.
arr = da.ones(shape=(1000, 100, 48_000), chunks=(1000, 100, 83))
arr.reshape(1000, 100, 4, 12_000)
Previously, this put the last dimension into a single chunk of size 12_000.
The new algorithm will ensure that the chunk-size between in- and output is kept the same. This will avoid large increases in chunk-size and fragmentation of chunks.
The scheduler previously created an inefficient execution graph for Xarray GroupBy-Reduction patterns that use the cohorts strategy:
import xarray as xr
arr = xr.open_zarr(...)
arr.chunk(time=TimeResampler("ME")).groupby("time.month").mean()
An issue in the algorithm that creates the execution order of the task graph lead to an inefficient execution strategy that accumulates a lot of unnecessary memory on the cluster. The improvement is very similar to the previous ordering improvement in 2024.08.0.
This release drops support for Python 3.9 in accordance with NEP 29. Python 3.10 is now the required minimum version to run Dask.
See 11245 and 8793 by `Patrick Hoefler`_ for more details.
Add changelog for dask order patch @phofl
See the Changelog for more information.
Performance improvement for slicing a Dask Array with a positional indexer. Random access patterns are now more stable and produce easier-to-use results.
x [ slice ( None ), [ 1 , 1 , 3 , 6 , 3 , 4 , 5 ]]
Using a positional indexer was previously prone to drastically increasing the number of output chunks and generating a very large task graph. This has been fixed with a more efficient algorithm.
The new algorithm will keep the chunk-sizes along the axis that is indexed the same to avoid fragmentation of chunks or a large increase in chunk-size.
See dask#11262 and dask#11267 by Patrick Hoefler for more details and performance benchmarks.
The scheduler previously created an inefficient execution graph for Xarray GroupBy-Reduction patterns like:
import xarray as xr arr = xr . open_zarr ( ... ) arr . groupby ( "time.month" ) . mean ()
An issue in the algorithm that creates the execution order of the task graph lead to an inefficient execution strategy that accumulates a lot of unneceessary memory on the cluster.
The operation itself is embarassingly parallel. Using the proper execution strategy the scheduler can now execute the operation with constant memory, avoiding spilling and allowing us to scale to larger datasets.
See distributed#8818 by Patrick Hoefler for more details and examples. Additional changes
Add changelog for dask order patch ( dask#11278 ) Patrick Hoefler
Add regression test for xarray map reduce ( dask#11277 ) Florian Jetter
Add changelog entry for take ( dask#11274 ) Patrick Hoefler
Revert “order: remove data task graph normalization” ( dask#11276 ) Patrick Hoefler
Use the shuffle algorithm for take ( dask#11267 ) Patrick Hoefler
Implement task-based array shuffle ( dask#11262 ) Patrick Hoefler
Remove data task graph normalization ( dask#11263 ) Florian Jetter
Update zoom link for monthly meeting ( dask#11265 ) Sarah Charlotte Johnson
Update data loading section of best practices ( dask#11247 ) Patrick Hoefler
Match default chunksize in docstring to actual default set in code ( dask#11254 ) Bernhard Raml
Fixup casting error in pandas 3 ( dask#11250 ) Patrick Hoefler
Skip new warning from pandas ( dask#11249 ) Patrick Hoefler
Fix pandas nightly bugs ( dask#11244 ) Patrick Hoefler
Run graph normalisation after dask order ( distributed#8818 ) Patrick Hoefler
Update large graph size warning to remove scatter recommendation ( distributed#8815 ) Patrick Hoefler
Fail tasks exceeding no-workers-timeout ( distributed#8806 ) Hendrik Makait
Fix exception handling for NannyPlugin.setup and NannyPlugin.teardown ( distributed#8811 ) Hendrik Makait
Fix exception handling for WorkerPlugin.setup and WorkerPlugin.teardown ( distributed#8810 ) Hendrik Makait
typo fix ( distributed#8812 ) alex-rakowski
Fix if / else for send_recv_from_rpc ( distributed#8809 ) Patrick Hoefler
Ensure that adaptive only stops once ( distributed#8807 ) Hendrik Makait
Reduce noise from GC-related logging ( distributed#8804 ) Hendrik Makait
Remove unused delete_interval and synchronize_worker_interval from Scheduler ( distributed#8801 ) Hendrik Makait
Change log level for Compute Failed log message ( distributed#8802 ) Patrick Hoefler
Add Prometheus metric for time spent on GC ( distributed#8803 ) Hendrik Makait
Add Prometheus metrics for dask_worker_{added|removed}_total ( distributed#8798 ) Hendrik Makait
Add log event for worker-ttl-timed-out ( distributed#8800 ) Hendrik Makait
Add Prometheus metrics for dask_client_connections_{added|removed}_total ( distributed#8799 ) Hendrik Makait
Fix PackageInstall plugin ( distributed#8794 ) Hendrik Makait
Make stealing more robust ( distributed#8788 ) Hendrik Makait
Leave a warning about future instantiation ( distributed#8782 ) Florian Jetter
Performance improvement for slicing a Dask Array with a positional indexer. Random access patterns are now more stable and produce easier-to-use results.
x[slice(None), [1, 1, 3, 6, 3, 4, 5]]
Using a positional indexer was previously prone to drastically increasing the number of output chunks and generating a very large task graph. This has been fixed with a more efficient algorithm.
The new algorithm will keep the chunk-sizes along the axis that is indexed the same to avoid fragmentation of chunks or a large increase in chunk-size.
See 11262 and 11267 by `Patrick Hoefler`_ for more details and performance benchmarks.
The scheduler previously created an inefficient execution graph for Xarray GroupBy-Reduction patterns like:
import xarray as xr
arr = xr.open_zarr(...)
arr.groupby("time.month").mean()
An issue in the algorithm that creates the execution order of the task graph lead to an inefficient execution strategy that accumulates a lot of unneceessary memory on the cluster.
The operation itself is embarassingly parallel. Using the proper execution strategy the scheduler can now execute the operation with constant memory, avoiding spilling and allowing us to scale to larger datasets.
See 8818 by `Patrick Hoefler`_ for more details and examples.
Fixes for d freq deprecation in pandas=3 @jrbourbeau
d freq deprecation in pandas=3 @jrbourbeau (#11228)See the Changelog for more information.
distributed.Lock is now resilient to worker failures. Previously deadlocks were possible in cases where a lock-holding worker was lost and/or failed to release the lock due to an error.
See distributed#8770 by Florian Jetter for more details. Additional changes
Remove and warn of persist usage ( dask#11237 ) Patrick Hoefler
Preserve timestamp unit during meta creation ( dask#11233 ) Patrick Hoefler
Ensure that dask-expr DataFrames are optimized when put into delayed ( dask#11231 ) Patrick Hoefler
Fixes for d freq deprecation in pandas=3 ( dask#11228 ) James Bourbeau
bump approx threshold for test_quantile ( dask#10720 ) Florian Jetter
Bump xarray-contrib/issue-from-pytest-log from 1.2.8 to 1.3.0 ( dask#11221 )
Bump JamesIves/github-pages-deploy-action from 4.6.1 to 4.6.3 ( dask#11222 )
Ensure Lock always register with scheduler ( distributed#8781 ) Florian Jetter
Temporarily pin setuptools < 71 ( distributed#8785 ) James Bourbeau
Restore len() on TaskPrefix ( distributed#8783 ) Hendrik Makait
Avoid false positives for p2p-failed log event ( distributed#8777 ) Hendrik Makait
Expose paused and retired workers separately in prometheus ( distributed#8613 ) Patrick Hoefler
Creating transitions-failures log event ( distributed#8776 ) alex-rakowski
Implement HLG layer for P2P rechunking ( distributed#8751 ) Hendrik Makait
Add another test for a possible deadlock scenario caused by ( distributed#8703 ) ( distributed#8769 ) Hendrik Makait
Raise an error if compute on persisted collection with released futures ( distributed#8764 ) Florian Jetter
Re-raise P2PConsistencyError from failed P2P tasks ( distributed#8748 ) Hendrik Makait
Robuster faster tests memory sampler ( distributed#8758 ) Florian Jetter
Fix scheduler_bokeh::test_shuffling ( distributed#8766 ) Florian Jetter
Increase timeouts for pubsub::test_client_worker ( distributed#8765 ) Florian Jetter
Factor out async taskgroup ( distributed#8756 ) Florian Jetter
Don’t sort keys lexicographically in worker table ( distributed#8753 ) Florian Jetter
Use functools.cache instead of functools.lru_cache for extremely often called functions ( distributed#8762 ) Jonas Dedden
Robuster deeply nested structures ( distributed#8730 ) Florian Jetter
Adding HLG to MAP ( distributed#8740 ) alex-rakowski
Add close worker button to worker info page ( distributed#8742 ) James Bourbeau
distributed.Lock is now resilient to worker failures. Previously deadlocks were possible in cases where a lock-holding worker was lost and/or failed to release the lock due to an error.
See 8770 by `Florian Jetter`_ for more details.
Only count data that is in memory for xarray sizeof @fjetter
See the Changelog for more information.
This release drops support for pandas<2 . pandas 2.0 is now the required minimum version to run Dask DataFrame.
The mimimum version of partd was also raised to 1.4.0. Versions before 1.4 are not compatible with pandas 2.
See dask#11199 by Patrick Hoefler for more details.
distributed.Pub and distributed.Sub have been deprecated and will be removed in a future release. Please switch to distributed.Client.log_event() and distributed.Worker.log_event() instead.
See distributed#8724 by Hendrik Makait for more details. Additional changes
Only count data that is in memory for xarray sizeof ( dask#11206 ) Florian Jetter
Fix botocore re-raising error ( dask#11209 ) Patrick Hoefler
Update Coiled links in documentation ( dask#11211 ) Sarah Charlotte Johnson
Add some array-expr methods ( dask#11210 ) Patrick Hoefler
Fix quantile for arrow dtypes ( dask#11202 ) Patrick Hoefler
Add utility to verify optional dependencies ( dask#11205 ) Patrick Hoefler
Implement array expression switch ( dask#11203 ) Patrick Hoefler
Remove no longer supported ipython reference ( dask#11196 ) Patrick Hoefler
Remove from_delayed references ( dask#11195 ) Patrick Hoefler
Add other IO connectors to docs ( dask#11189 ) Patrick Hoefler
Fix assert_eq import from cudf ( distributed#8747 ) James Bourbeau
Log traceback upon task error ( distributed#8746 ) Hendrik Makait
Update system monitor when polling Prometheus metrics ( distributed#8745 ) Hendrik Makait
Bump pandas to 2.0 in mindeps build ( distributed#8743 ) James Bourbeau
Refactor event logging functionality into broker ( distributed#8731 ) Hendrik Makait
Drop support for pandas 1.X ( distributed#8741 ) Hendrik Makait
Remove is_python_shutting_down ( distributed#8492 ) Hendrik Makait
Fix test_task_state_instance_are_garbage_collected ( distributed#8735 ) Hendrik Makait
Fix floating-point inaccuracy ( distributed#8736 ) Hendrik Makait
Fix pynvml handles ( distributed#8693 ) Benjamin Zaitlen
get_ip : handle getting 0.0.0.0 ( distributed#8712 ) Adam Williamson
Remove FutureWarning in test_task_state_instance_are_garbage_collected ( distributed#8734 ) Hendrik Makait
Fix mindeps -testing on CI ( distributed#8728 ) Hendrik Makait
Extract tests related to event-logging into separate file ( distributed#8733 ) Hendrik Makait
Use safer context for ProcessPoolExecutor ( distributed#8715 ) Elliott Sales de Andrade
Cache URL encoding of worker addresses in dashboard ( distributed#8725 ) Florian Jetter
More robust bokeh test_shuffling ( distributed#8727 ) Florian Jetter
Fix type in actor docs ( distributed#8711 ) Sultan Orazbayev
More useful warning if a plugin type is provided instead of instance ( distributed#8689 ) Florian Jetter
Improve error on cancelled tasks due to disconnect ( distributed#8705 ) Hendrik Makait
Fix wait condition on test_forget_errors ( distributed#8714 ) Elliott Sales de Andrade
Skip test_deadlock_dependency_of_queued_released ( distributed#8723 ) Hendrik Makait
Fix test_quiet_client_close ( distributed#8722 ) Hendrik Makait
Fix cleanup iteration in save_sys_modules ( distributed#8713 ) Elliott Sales de Andrade
Add quotes to missing bokeh installation commands ( distributed#8717 ) James Bourbeau
This release drops support for pandas<2. pandas 2.0 is now the required minimum version to run Dask DataFrame.
The mimimum version of partd was also raised to 1.4.0. Versions before 1.4 are not compatible with pandas 2.
See 11199 by `Patrick Hoefler`_ for more details.
distributed.Pub and distributed.Sub have been deprecated and will be removed in a future release. Please switch to distributed.Client.log_event and distributed.Worker.log_event instead.
See 8724 by `Hendrik Makait`_ for more details.
Get docs build passing @jrbourbeau
This is a patch release to update an issue with dask and distributed version pinning in the 2024.6.1 release. Additional changes
Get docs build passing ( dask#11184 ) James Bourbeau
profile._f_lineno : handle next_line being None in Python 3.13 ( dask#8710 ) Adam Williamson
This is a patch release to update an issue with dask and distributed version pinning in the 2024.6.1 release.
Cache global query-planning config @rjzamora
test\_map\_freq\_to\_period\_start for pandas=3 @jrbourbeau (#11181)See the Changelog for more information.
This release includes a critical fix that fixes a deadlock that can arise when dependencies of root-ish tasks are rescheduled, e.g. due to a worker being lost.
See distributed#8703 by Hendrik Makait for more details. Additional changes
Cache global query-planning config ( dask#11183 ) Richard (Rick) Zamora
Python 3.13 fixes ( dask#11185 ) Adam Williamson
Fix test_map_freq_to_period_start for pandas=3 ( dask#11181 ) James Bourbeau
Bump release-drafter/release-drafter from 5 to 6 ( distributed#8699 )
This release includes a critical fix that fixes a deadlock that can arise when dependencies of root-ish tasks are rescheduled, e.g. due to a worker being lost.
See 8703 by `Hendrik Makait`_ for more details.
Remove deprecated dask.compatibility module @jrbourbeau
test\_dt\_accessor with query planning disabled @jrbourbeau (#11177)packaging.version.Version @jrbourbeau (#11171)dask.compatibility module @jrbourbeau (#11172)xarray.NamedArray @hendrikmakait (#11168)See the Changelog for more information.
Tokenizing memmap arrays will now avoid materializing the array into memory.
See dask#11161 by Florian Jetter for more details. Additional changes
Fix test_dt_accessor with query planning disabled ( dask#11177 ) James Bourbeau
Use packaging.version.Version ( dask#11171 ) James Bourbeau
Remove deprecated dask.compatibility module ( dask#11172 ) James Bourbeau
Ensure compatibility for xarray.NamedArray ( dask#11168 ) Hendrik Makait
Estimate sizes of xarray collections ( dask#11166 ) Florian Jetter
Add section about futures and variables ( dask#11164 ) Florian Jetter
Update docs for combined Dask community meeting info ( dask#11159 ) Sarah Charlotte Johnson
Avoid rounding error in test_prometheus_collect_count_total_by_cost_multipliers ( distributed#8687 ) Hendrik Makait
Log key collision count in update_graph log event ( distributed#8692 ) Hendrik Makait
Automate GitHub Releases when new tags are pushed ( distributed#8626 ) Jacob Tomlinson
Fix log event with multiple topics ( distributed#8691 ) Hendrik Makait
Rename safe to expected in Scheduler.remove_worker ( distributed#8686 ) Hendrik Makait
Log event during failure ( distributed#8663 ) Hendrik Makait
Eagerly update aggregate statistics for TaskPrefix instead of calculating them on-demand ( distributed#8681 ) Hendrik Makait
Improve graph submission time for P2P rechunking by avoiding unpack recursion into indices ( distributed#8672 ) Florian Jetter
Add safe keyword to remove-worker event ( distributed#8647 ) alex-rakowski
Improved errors and reduced logging for P2P RPC calls ( distributed#8666 ) Hendrik Makait
Adjust P2P tests for dask-expr ( distributed#8662 ) Hendrik Makait
Iterate over copy of Server.digests_total_since_heartbeat to avoid RuntimeError ( distributed#8670 ) Hendrik Makait
Log task state in Compute Failed ( distributed#8668 ) Hendrik Makait
Add Prometheus gauge for task groups ( distributed#8661 ) Hendrik Makait
Fix too strict assertion in shuffle code for pandas subclasses ( distributed#8667 ) Joris Van den Bossche
Reduce noise from erring tasks that are not supposed to be running ( distributed#8664 ) Hendrik Makait
Tokenizing memmap arrays will now avoid materializing the array into memory.
See 11161 by `Florian Jetter`_ for more details.
Fix nightly Zarr installation in CI @jrbourbeau
RAPIDS_VER to 24.08 @github-actions (#11141)test\_groupby\_grouper\_dispatch @rjzamora (#11144)test\_array\_function\_sparse with new sparse release @jrbourbeau (#11139)test\_parse\_dates\_multi\_column on pandas=3 @jrbourbeau (#11132)See the Changelog for more information.
This release primarily contains minor bug fixes. Additional changes
Fix nightly Zarr installation in CI ( dask#11151 ) James Bourbeau
Add python 3.11 build to GPU CI ( dask#11135 ) Charles Blackmon-Luca
Update gpuCI RAPIDS_VER to 24.08 ( dask#11141 )
Update test_groupby_grouper_dispatch ( dask#11144 ) Richard (Rick) Zamora
Bump JamesIves/github-pages-deploy-action from 4.6.0 to 4.6.1 ( dask#11136 )
Unskip test_array_function_sparse with new sparse release ( dask#11139 ) James Bourbeau
Fix test_parse_dates_multi_column on pandas=3 ( dask#11132 ) James Bourbeau
Don’t draft release notes for tagged commits ( dask#11138 ) Jacob Tomlinson
Reduce task group count for partial P2P rechunks ( distributed#8655 ) Hendrik Makait
Update gpuCI RAPIDS_VER to 24.08 ( distributed#8652 )
Submit collections metadata to scheduler ( distributed#8612 ) Florian Jetter
Fix indent in code example in task-launch.rst ( distributed#8650 ) Ray Bell
Avoid multiple WorkerState sphinx error ( distributed#8643 ) James Bourbeau
Minor updates to ML page @jrbourbeau
sparse test on 0.15.2 @jrbourbeau (#11131)pyarrow is installed in upstream CI build @jrbourbeau (#11121)zarr stores in to\_zarr with distributed @GFleishman (#10422)pytest error when skipping NumPy 2.0 tests @jrbourbeau (#11110)h5py in upstream CI build @jrbourbeau (#11108)scikit-image in upstream CI build @jrbourbeau (#11107)meshgrid and atleast\_\*d NumPy 2 updates @jrbourbeau (#11106)percentiles\_summary logic @rjzamora (#11094)See the Changelog for more information.
This release contains compatibility updates for the upcoming NumPy 2.0 release.
See dask#11096 by Benjamin Zaitlen and dask#11106 by James Bourbeau for more details.
This release contains adds support for MutableMapping -backed Zarr stores like zarr.storage.DirectoryStore , etc.
See dask#10422 by Greg M. Fleishman for more details. Additional changes
Minor updates to ML page ( dask#11129 ) James Bourbeau
Skip failing sparse test on 0.15.2 ( dask#11131 ) James Bourbeau
Make sure nightly pyarrow is installed in upstream CI build ( dask#11121 ) James Bourbeau
Add initial draft of ML overview document ( dask#11114 ) Matthew Rocklin
Test query-planning in gpuCI ( dask#11060 ) Richard (Rick) Zamora
Avoid pytest error when skipping NumPy 2.0 tests ( dask#11110 ) James Bourbeau
Use nightly h5py in upstream CI build ( dask#11108 ) James Bourbeau
Use nightly scikit-image in upstream CI build ( dask#11107 ) James Bourbeau
Bump actions/checkout from 4.1.4 to 4.1.5 ( dask#11105 )
Enable parquet append tests after fix ( dask#11104 ) Patrick Hoefler
Skip fastparquet tests for numpy 2 ( dask#11103 ) Patrick Hoefler
Fix misspelling found by codespell ( dask#11097 ) Dimitri Papadopoulos Orfanos
Fix doc build ( dask#11099 ) Patrick Hoefler
Clean up percentiles_summary logic ( dask#11094 ) Richard (Rick) Zamora
Apply ruff/flake8-implicit-str-concat rule ISC001 ( dask#11098 ) Dimitri Papadopoulos Orfanos
Fix clocks on Windows with Python 3.13 ( distributed#8642 ) Victor Stinner
Fix “Print host info” CI step on Mac OS (arm64) ( distributed#8638 ) Hendrik Makait
This release contains compatibility updates for the upcoming NumPy 2.0 release.
See 11096 by `Benjamin Zaitlen`_ and 11106 by `James Bourbeau`_ for more details.
This release contains adds support for MutableMapping-backed Zarr stores like zarr.storage.DirectoryStore, etc.
See 10422 by `Greg M. Fleishman`_ for more details.
DOC: intersphinx, don't link to click dev version. @Carreau
melt support when query-planning is enabled @rjzamora (#11088)clear\_known\_categories utility @rjzamora (#11059)RAPIDS_VER to 24.06, disable query planning @charlesbluca (#11045)See the Changelog for more information.
This release primarily contains minor bugfixes. Additional changes
Don’t link to click intersphinx dev version ( dask#11091 ) M Bussonnier
Fix API doc links for some dask-expr expressions ( dask#11092 ) Patrick Hoefler
Add dask-expr to upstream build ( dask#11086 ) Patrick Hoefler
Add melt support when query-planning is enabled ( dask#11088 ) Richard (Rick) Zamora
Skip dataframe/product when in numpy 2 envs ( dask#11089 ) Benjamin Zaitlen
Add plots to illustrate what the optimizer does ( dask#11072 ) Patrick Hoefler
Fixup pandas upstream tests ( dask#11085 ) Patrick Hoefler
Bump conda-incubator/setup-miniconda from 3.0.3 to 3.0.4 ( dask#11084 )
Bump actions/checkout from 4.1.3 to 4.1.4 ( dask#11083 )
Fix CI after pytest changes ( dask#11082 ) Patrick Hoefler
Fixup tests for more efficient dask-expr implementation ( dask#11071 ) Patrick Hoefler
Generalize clear_known_categories utility ( dask#11059 ) Richard (Rick) Zamora
Bump JamesIves/github-pages-deploy-action from 4.5.0 to 4.6.0 ( dask#11062 )
Bump release-drafter/release-drafter from 5 to 6 ( dask#11063 )
Bump actions/checkout from 4.1.2 to 4.1.3 ( dask#11061 )
Update GPU CI RAPIDS_VER to 24.06, disable query planning ( dask#11045 ) Charles Blackmon-Luca
Move tests ( distributed#8631 ) Hendrik Makait
Bump actions/checkout from 4.1.2 to 4.1.3 ( distributed#8628 )
Add GitHub Releases automation @jacobtomlinson
test_set_index for "cudf" backend @rjzamora (#11029)to/from\_legacy\_dataframe instead of to/from\_dask\_dataframe @rjzamora (#11025)See the Changelog for more information.
The Query Optimizer will inspect quires to determine if a merge(...) or groupby(...).apply(...) requires a shuffle. A shuffle can be avoided, if the DataFrame was shuffled on the same columns in a previous step without any operations in between that change the partitioning layout or the relevant values in each partition.
result = df . merge ( df2 , on = "a" ) >>> result = result . merge ( df3 , on = "a" )
The Query optimizer will identify that result was previously shuffled on "a" as well and thus only shuffle df3 in the second merge operation before doing a blockwise merge.
The Query Optimizer will automatically repartition datasets read from Parquet files if individual partitions are too small. This will reduce the number of partitions in consequentially also the size of the task graph.
The Optimizer aims to produce partitions of at least 75MB and will combine multiple files together if necessary to reach this threshold. The value can be configured by using
dask . config . set ({ "dataframe.parquet.minimum-partition-size" : 100_000_000 })
The value is given in bytes. The default threshold is relatively conservative to avoid memory issues on worker nodes with a relatively small amount of memory per thread. Additional changes
Add GitHub Releases automation ( dask#11057 ) Jacob Tomlinson
Add changelog entries for new release ( dask#11058 ) Patrick Hoefler
Reinstate try/except block in _bind_property ( dask#11049 ) Lawrence Mitchell
Fix link for query planning docs ( dask#11054 ) Patrick Hoefler
Add config parameter for parquet file size ( dask#11052 ) Patrick Hoefler
Update percentile docstring ( dask#11053 ) Abel Aoun
Add docs for query optimizer ( dask#11043 ) Patrick Hoefler
Assignment of np.ma.masked to obect-type Array ( dask#9627 ) David Hassell
Don’t error if dask_expr is not installed ( dask#11048 ) Simon Høxbro Hansen
Adjust test_set_index for “cudf” backend ( dask#11029 ) Richard (Rick) Zamora
Use to/from_legacy_dataframe instead of to/from_dask_dataframe ( dask#11025 ) Richard (Rick) Zamora
Tokenize bag groupby keys ( dask#10734 ) Charles Stern
Add lazy “cudf” registration for p2p-related dispatch functions ( dask#11040 ) Richard (Rick) Zamora
Collect memray profiles on exception ( distributed#8625 ) Florian Jetter
Ensure inproc properly emulates serialization protocol ( distributed#8622 ) Florian Jetter
Relax test stats profiling2 ( distributed#8621 ) Florian Jetter
Restart workers when worker-ttl expires ( distributed#8538 ) crusaderky
Use monotonic for deadline test ( distributed#8620 ) Florian Jetter
Fix race condition for published futures with annotations ( distributed#8577 ) Florian Jetter
Scatter by worker instead of worker -> nthreads ( distributed#8590 ) Miles
Send log-event if worker is restarted because of memory pressure ( distributed#8617 ) Patrick Hoefler
Do not print xfailed tests in CI ( distributed#8619 ) Florian Jetter
ensure workers are not downscaled when participating in p2p ( distributed#8610 ) Florian Jetter
Run against stable fsspec ( distributed#8615 ) Florian Jetter
The Query Optimizer will inspect quires to determine if a merge(...) or groupby(...).apply(...) requires a shuffle. A shuffle can be avoided, if the DataFrame was shuffled on the same columns in a previous step without any operations in between that change the partitioning layout or the relevant values in each partition.
>>> result = df.merge(df2, on="a")
>>> result = result.merge(df3, on="a")
The Query optimizer will identify that result was previously shuffled on "a" as well and thus only shuffle df3 in the second merge operation before doing a blockwise merge.
The Query Optimizer will automatically repartition datasets read from Parquet files if individual partitions are too small. This will reduce the number of partitions in consequentially also the size of the task graph.
The Optimizer aims to produce partitions of at least 75MB and will combine multiple files together if necessary to reach this threshold. The value can be configured by using
>>> dask.config.set({"dataframe.parquet.minimum-partition-size": 100_000_000})
The value is given in bytes. The default threshold is relatively conservative to avoid memory issues on worker nodes with a relatively small amount of memory per thread.
This is a minor bugfix release that that fixes an error when importing dask.dataframe with Python 3.11.9.
This is a minor bugfix release that that fixes an error when importing dask.dataframe with Python 3.11.9.
See dask#11035 and dask#11039 from Richard (Rick) Zamora for details. Additional changes
Remove skips for named aggregations ( dask#11036 ) Patrick Hoefler
Don’t deep-copy read-only buffers on unpickle ( distributed#8609 ) crusaderky
Add dask-expr to dask conda recipe ( distributed#8601 ) Charles Blackmon-Luca
This is a minor bugfix release that that fixes an error when importing dask.dataframe with Python 3.11.9.
See 11035 and 11039 from `Richard (Rick) Zamora`_ for details.
- Fixup deprecation warning from pandas ( distributed#8564 ) Patrick Hoefler
This release contains a variety of bugfixes in Dask DataFrame’s new query planner.
GPU memory and utilization dashboard functionality has been restored. Previously these plots were unintentionally left blank.
See distributed#8572 from Benjamin Zaitlen for details. Additional changes
Build nightlies on tag releases ( dask#11014 ) Charles Blackmon-Luca
Remove xfail tracebacks from test suite ( dask#11028 ) Patrick Hoefler
Fix CI for upstream pandas changes ( dask#11027 ) Patrick Hoefler
Fix value_counts raising if branch exists of nans only ( dask#11023 ) Patrick Hoefler
Enable custom expressions in dask_cudf ( dask#11013 ) Richard (Rick) Zamora
Raise ImportError instead of ValueError when dask-expr cannot be imported ( dask#11007 ) James Lamb
Add HypersSpy to ecosystem.rst ( dask#11008 ) Jonas Lähnemann
Add Hugging Face hf:// to the list of fsspec compatible remote services ( dask#11012 ) Quentin Lhoest
Bump actions/checkout from 4.1.1 to 4.1.2 ( dask#11009 )
Refresh documentation for annotations and spans ( distributed#8593 ) crusaderky
Fixup deprecation warning from pandas ( distributed#8564 ) Patrick Hoefler
Add Python 3.11 to GPU CI matrix ( distributed#8598 ) Charles Blackmon-Luca
Deadline to use a monotonic timer ( distributed#8597 ) crusaderky
Update gpuCI RAPIDS_VER to 24.06 ( distributed#8588 )
Refactor restart() and restart_workers() ( distributed#8550 ) crusaderky
Bump actions/checkout from 4.1.1 to 4.1.2 ( distributed#8587 )
Fix bokeh deprecations ( distributed#8594 ) Miles
Fix flaky test: test_shutsdown_cleanly ( distributed#8582 ) Miles
Include type in failed sizeof warning ( distributed#8580 ) James Bourbeau
This release contains a variety of bugfixes in Dask DataFrame's new query planner.
GPU memory and utilization dashboard functionality has been restored. Previously these plots were unintentionally left blank.
See 8572 from `Benjamin Zaitlen`_ for details.
This is a minor release that primarily demotes an exception to a warning if dask-expr is not installed when upgrading. Additional changes
This is a minor release that primarily demotes an exception to a warning if dask-expr is not installed when upgrading. Additional changes
Only warn if dask-expr is not installed ( dask#11003 ) Florian Jetter
Fix typos found by codespell ( dask#10993 ) Dimitri Papadopoulos Orfanos
Extra CI job with dask-expr disabled ( distributed#8583 ) crusaderky
Fix worker dashboard proxy ( distributed#8528 ) Miles
Fix flaky test_restart_waits_for_new_workers ( distributed#8573 ) crusaderky
Fix flaky test_raise_on_incompatible_partitions ( distributed#8571 ) crusaderky
This is a minor release that primarily demotes an exception to a warning if dask-expr is not installed when upgrading.
This release is enabling query planning by default for all users of dask.dataframe .
Released on March 11, 2024
This release is enabling query planning by default for all users of dask.dataframe .
The query planning functionality represents a rewrite of the DataFrame using dask-expr . This is a drop-in replacement and we expect that most users will not have to adjust any of their code. Any feedback can be reported on the Dask issue tracker or on the query planning feedback issue .
If you are encountering any issues you are still able to opt-out by setting
import dask >>> dask . config . set ({ 'dataframe.query-planning' : False })
The new query planning backend is requiring at least pandas 2.0 . This pandas version will automatically be installed if you are installing from conda or if you are installing using dask[complete] or dask[dataframe] from pip.
The legacy DataFrame implementation is still supporting pandas 1.X if you install dask without extras. Additional changes
Update tests for pandas nightlies with dask-expr ( dask#10989 ) Patrick Hoefler
Use dask-expr docs as main reference docs for DataFrames ( dask#10990 ) Patrick Hoefler
Adjust from_array test for dask-expr ( dask#10988 ) Patrick Hoefler
Unskip to_delayed test ( dask#10985 ) Patrick Hoefler
Bump conda-incubator/setup-miniconda from 3.0.1 to 3.0.3 ( dask#10978 )
Fix bug when enabling dask-expr ( dask#10977 ) Patrick Hoefler
Update docs and requirements for dask-expr and remove warning ( dask#10976 ) Patrick Hoefler
Fix numpy 2 compatibility with ogrid usage ( dask#10929 ) David Hoese
Turn on dask-expr switch ( dask#10967 ) Patrick Hoefler
Force initializing the random seed with the same byte order interpret… ( dask#10970 ) Elliott Sales de Andrade
Use correct encoding for line terminator when reading CSV ( dask#10972 ) Elliott Sales de Andrade
perf: do not unnecessarily recalculate input/output indices in _optimize_blockwise ( dask#10966 ) Lindsey Gray
Adjust tests for string option in dask-expr ( dask#10968 ) Patrick Hoefler
Adjust tests for array conversion in dask-expr ( dask#10973 ) Patrick Hoefler
TST: Fix sizeof tests on 32bit ( dask#10971 ) Elliott Sales de Andrade
TST: Add missing skip for pyarrow ( dask#10969 ) Elliott Sales de Andrade
Implement dask-expr conversion for bag.to_dataframe ( dask#10963 ) Patrick Hoefler
Fix dask-expr import errors ( dask#10964 ) Miles
Clean up Sphinx documentation for dask.config ( dask#10959 ) crusaderky
Use stdlib importlib.metadata on Python 3.12+ ( dask#10955 ) wim glenn
Cast partitioning_index to smaller size ( dask#10953 ) Florian Jetter
Reuse dask/dask groupby Aggregation ( dask#10952 ) Patrick Hoefler
ensure tokens on futures are unique ( distributed#8569 ) Florian Jetter
Don’t obfuscate fine performance metrics failures ( distributed#8568 ) crusaderky
Mark shuffle fast tasks in dask-expr ( distributed#8563 ) crusaderky
Weigh gilknocker Prometheus metric by duration ( distributed#8558 ) crusaderky
Fix scheduler transition error on memory->erred ( distributed#8549 ) Hendrik Makait
Make CI happy again ( distributed#8560 ) Miles
Fix flaky test_Future_release_sync ( distributed#8562 ) crusaderky
Fix flaky test_flaky_connect_recover_with_retry ( distributed#8556 ) Hendrik Makait
typing tweaks in scheduler.py ( distributed#8551 ) crusaderky
Bump conda-incubator/setup-miniconda from 3.0.2 to 3.0.3 ( distributed#8553 )
Install dask-expr on CI ( distributed#8552 ) Hendrik Makait
P2P shuffle can drop partition column before writing to disk ( distributed#8531 ) Hendrik Makait
Better logging for worker removal ( distributed#8517 ) crusaderky
Add indicator support to merge ( distributed#8539 ) Patrick Hoefler
Bump conda-incubator/setup-miniconda from 3.0.1 to 3.0.2 ( distributed#8535 )
Avoid iteration error when getting module path ( distributed#8533 ) James Bourbeau
Ignore stdlib threading module in code collection ( distributed#8532 ) James Bourbeau
Fix excessive logging on P2P retry ( distributed#8511 ) Hendrik Makait
Prevent typos in retire_workers parameters ( distributed#8524 ) crusaderky
Cosmetic cleanup of test_steal (backport from #8185) ( distributed#8509 ) crusaderky
Fix flaky test_compute_per_key ( distributed#8521 ) crusaderky
Fix flaky test_no_workers_timeout_queued ( distributed#8523 ) crusaderky
Released on March 11, 2024
This release is enabling query planning by default for all users of dask.dataframe.
The query planning functionality represents a rewrite of the DataFrame using dask-expr. This is a drop-in replacement and we expect that most users will not have to adjust any of their code. Any feedback can be reported on the Dask issue tracker or on the query planning feedback issue.
If you are encountering any issues you are still able to opt-out by setting
>>> import dask
>>> dask.config.set({'dataframe.query-planning': False})
The new query planning backend is requiring at least pandas 2.0. This pandas version will automatically be installed if you are installing from conda or if you are installing using dask[complete] or dask[dataframe] from pip.
The legacy DataFrame implementation is still supporting pandas 1.X if you install dask without extras.
The last release contained a DeprecationWarning that alerts users to an upcoming switch of dask.dafaframe to use the new backend with support for quer…
Released on February 23, 2024
The last release contained a DeprecationWarning that alerts users to an upcoming switch of dask.dafaframe to use the new backend with support for query planning (see also dask#10934 ).
This DeprecationWarning is triggered in import of the dask.dataframe module and the community raised concerns about this being to verbose.
It is now possible to silence this warning
See dask#10936 and dask#10925 from Miles for details.
Blockwise fusion optimization can cause a task key collision that is not being handled properly by the distributed scheduler (see dask#9888 ). Users will typically notice this by seeing one of various internal exceptions that cause a system deadlock or critical failure. While this issue could not be fixed, the scheduler now implements a mechanism that should mitigate most occurences and issues a warning if the issue is detected.
See distributed#8185 from crusaderky and Florian Jetter for details.
Over the course of this, various improvements to tokenization have been implemented. See dask#10913 , dask#10884 , dask#10919 , dask#10896 and primarily dask#10883 from crusaderky for more details.
Adaptive scaling could previously lose data during downscaling if many tasks had to be moved. This typically, but not exclusively, occured on large clusters and would manifest as a recomputation of tasks and could cause clusters to oscillate between up- and downscaling without ever finishing.
See distributed#8522 from crusaderky for more details. Additional changes
Remove flaky fastparquet test ( dask#10948 ) Patrick Hoefler
Enable Aggregation from dask-expr ( dask#10947 ) Patrick Hoefler
Update tests for assign change in dask-expr ( dask#10944 ) Patrick Hoefler
Adjust for pandas large string change ( dask#10942 ) Patrick Hoefler
Fix flaky test_describe_empty ( dask#10943 ) crusaderky
Use Python 3.12 as reference environment ( dask#10939 ) crusaderky
[Cosmetic] Clean up temp paths in test_config.py ( dask#10938 ) crusaderky
[CLI] dask config set and dask config find updates. ( dask#10930 ) Miles
combine_first when a chunk is full of NaNs ( dask#10932 ) crusaderky
Correctly parse lowercase true/false config from CLI ( dask#10926 ) crusaderky
dask config get fix when printing None values ( dask#10927 ) crusaderky
query-planning can’t be None ( dask#10928 ) crusaderky
Add dask config set ( dask#10921 ) Miles
Make nunique faster again ( dask#10922 ) Patrick Hoefler
Clean up some Cython warnings handling ( dask#10924 ) crusaderky
Bump pre-commit/action from 3.0.0 to 3.0.1 ( dask#10920 )
Raise and avoid data loss of meta provided to P2P shuffle is wrong ( distributed#8520 ) Florian Jetter
Fix gpuci: np.product is deprecated ( distributed#8518 ) crusaderky
Update gpuCI RAPIDS_VER to 24.04 ( distributed#8471 )
Unpin ipywidgets on Python 3.12 ( distributed#8516 ) crusaderky
Keep old dependencies on run_spec collision ( distributed#8512 ) crusaderky
Trivial mypy fix ( distributed#8513 ) crusaderky
Ensure large payload can be serialized and sent over comms ( distributed#8507 ) Florian Jetter
Allow large graph warning threshold to be configured ( distributed#8508 ) Florian Jetter
Tokenization-related test tweaks (backport from #8185) ( distributed#8499 ) crusaderky
Tweaks to update_graph (backport from #8185) ( distributed#8498 ) crusaderky
AMM: test incremental retirements ( distributed#8501 ) crusaderky
Suppress dask-expr warning in CI ( distributed#8505 ) crusaderky
Ignore dask-expr warning in CI ( distributed#8504 ) James Bourbeau
Improve tests for P2P stable ordering ( distributed#8458 ) Hendrik Makait
Bump pre-commit/action from 3.0.0 to 3.0.1 ( distributed#8503 )
Released on February 23, 2024
The last release contained a DeprecationWarning that alerts users to an upcoming switch of dask.dafaframe to use the new backend with support for query planning (see also 10934).
This DeprecationWarning is triggered in import of the dask.dataframe module and the community raised concerns about this being to verbose.
It is now possible to silence this warning
# via Python
>>> dask.config.set({'dataframe.query-planning-warning': False})
# via CLI
dask config set dataframe.query-planning-warning False
See 10936 and 10925 from `Miles`_ for details.
Blockwise fusion optimization can cause a task key collision that is not being handled properly by the distributed scheduler (see 9888). Users will typically notice this by seeing one of various internal exceptions that cause a system deadlock or critical failure. While this issue could not be fixed, the scheduler now implements a mechanism that should mitigate most occurences and issues a warning if the issue is detected.
See 8185 from `crusaderky`_ and `Florian Jetter`_ for details.
Over the course of this, various improvements to tokenization have been implemented. See 10913, 10884, 10919, 10896 and primarily 10883 from `crusaderky`_ for more details.
Adaptive scaling could previously lose data during downscaling if many tasks had to be moved. This typically, but not exclusively, occured on large clusters and would manifest as a recomputation of tasks and could cause clusters to oscillate between up- and downscaling without ever finishing.
See 8522 from `crusaderky`_ for more details.
The current Dask DataFrame implementation is deprecated. In a future release, Dask DataFrame will use new implementation that contains several improve…
Released on February 9, 2024
The current Dask DataFrame implementation is deprecated. In a future release, Dask DataFrame will use new implementation that contains several improvements including a logical query planning. The user-facing DataFrame API will remain unchanged.
The new implementation is already available and can be enabled by installing the dask-expr library:
$ pip install dask-expr
and turning the query planning option on:
import dask >>> dask . config . set ({ 'dataframe.query-planning' : True }) >>> import dask.dataframe as dd
API documentation for the new implementation is available at https://docs.dask.org/en/stable/dataframe-api.html
Any feedback can be reported on the Dask issue tracker dask/dask#issues
See dask#10912 from Patrick Hoefler for details.
This release contains several improvements to Dask’s object tokenization logic. More objects now produce deterministic tokens, which can lead to improved performance through caching of intermediate results.
See dask#10898 , dask#10904 , dask#10876 , dask#10874 , and dask#10865 from crusaderky for details. Additional changes
Fix inplace modification on read-only arrays for string conversion ( dask#10886 ) Patrick Hoefler
Add changelog entry for dask-expr ( dask#10915 ) Patrick Hoefler
Fix leftsemi merge for cudf ( dask#10914 ) Patrick Hoefler
Slight update to dask-expr warning ( dask#10916 ) James Bourbeau
Improve performance for groupby.nunique ( dask#10910 ) Patrick Hoefler
Add configuration for leftsemi merges in dask-expr ( dask#10908 ) Patrick Hoefler
Adjust assign test for dask-expr ( dask#10907 ) Patrick Hoefler
Avoid pytest.warns in test_to_datetime for GPU CI ( dask#10902 ) Richard (Rick) Zamora
Update deployment options in docs homepage ( dask#10901 ) James Bourbeau
Fix typo in dataframe docs ( dask#10900 ) Matthew Rocklin
Bump peter-evans/create-pull-request from 5 to 6 ( dask#10894 )
Fix mimesis API >=13.1.0 - use random.randint ( dask#10888 ) Miles
Adjust invalid test ( dask#10897 ) Patrick Hoefler
Pickle da.argwhere and da.count_nonzero ( dask#10885 ) crusaderky
Fix dask-expr tests after singleton pr ( dask#10892 ) Patrick Hoefler
Set lower bound version for s3fs ( dask#10889 ) Miles
Add a couple of dask-expr fixes for new parquet cache ( dask#10880 ) Florian Jetter
Update deployment documentation ( dask#10882 ) Matthew Rocklin
Start with dask-expr doc build ( dask#10879 ) Patrick Hoefler
Test tokenization of static and class methods ( dask#10872 ) crusaderky
Add distributed.print and distributed.warn to API docs ( dask#10878 ) James Bourbeau
Run macos ci on M1 architecture ( dask#10877 ) Patrick Hoefler
Update tests for dask-expr ( dask#10838 ) Patrick Hoefler
Update parquet tests to align with dask-expr fixes ( dask#10851 ) Richard (Rick) Zamora
Fix regression in test_graph_manipulation ( dask#10873 ) crusaderky
Adjust pytest errors for dask-expr ci ( dask#10871 ) Patrick Hoefler
Set upper bound version for numba when pandas<2.1 ( dask#10890 ) Miles
Deprecate method parameter in DataFrame.fillna ( dask#10846 ) Miles
Remove warning filter from pyproject.toml ( dask#10867 ) Patrick Hoefler
Skip test_append_with_partition for fastparquet ( dask#10828 ) Patrick Hoefler
Fix pytest 8 issues ( dask#10868 ) Patrick Hoefler
Adjust test for support of median in Groupby.aggregate in dask-expr (2/2) ( dask#10870 ) Hendrik Makait
Allow length of ascending to be larger than one in sort_values ( dask#10864 ) Florian Jetter
Allow other message raised in Python 3.9 ( dask#10862 ) Hendrik Makait
Don’t crash when getting computation code in pathological cases ( distributed#8502 ) James Bourbeau
Bump peter-evans/create-pull-request from 5 to 6 ( distributed#8494 )
fix test of cudf spilling metrics ( distributed#8478 ) Mads R. B. Kristensen
Upgrade to pytest 8 ( distributed#8482 ) crusaderky
Fix test_two_consecutive_clients_share_results ( distributed#8484 ) crusaderky
Client word mix-up ( distributed#8481 ) templiert
Released on February 9, 2024
The current Dask DataFrame implementation is deprecated. In a future release, Dask DataFrame will use new implementation that contains several improvements including a logical query planning. The user-facing DataFrame API will remain unchanged.
The new implementation is already available and can be enabled by installing the dask-expr library:
$ pip install dask-expr
and turning the query planning option on:
>>> import dask
>>> dask.config.set({'dataframe.query-planning': True})
>>> import dask.dataframe as dd
API documentation for the new implementation is available at https://docs.dask.org/en/stable/dataframe-api.html
Any feedback can be reported on the Dask issue tracker https://github.com/dask/dask/issues
See 10912 from `Patrick Hoefler`_ for details.
This release contains several improvements to Dask's object tokenization logic. More objects now produce deterministic tokens, which can lead to improved performance through caching of intermediate results.
See 10898, 10904, 10876, 10874, and 10865 from `crusaderky`_ for details.
- Deprecate convert_dtype in apply ( dask#10827 ) Miles
Released on January 26, 2024
This release contains compatibility updates for the latest pandas and scipy releases.
See dask#10834 , dask#10849 , dask#10845 , and distributed#8474 from crusaderky for details.
Deprecate convert_dtype in apply ( dask#10827 ) Miles
Deprecate axis in DataFrame.rolling ( dask#10803 ) Miles
Deprecate out= and dtype= parameter in most DataFrame methods ( dask#10800 ) crusaderky
Deprecate axis in groupby cumulative transformers ( dask#10796 ) Miles
Rename shuffle to shuffle_method in remaining methods ( dask#10797 ) Miles
Additional changes
Add recommended deployment options to deployment docs ( dask#10866 ) James Bourbeau
Improve _agg_finalize to confirm to output expectation ( dask#10835 ) Hendrik Makait
Implement deterministic tokenization for hlg ( dask#10817 ) Patrick Hoefler
Refactor: move tests for tokenize() to its own module ( dask#10863 ) crusaderky
Update DataFrame examples section ( dask#10856 ) James Bourbeau
Temporarily pin mimesis<13.1.0 ( dask#10860 ) James Bourbeau
Trivial cosmetic tweaks to _testing.py ( dask#10857 ) crusaderky
Unskip and adjust tests for groupby -aggregate with median using dask-expr ( dask#10832 ) Hendrik Makait
Fix test for sizeof(pd.MultiIndex) in upstream CI ( dask#10850 ) crusaderky
numpy 2.0: fix slicing by uint64 array ( dask#10854 ) crusaderky
Rename numpy version constants to match pandas ( dask#10843 ) crusaderky
Bump actions/cache from 3 to 4 ( dask#10852 )
Update gpuCI RAPIDS_VER to 24.04 ( dask#10841 )
Fix deprecations in doctest ( dask#10844 ) crusaderky
Changed dtype arithmetics in numpy 2.x ( dask#10831 ) crusaderky
Adjust tests for median support in dask-expr ( dask#10839 ) Patrick Hoefler
Adjust tests for median support in groupby-aggregate in dask-expr ( dask#10840 ) Hendrik Makait
numpy 2.x: fix std() on MaskedArray ( dask#10837 ) crusaderky
Fail dask-expr ci if tests fail ( dask#10829 ) Patrick Hoefler
Activate query_planning when exporting tests ( dask#10833 ) Patrick Hoefler
Expose dataframe tests ( dask#10830 ) Patrick Hoefler
numpy 2: deprecations in n-dimensional fft functions ( dask#10821 ) crusaderky
Generalize CreationDispatch for dask-expr ( dask#10794 ) Richard (Rick) Zamora
Remove circular import when dask-expr enabled ( dask#10824 ) Miles
Minor[CI]: publish-test-results not marked as failed ( dask#10825 ) Miles
Fix more tests to use pytest.warns() ( dask#10818 ) Michał Górny
np.unique() : inverse is shaped in numpy 2 ( dask#10819 ) crusaderky
Pin test_split_adaptive_files to pyarrow engine ( dask#10820 ) Patrick Hoefler
Adjust remaining tests in dask/dask ( dask#10813 ) Patrick Hoefler
Restrict test to Arrow only ( dask#10814 ) Patrick Hoefler
Filter warnings from std test ( dask#10815 ) Patrick Hoefler
Adjust mostly indexing tests ( dask#10790 ) Patrick Hoefler
Updates to deployment docs ( dask#10778 ) Sarah Charlotte Johnson
Unblock documentation build ( dask#10807 ) Miles
Adjust test_to_datetime for dask-expr compatibility Hendrik Makait
Upstream CI tweaks ( dask#10806 ) crusaderky
Improve tests for to_numeric ( dask#10804 ) Hendrik Makait
Fix test-report cache key indent ( dask#10798 ) Miles
Add test-report workflow ( dask#10783 ) Miles
Handle matrix subclass serialization ( distributed#8480 ) Florian Jetter
Use smallest data type for partition column in P2P ( distributed#8479 ) Florian Jetter
pandas 2.2: fix test_dataframe_groupby_tasks ( distributed#8475 ) crusaderky
Bump actions/cache from 3 to 4 ( distributed#8477 )
pandas 2.2 vs. pyarrow 14: deprecated DatetimeTZBlock ( distributed#8476 ) crusaderky
pandas 2.2.0: Deprecated frequency alias M in favor of ME ( distributed#8473 ) Hendrik Makait
Fix docs build ( distributed#8472 ) Hendrik Makait
Fix P2P-based joins with explicit npartitions ( distributed#8470 ) Hendrik Makait
Ignore dask-expr in test_report.py script ( distributed#8464 ) Miles
Nit: hardcode Python version in test report environment ( distributed#8462 ) crusaderky
Change test_report.py - skip bad artifacts in dask/dask ( distributed#8461 ) Miles
Replace all occurrences of sys.is_finalizing ( distributed#8449 ) Florian Jetter
Released on January 26, 2024
This release contains compatibility updates for the latest pandas and scipy releases.
See 10834, 10849, 10845, and 8474 from `crusaderky`_ for details.
Deprecate convert_dtype in apply (10827) `Miles`_
Deprecate axis in DataFrame.rolling (10803) `Miles`_
Deprecate out= and dtype= parameter in most DataFrame methods (10800) `crusaderky`_
Deprecate axis in groupby cumulative transformers (10796) `Miles`_
Rename shuffle to shuffle_method in remaining methods (10797) `Miles`_
The fastparquet Parquet engine has been deprecated. Users should migrate to the pyarrow engine by installing PyArrow and removing engine="fastparquet"…
Released on January 12, 2024
P2P rechunking now utilizes the relationships between input and output chunks. For situations that do not require all-to-all data transfer, this may significantly reduce the runtime and memory/disk footprint. It also enables task culling.
See distributed#8330 from Hendrik Makait for details.
The fastparquet Parquet engine has been deprecated. Users should migrate to the pyarrow engine by installing PyArrow and removing engine="fastparquet" in read_parquet or to_parquet calls.
See dask#10743 from crusaderky for details.
This release improves serialization robustness for arbitrary data. Previously there were some cases where serialization could fail for non- msgpack serializable data. In those cases we now fallback to using pickle .
See dask#8447 from Hendrik Makait for details.
Deprecate shuffle keyword in favour of shuffle_method for DataFrame methods ( dask#10738 ) Hendrik Makait
Deprecate automatic argument inference in repartition ( dask#10691 ) Patrick Hoefler
Deprecate compute parameter in set_index ( dask#10784 ) Miles
Deprecate inplace in eval ( dask#10785 ) Miles
Deprecate Series.view ( dask#10754 ) Miles
Deprecate npartitions="auto" for set_index & sort_values ( dask#10750 ) Miles
Additional changes
Avoid shortcut in tasks shuffle that let to data loss ( dask#10763 ) Patrick Hoefler
Ignore data tasks when ordering ( dask#10706 ) Florian Jetter
Add get_dummies from dask-expr ( dask#10791 ) Patrick Hoefler
Adjust IO tests for dask-expr migration ( dask#10776 ) Patrick Hoefler
Remove deprecation warning about sort and split_out in groupby ( dask#10788 ) Patrick Hoefler
Address pandas deprecations ( dask#10789 ) Patrick Hoefler
Import distributed only once in get_scheduler ( dask#10771 ) Florian Jetter
Simplify GitHub actions ( dask#10781 ) crusaderky
Add unit test overview ( dask#10769 ) Miles
Clean up redundant bits in CI ( dask#10768 ) crusaderky
Update tests for ufunc ( dask#10773 ) Patrick Hoefler
Use pytest.mark.skipif(DASK_EXPR_ENABLED) ( dask#10774 ) crusaderky
Adjust shuffle tests for dask-expr ( dask#10759 ) Patrick Hoefler
Fix some deprecation warnings from pandas ( dask#10749 ) Patrick Hoefler
Adjust shuffle tests for dask-expr ( dask#10762 ) Patrick Hoefler
Update pre-commit ( dask#10767 ) Hendrik Makait
Clean up config switches in CI ( dask#10766 ) crusaderky
Improve exception for validate_key ( dask#10765 ) Hendrik Makait
Handle datetimeindexes in set_index with unknown divisions ( dask#10757 ) Patrick Hoefler
Add hashing for decimals ( dask#10758 ) Patrick Hoefler
Review tests for is_monotonic ( dask#10756 ) crusaderky
Change argument order in value_counts_aggregate ( dask#10751 ) Patrick Hoefler
Adjust some groupby tests for dask-expr ( dask#10752 ) Patrick Hoefler
Restrict mimesis to < 12 for 3.9 build ( dask#10755 ) Patrick Hoefler
Don’t evaluate config in skip condition ( dask#10753 ) Patrick Hoefler
Adjust some tests to be compatible with dask-expr ( dask#10714 ) Patrick Hoefler
Make dask.array.utils functions more generic to other Dask Arrays ( dask#10676 ) Matthew Rocklin
Remove duplciate “single machine” section ( dask#10747 ) Matthew Rocklin
Tweak ORC engine= parameter ( dask#10746 ) crusaderky
Add pandas 3.0 deprecations and migration prep for dask-expr ( dask#10723 ) Miles
Add task graph animation to docs homepage ( dask#10730 ) Sarah Charlotte Johnson
Use new Xarray logo ( dask#10729 ) James Bourbeau
Update tab styling on “10 Minutes to Dask” page ( dask#10728 ) James Bourbeau
Update environment file upload step in CI ( dask#10726 ) James Bourbeau
Don’t duplicate unobserved categories in GroupBy.nunqiue if split_out>1 ( dask#10716 ) Patrick Hoefler
Changelog entry for dask.order update ( dask#10715 ) Florian Jetter
Relax redundant-key check in _check_dsk ( dask#10701 ) Richard (Rick) Zamora
Fix test_report.py ( distributed#8459 ) Miles
Revert pickle change ( distributed#8456 ) Florian Jetter
Adapt test_report.py to support dask/dask repository ( distributed#8450 ) Miles
Maintain stable ordering for P2P shuffling ( distributed#8453 ) Hendrik Makait
Add no worker timeout for scheduler ( distributed#8371 ) FTang21
Allow tests workflow to be dispatched manually by maintainers ( distributed#8445 ) Erik Sundell
Make scheduler-related transition functionality private ( distributed#8448 ) Hendrik Makait
Update pre-commit hooks ( distributed#8444 ) Hendrik Makait
Do not always check if main in result when pickling ( distributed#8443 ) Florian Jetter
Delegate wait_for_workers to cluster instances only when implemented ( distributed#8441 ) Erik Sundell
Extend sleep in test_pandas ( distributed#8440 ) Julian Gilbey
Avoid deprecated shuffle keyword ( distributed#8439 ) Hendrik Makait
Shuffle metrics 4/4: Remove bespoke diagnostics ( distributed#8367 ) crusaderky
Do not run gilknocker in testsuite ( distributed#8423 ) Florian Jetter
Tweak abstractmethods ( distributed#8427 ) crusaderky
Shuffle metrics 3/4: Capture background metrics ( distributed#8366 ) crusaderky
Shuffle metrics 2/4: Add background metrics ( distributed#8365 ) crusaderky
Shuffle metrics 1/4: Add foreground metrics ( distributed#8364 ) crusaderky
Bump actions/upload-artifact from 3 to 4 ( distributed#8420 )
Fix test_merge_p2p_shuffle_reused_dataframe_with_different_parameters ( distributed#8422 ) Hendrik Makait
Expand Client.upload_file docs example ( distributed#8313 ) Miles
Improve logging in P2P’s scheduler plugin ( distributed#8410 ) Hendrik Makait
Re-enable test_decide_worker_coschedule_order_neighbors ( distributed#8402 ) Florian Jetter
Add cuDF spilling statistics to RMM/GPU memory plot ( distributed#8148 ) Charles Blackmon-Luca
Fix inconsistent hashing for Nanny-spawned workers ( distributed#8400 ) Charles Stern
Do not allow workers to downscale if they are running long-running tasks (e.g. worker_client ) ( distributed#7481 ) Florian Jetter
Fix flaky test_subprocess_cluster_does_not_depend_on_logging ( distributed#8417 ) crusaderky
Released on January 12, 2024
P2P rechunking now utilizes the relationships between input and output chunks. For situations that do not require all-to-all data transfer, this may significantly reduce the runtime and memory/disk footprint. It also enables task culling.
See 8330 from `Hendrik Makait`_ for details.
The fastparquet Parquet engine has been deprecated. Users should migrate to the pyarrow engine by installing PyArrow and removing engine="fastparquet" in read_parquet or to_parquet calls.
See 10743 from `crusaderky`_ for details.
This release improves serialization robustness for arbitrary data. Previously there were some cases where serialization could fail for non-msgpack serializable data. In those cases we now fallback to using pickle.
See 8447 from `Hendrik Makait`_ for details.
Deprecate shuffle keyword in favour of shuffle_method for DataFrame methods (10738) `Hendrik Makait`_
Deprecate automatic argument inference in repartition (10691) `Patrick Hoefler`_
Deprecate compute parameter in set_index (10784) `Miles`_
Deprecate inplace in eval (10785) `Miles`_
Deprecate Series.view (10754) `Miles`_
Deprecate npartitions="auto" for set_index & sort_values (10750) `Miles`_
This feature is still under active development and the API isn’t stable yet, so breaking changes can occur. We expect to make the query optimizer the…
Released on December 15, 2023
Dask DataFrames are now much more performant by using a logical query planner. This feature is currently off by default, but can be turned on with:
dask . config . set ({ "dataframe.query-planning" : True })
You also need to have dask-expr installed:
pip install dask-expr
We’ve seen promising performance improvements so far, see this blog post and these regularly updated benchmarks for more information. A more detailed explanation of how the query optimizer works can be found in this blog post .
This feature is still under active development and the API isn’t stable yet, so breaking changes can occur. We expect to make the query optimizer the default early next year.
See dask#10634 from Patrick Hoefler for details.
read_parquet will now infer the Arrow types pa.date32() , pa.date64() and pa.decimal() as a ArrowDtype in pandas. These dtypes are backed by the original Arrow array, and thus avoid the conversion to NumPy object. Additionally, read_parquet will no longer infer nested and binary types as strings, they will be stored in NumPy object arrays.
See dask#10698 and dask#10705 from Patrick Hoefler for details.
This release includes a major rewrite to a core part of our scheduling logic. It includes a new approach to the topological sorting algorithm in dask.order which determines the order in which tasks are run. Improper ordering is known to be a major contributor to too large cluster memory pressure.
Updates in this release fix a couple of performance regressions that were introduced in the release 2023.10.0 (see dask#10535 ). Generally, computations should now be much more eager to release data if it is no longer required in memory.
See dask#10660 , dask#10697 from Florian Jetter for details.
This release contains several updates that fix a possible deadlock introduced in 2023.9.2 and improve the robustness of P2P-based merging when the cluster is dynamically scaling up.
See distributed#8415 , distributed#8416 , and distributed#8414 from Hendrik Makait for details.
The distributed.scheduler.pickle configuration option is no longer supported. As of the 2023.4.0 release, pickle is used to transmit task graphs, so can no longer be disabled. We now raise an informative error when distributed.scheduler.pickle is set to False .
See distributed#8401 from Florian Jetter for details. Additional changes
Add changelog entry for recent P2P merge fixes ( dask#10712 ) Hendrik Makait
Update DataFrame page ( dask#10710 ) Matthew Rocklin
Add changelog entry for dask-expr switch ( dask#10704 ) Patrick Hoefler
Improve changelog entry for PipInstall changes ( dask#10711 ) Hendrik Makait
Remove PR labeler ( dask#10709 ) James Bourbeau
Add .wrapped to Delayed object ( dask#10695 ) Andrew S. Rosen
Bump actions/labeler from 4.3.0 to 5.0.0 ( dask#10689 )
Bump actions/stale from 8 to 9 ( dask#10690 )
[Dask.order] Remove non-runnable leaf nodes from ordering ( dask#10697 ) Florian Jetter
Update installation docs ( dask#10699 ) Matthew Rocklin
Fix software environment link in docs ( dask#10700 ) James Bourbeau
Avoid converting non-strings to arrow strings for read_parquet ( dask#10692 ) Patrick Hoefler
Bump xarray-contrib/issue-from-pytest-log from 1.2.7 to 1.2.8 ( dask#10687 )
Fix tokenize for pd.DateOffset ( dask#10664 ) jochenott
Bugfix for writing empty array to zarr ( dask#10506 ) Ben
Docs update, fixup styling, mention free ( dask#10679 ) Matthew Rocklin
Update deployment docs ( dask#10680 ) Matthew Rocklin
Dask.order rewrite using a critical path approach ( dask#10660 ) Florian Jetter
Avoid substituting keys that occur multiple times ( dask#10646 ) Florian Jetter
Add missing image to docs ( dask#10694 ) Matthew Rocklin
Bump actions/setup-python from 4 to 5 ( dask#10688 )
Update landing page ( dask#10674 ) Matthew Rocklin
Make meta check simpler in dispatch ( dask#10638 ) Patrick Hoefler
Pin PR Labeler ( dask#10675 ) Matthew Rocklin
Reorganize docs index a bit ( dask#10669 ) Matthew Rocklin
Bump actions/setup-java from 3 to 4 ( dask#10667 )
Bump conda-incubator/setup-miniconda from 2.2.0 to 3.0.1 ( dask#10668 )
Bump xarray-contrib/issue-from-pytest-log from 1.2.6 to 1.2.7 ( dask#10666 )
Fix test_categorize_info with nightly pyarrow ( dask#10662 ) James Bourbeau
Rewrite test_subprocess_cluster_does_not_depend_on_logging ( distributed#8409 ) Hendrik Makait
Avoid RecursionError when failing to pickle key in SpillBuffer and using tblib=3 ( distributed#8404 ) Hendrik Makait
Allow tasks to override is_rootish heuristic ( distributed#8412 ) Hendrik Makait
Remove GPU executor ( distributed#8399 ) Hendrik Makait
Do not rely on logging for subprocess cluster ( distributed#8398 ) Hendrik Makait
Update gpuCI RAPIDS_VER to 24.02 ( distributed#8384 )
Bump actions/setup-python from 4 to 5 ( distributed#8396 )
Ensure output chunks in P2P rechunking are distributed homogeneously ( distributed#8207 ) Florian Jetter
Trivial: fix typo ( distributed#8395 ) crusaderky
Bump JamesIves/github-pages-deploy-action from 4.4.3 to 4.5.0 ( distributed#8387 )
Bump conda-incubator/setup-miniconda from 3.0.0 to 3.0.1 ( distributed#8388 )
Released on December 15, 2023
Dask DataFrames are now much more performant by using a logical query planner. This feature is currently off by default, but can be turned on with:
dask.config.set({"dataframe.query-planning": True})
You also need to have dask-expr installed:
pip install dask-expr
We've seen promising performance improvements so far, see this blog post and these regularly updated benchmarks for more information. A more detailed explanation of how the query optimizer works can be found in this blog post.
This feature is still under active development and the API isn't stable yet, so breaking changes can occur. We expect to make the query optimizer the default early next year.
See 10634 from `Patrick Hoefler`_ for details.
read_parquet will now infer the Arrow types pa.date32(), pa.date64() and pa.decimal() as a ArrowDtype in pandas. These dtypes are backed by the original Arrow array, and thus avoid the conversion to NumPy object. Additionally, read_parquet will no longer infer nested and binary types as strings, they will be stored in NumPy object arrays.
See 10698 and 10705 from `Patrick Hoefler`_ for details.
This release includes a major rewrite to a core part of our scheduling logic. It includes a new approach to the topological sorting algorithm in dask.order which determines the order in which tasks are run. Improper ordering is known to be a major contributor to too large cluster memory pressure.
Updates in this release fix a couple of performance regressions that were introduced in the release 2023.10.0 (see 10535). Generally, computations should now be much more eager to release data if it is no longer required in memory.
See 10660, 10697 from `Florian Jetter`_ for details.
This release contains several updates that fix a possible deadlock introduced in 2023.9.2 and improve the robustness of P2P-based merging when the cluster is dynamically scaling up.
See 8415, 8416, and 8414 from `Hendrik Makait`_ for details.
The distributed.scheduler.pickle configuration option is no longer supported. As of the 2023.4.0 release, pickle is used to transmit task graphs, so can no longer be disabled. We now raise an informative error when distributed.scheduler.pickle is set to False.
See 8401 from `Florian Jetter`_ for details.
The distributed.PipInstall plugin now has more robust restart logic and also supports environment variables .
Released on December 1, 2023
The distributed.PipInstall plugin now has more robust restart logic and also supports environment variables .
Below shows how users can use the distributed.PipInstall plugin and a TOKEN environment variable to securely install a package from a private repository:
from dask.distributed import PipInstall plugin = PipInstall ( packages = [ "private_package@git+https://$ {TOKEN} @github.com/dask/private_package.git]) client . register_plugin ( plugin )
See distributed#8374 , distributed#8357 , and distributed#8343 from Hendrik Makait for details.
This release contains compatibility updates for using bokeh>=3.3.0 with proxied Dask dashboards. Previously the contents of dashboard plots wouldn’t be displayed.
See distributed#8347 and distributed#8381 from Jacob Tomlinson for details. Additional changes
Add network marker to test_pyarrow_filesystem_option_real_data ( dask#10653 ) Richard (Rick) Zamora
Bump GPU CI to CUDA 11.8 ( dask#10656 ) Charles Blackmon-Luca
Tokenize pandas offsets deterministically ( dask#10643 ) Patrick Hoefler
Add tokenize pd.NA functionality ( dask#10640 ) Patrick Hoefler
Update gpuCI RAPIDS_VER to 24.02 ( dask#10636 )
Fix precision handling in array.linalg.norm ( dask#10556 ) joanrue
Add axis argument to DataFrame.clip and Series.clip ( dask#10616 ) Richard (Rick) Zamora
Update changelog entry for in-memory rechunking ( dask#10630 ) Florian Jetter
Fix flaky test_resources_reset_after_cancelled_task ( distributed#8373 ) crusaderky
Bump GPU CI to CUDA 11.8 ( distributed#8376 ) Charles Blackmon-Luca
Bump conda-incubator/setup-miniconda from 2.2.0 to 3.0.0 ( distributed#8372 )
Add debug logs to P2P scheduler plugin ( distributed#8358 ) Hendrik Makait
O(1) access for /info/task/ endpoint ( distributed#8363 ) crusaderky
Remove stringification from shuffle annotations ( distributed#8362 ) crusaderky
Don’t cast int metrics to float ( distributed#8361 ) crusaderky
Drop asyncio TCP backend ( distributed#8355 ) Florian Jetter
Add offload support to context_meter.add_callback ( distributed#8360 ) crusaderky
Test that sync() propagates contextvars ( distributed#8354 ) crusaderky
captured_context_meter ( distributed#8352 ) crusaderky
context_meter.clear_callbacks ( distributed#8353 ) crusaderky
Use @log_errors decorator ( distributed#8351 ) crusaderky
Fix test_statistical_profiling_cycle ( distributed#8356 ) Florian Jetter
Shuffle: don’t parse dask.config at every RPC ( distributed#8350 ) crusaderky
Replace Client.register_plugin s idempotent argument with .idempotent attribute on plugins ( distributed#8342 ) Hendrik Makait
Fix test report generation ( distributed#8346 ) Hendrik Makait
Install pyarrow-hotfix on mindeps-pandas CI ( distributed#8344 ) Hendrik Makait
Reduce memory usage of scheduler process - optimize scheduler.py::TaskState class ( distributed#8331 ) Miles
Bump pre-commit linters ( distributed#8340 ) crusaderky
Update cuDF test with explicit dtype=object ( distributed#8339 ) Peter Andreas Entschev
Fix Cluster / SpecCluster calls to async close methods ( distributed#8327 ) Peter Andreas Entschev
Released on December 1, 2023
The distributed.PipInstall plugin now has more robust restart logic and also supports environment variables.
Below shows how users can use the distributed.PipInstall plugin and a TOKEN environment variable to securely install a package from a private repository:
from dask.distributed import PipInstall
plugin = PipInstall(packages=["private_package@git+https://${TOKEN}@github.com/dask/private_package.git])
client.register_plugin(plugin)
See 8374, 8357, and 8343 from `Hendrik Makait`_ for details.
This release contains compatibility updates for using bokeh>=3.3.0 with proxied Dask dashboards. Previously the contents of dashboard plots wouldn't be displayed.
See 8347 and 8381 from `Jacob Tomlinson`_ for details.
pyarrow<14.0.1 usage is deprecated starting in this release. It’s recommended for all users to upgrade their version of pyarrow or install pyarrow-hot…
Released on November 10, 2023
Users should see significant performance improvements when using in-memory P2P array rechunking. This is due to no longer copying underlying data buffers.
Below shows a simple example where we compare performance of different rechunking methods.
shape = ( 30_000 , 6_000 , 150 ) # 201.17 GiB input_chunks = ( 60 , - 1 , - 1 ) # 411.99 MiB output_chunks = ( - 1 , 6 , - 1 ) # 205.99 MiB arr = da . random . random ( size , chunks = input_chunks ) with dask . config . set ({ "array.rechunk.method" : "p2p" , "distributed.p2p.disk" : True , }): ( da . random . random ( size , chunks = input_chunks ) . rechunk ( output_chunks ) . sum () . compute () )
See distributed#8282 , distributed#8318 , distributed#8321 from crusaderky and ( distributed#8322 ) from Hendrik Makait for details.
pyarrow<14.0.1 usage is deprecated starting in this release. It’s recommended for all users to upgrade their version of pyarrow or install pyarrow-hotfix . See this CVE for full details.
See dask#10622 from Florian Jetter for details.
Using filesystem="arrow" when reading Parquet datasets now properly inferrs the correct cloud region when accessing remote, cloud-hosted data.
See dask#10590 from Richard (Rick) Zamora for details.
See distributed#8332 from Hendrik Makait for details. Additional changes
Fix sporadic failure of test_dataframe::test_quantile ( dask#10625 ) Miles
Bump minimum click to >=8.1 ( dask#10623 ) Jacob Tomlinson
Refactor test_quantile ( dask#10620 ) Miles
Avoid PerformanceWarning for fragmented DataFrame ( dask#10621 ) Patrick Hoefler
Generalize computation of NEW_*_VER in GPU CI updating workflow ( dask#10610 ) Charles Blackmon-Luca
Switch to newer GPU CI images ( dask#10608 ) Charles Blackmon-Luca
Remove double slash in fsspec tests ( dask#10605 ) Mario Šaško
Reenable test_ucx_config_w_env_var ( distributed#8272 ) Peter Andreas Entschev
Don’t share host_array when receiving from network ( distributed#8308 ) crusaderky
Generalize computation of NEW_*_VER in GPU CI updating workflow ( distributed#8319 ) Charles Blackmon-Luca
Switch to newer GPU CI images ( distributed#8316 ) Charles Blackmon-Luca
Minor updates to shuffle dashboard ( distributed#8315 ) Matthew Rocklin
Don’t use bytearray().join ( distributed#8312 ) crusaderky
Reuse identical shuffles in P2P hash join ( distributed#8306 ) Hendrik Makait
Released on November 10, 2023
Users should see significant performance improvements when using in-memory P2P array rechunking. This is due to no longer copying underlying data buffers.
Below shows a simple example where we compare performance of different rechunking methods.
shape = (30_000, 6_000, 150) # 201.17 GiB
input_chunks = (60, -1, -1) # 411.99 MiB
output_chunks = (-1, 6, -1) # 205.99 MiB
arr = da.random.random(size, chunks=input_chunks)
with dask.config.set({
"array.rechunk.method": "p2p",
"distributed.p2p.disk": True,
}):
(
da.random.random(size, chunks=input_chunks)
.rechunk(output_chunks)
.sum()
.compute()
)
See 8282, 8318, 8321 from `crusaderky`_ and (8322) from `Hendrik Makait`_ for details.
pyarrow<14.0.1 usage is deprecated starting in this release. It's recommended for all users to upgrade their version of pyarrow or install pyarrow-hotfix. See this CVE for full details.
See 10622 from `Florian Jetter`_ for details.
Using filesystem="arrow" when reading Parquet datasets now properly inferrs the correct cloud region when accessing remote, cloud-hosted data.
See 10590 from `Richard (Rick) Zamora`_ for details.
See 8332 from `Hendrik Makait`_ for details.
- Unignore and fix deprecated freq aliases ( dask#10577 ) Thomas Grainger
Released on October 27, 2023
This release adds official support for Python 3.12.
See dask#10544 and distributed#8223 from Thomas Grainger for details. Additional changes
Avoid splitting parquet files to row groups as aggressively ( dask#10600 ) Matthew Rocklin
Speed up normalize_chunks for common case ( dask#10579 ) Martin Durant
Use Python 3.11 for upstream and doctests CI build ( dask#10596 ) Thomas Grainger
Bump actions/checkout from 4.1.0 to 4.1.1 ( dask#10592 )
Switch to PyTables HEAD ( dask#10580 ) Thomas Grainger
Remove numpy.core warning filter, link to issue on pyarrow caused BlockManager warning ( dask#10571 ) Thomas Grainger
Unignore and fix deprecated freq aliases ( dask#10577 ) Thomas Grainger
Move register_assert_rewrite earlier in conftest to fix warnings ( dask#10578 ) Thomas Grainger
Upgrade versioneer to 0.29 ( dask#10575 ) Thomas Grainger
change test_concat_categorical to be non-strict ( dask#10574 ) Thomas Grainger
Enable SciPy tests with NumPy 2.0 Thomas Grainger
Enable tests for scikit-image with NumPy 2.0 ( dask#10569 ) Thomas Grainger
Fix upstream build ( dask#10549 ) Thomas Grainger
Add optimized code paths for drop_duplicates ( dask#10542 ) Richard (Rick) Zamora
Support cudf backend in dd.DataFrame.sort_values ( dask#10551 ) Richard (Rick) Zamora
Rename “GIL Contention” to just GIL in chart labels ( distributed#8305 ) Matthew Rocklin
Bump actions/checkout from 4.1.0 to 4.1.1 ( distributed#8299 )
Fix dashboard ( distributed#8293 ) Hendrik Makait
@log_errors for async tasks ( distributed#8294 ) crusaderky
Annotations and better tests for serialize_bytes ( distributed#8300 ) crusaderky
Temporarily xfail test_decide_worker_coschedule_order_neighbors to unblock CI ( distributed#8298 ) James Bourbeau
Skip xdist and matplotlib in code samples ( distributed#8290 ) Matthew Rocklin
Use numpy._core on numpy>=2.dev0 ( distributed#8291 ) Thomas Grainger
Fix calculation of MemoryShardsBuffer.bytes_read ( distributed#8289 ) crusaderky
Allow P2P to store data in-memory ( distributed#8279 ) Hendrik Makait
Upgrade versioneer to 0.29 ( distributed#8288 ) Thomas Grainger
Allow ResourceLimiter to be unlimited ( distributed#8276 ) Hendrik Makait
Run pre-commit autoupdate ( distributed#8281 ) Thomas Grainger
Annotate instance variables for P2P layers ( distributed#8280 ) Hendrik Makait
Remove worker gracefully should not mark tasks as suspicious ( distributed#8234 ) Thomas Grainger
Add signal handling to dask spec ( distributed#8261 ) Thomas Grainger
Add typing for sync ( distributed#8275 ) Hendrik Makait
Better annotations for shuffle offload ( distributed#8277 ) crusaderky
Test minimum versions for p2p shuffle ( distributed#8270 ) crusaderky
Run coverage on test failures ( distributed#8269 ) crusaderky
Use aiohttp with extensions ( distributed#8274 ) Thomas Grainger
Released on October 27, 2023
This release adds official support for Python 3.12.
See 10544 and 8223 from `Thomas Grainger`_ for details.
This release contains major updates to Dask’s task graph scheduling logic. The updates here significantly reduce memory pressure on array reductions.
Released on October 13, 2023
This release contains major updates to Dask’s task graph scheduling logic. The updates here significantly reduce memory pressure on array reductions. We anticipate this will have a strong impact on the array computing community.
See dask#10535 from Florian Jetter for details.
There are several updates (listed below) that make P2P shuffling much more robust and less likely to fail.
See distributed#8262 , distributed#8264 , distributed#8242 , distributed#8244 , and distributed#8235 from Hendrik Makait and distributed#8124 from Charles Blackmon-Luca for details.
Users should see reduced CPU load on their scheduler when computing large task graphs.
See distributed#8238 and dask#10547 from Florian Jetter and distributed#8240 from crusaderky for details. Additional changes
Dispatch the partd.Encode class used for disk-based shuffling ( dask#10552 ) Richard (Rick) Zamora
Add documentation for hive partitioning ( dask#10454 ) Richard (Rick) Zamora
Add typing to dask.order ( dask#10553 ) Florian Jetter
Allow passing index_col=False in dd.read_csv ( dask#9961 ) Michael Leslie
Tighten HighLevelGraph annotations ( dask#10524 ) crusaderky
Support for latest ipykernel / ipywidgets ( distributed#8253 ) crusaderky
Check minimal pyarrow version for P2P merge ( distributed#8266 ) Hendrik Makait
Support for Python 3.12 ( distributed#8223 ) Thomas Grainger
Use memoryview.nbytes when warning on large graph send ( distributed#8268 ) crusaderky
Run tests without gilknocker ( distributed#8263 ) crusaderky
Disable ipv6 on MacOS CI ( distributed#8254 ) crusaderky
Clean up redundant minimum versions ( distributed#8251 ) crusaderky
Clean up use of BARRIER_PREFIX in scheduler plugin ( distributed#8252 ) crusaderky
Improve shuffle run handling in P2P’s worker plugin ( distributed#8245 ) Hendrik Makait
Explicitly set charset=utf-8 ( distributed#8250 ) crusaderky
Typing tweaks to distributed#8239 ( distributed#8247 ) crusaderky
Simplify scheduler assertion ( distributed#8246 ) crusaderky
Improve typing ( distributed#8239 ) Hendrik Makait
Respect cgroups v2 “low” memory limit ( distributed#8243 ) Samantha Hughes
Fix PackageInstall by making it a scheduler plugin ( distributed#8142 ) Hendrik Makait
Xfail test_ucx_config_w_env_var ( distributed#8241 ) crusaderky
SpecCluster resilience to broken workers ( distributed#8233 ) crusaderky
Suppress SpillBuffer stack traces for cancelled tasks ( distributed#8232 ) crusaderky
Update annotations after stringification changes ( distributed#8195 ) crusaderky
Reduce max recursion depth of profile ( distributed#8224 ) crusaderky
Offload deeply nested objects ( distributed#8214 ) crusaderky
Fix flaky test_close_connections ( distributed#8231 ) crusaderky
Fix flaky test_popen_timeout ( distributed#8229 ) crusaderky
Fix flaky test_adapt_then_manual ( distributed#8228 ) crusaderky
Prevent collisions in SpillBuffer ( distributed#8226 ) crusaderky
Allow retire_workers to run concurrently ( distributed#8056 ) Florian Jetter
Fix HTML repr for TaskState objects ( distributed#8188 ) Florian Jetter
Fix AttributeError for builtin_function_or_method in profile.py ( distributed#8181 ) Florian Jetter
Fix flaky test_spans (v2) ( distributed#8222 ) crusaderky
Released on October 13, 2023
This release contains major updates to Dask's task graph scheduling logic. The updates here significantly reduce memory pressure on array reductions. We anticipate this will have a strong impact on the array computing community.
See 10535 from `Florian Jetter`_ for details.
There are several updates (listed below) that make P2P shuffling much more robust and less likely to fail.
See 8262, 8264, 8242, 8244, and 8235 from `Hendrik Makait`_ and 8124 from `Charles Blackmon-Luca`_ for details.
Users should see reduced CPU load on their scheduler when computing large task graphs.
See 8238 and 10547 from `Florian Jetter`_ and 8240 from `crusaderky`_ for details.
The 2023.9.2 release introduced an unintentional breaking change in how configuration options are overriden in dask.config.get with the override_with=…
Released on September 29, 2023
The 2023.9.2 release introduced an unintentional breaking change in how configuration options are overriden in dask.config.get with the override_with= keyword (see dask#10519 ). This release restores the previous behavior.
See dask#10521 from crusaderky for details.
This release includes improved support for using common reductions in Dask Array (e.g. var , std , moment ) with complex dtypes.
See dask#10009 from wkrasnicki for details. Additional changes
Bump actions/checkout from 4.0.0 to 4.1.0 ( dask#10532 )
Match pandas reverting apply deprecation ( dask#10531 ) James Bourbeau
Update gpuCI RAPIDS_VER to 23.12 ( dask#10526 )
Temporarily skip failing tests with fsspec==2023.9.1 ( dask#10520 ) James Bourbeau
Released on September 29, 2023
The 2023.9.2 release introduced an unintentional breaking change in how configuration options are overriden in dask.config.get with the override_with= keyword (see 10519). This release restores the previous behavior.
See 10521 from `crusaderky`_ for details.
This release includes improved support for using common reductions in Dask Array (e.g. var, std, moment) with complex dtypes.
See 10009 from `wkrasnicki`_ for details.
The 2023.9.0 release modified the admin.traceback.shorten configuration option without introducing a deprecation cycle. This resulted in failures to c…
Released on September 15, 2023
Previously the default shuffling method would silently fallback from P2P to task-based shuffling if an older version of pyarrow was installed. Now we raise an informative error with the minimum required pyarrow version for P2P instead of silently falling back.
See dask#10496 from Hendrik Makait for details.
The 2023.9.0 release modified the admin.traceback.shorten configuration option without introducing a deprecation cycle. This resulted in failures to create Dask clusters in some cases. This release introduces a deprecation cycle for this configuration change.
See dask#10509 from crusaderky for details. Additional changes
Avoid materializing all iterators in delayed tasks ( dask#10498 ) James Bourbeau
Overhaul deprecations system in dask.config ( dask#10499 ) crusaderky
Remove unnecessary check in timeseries ( dask#10447 ) Patrick Hoefler
Use register_plugin in tests ( dask#10503 ) James Bourbeau
Make preserve_index explicit in pyarrow_schema_dispatch ( dask#10501 ) Hendrik Makait
Add **kwargs support for pyarrow_schema_dispatch ( dask#10500 ) Hendrik Makait
Centralize and type no_default ( dask#10495 ) crusaderky
Released on September 15, 2023
Previously the default shuffling method would silently fallback from P2P to task-based shuffling if an older version of pyarrow was installed. Now we raise an informative error with the minimum required pyarrow version for P2P instead of silently falling back.
See 10496 from `Hendrik Makait`_ for details.
The 2023.9.0 release modified the admin.traceback.shorten configuration option without introducing a deprecation cycle. This resulted in failures to create Dask clusters in some cases. This release introduces a deprecation cycle for this configuration change.
See 10509 from `crusaderky`_ for details.
Your coding agent can read these notes before it upgrades. Set up the MCP server →