dask unstable; multiple race condition errors: P2PConsistencyError, RuntimeError, FutureCancelledError ...
Nadie ha tomado este issue todavía.
Evaluación
- Dificultad
- 5/5
- Tiempo estimado
- Más de una semana
- Aptitud para principiantes
- 20/100
- Tipo de issue
- Error
- Claridad
- Necesita aclaración
- Estado de actividad
- Estancado
- Stack tecnológico
- python
- Área
- distributed-systems
Línea de trabajo
Comienza convirtiendo el flujo de trabajo reportado de xarray.Dataset.to_zarr() en un reproductor autónomo y completa los detalles que faltan sobre Python, el sistema operativo, la instalación y el entorno. Aísla la ruta de Dask shuffle/rechunk y determina cuáles de los errores reportados P2PConsistencyError, RuntimeError, FutureCancelledError, de tiempo de espera de conexión y de deserialización se pueden reproducir; el trabajo estará terminado cuando haya un fallo acotado con una prueba de regresión adecuada o una causa raíz claramente identificada.
Escrito por el modelo de indexación a partir del texto del issue.
Descripción
Hi,
I'm encountering repeated failures during xarray.Dataset.to_zarr() when using Dask Distributed.
The failures occur during the Dask shuffle/rechunk phase and appear to be caused by instability in the P2P shuffle as well as many other dask issues.
While processing a batch of NetCDF to zarr, I end up having to create a horrible code which has to catch all possible dask issues, so that, the cluster/client gets destroyed, recreated, and the failed batch of files reprocessed.
Even though recreating the cluster allows the batch to succeed on retry, it shouldn't be required.
Below are all the dask errors I have to catch in my code.
batch_is_processed = False
while not batch_is_processed:
try:
with self.lock:
ds.to_zarr(
self.store,
mode="w", # Overwrite mode for the first batch
write_empty_chunks=self.write_empty_chunks,
compute=True, # Compute the result immediately
consolidated=self.consolidated,
safe_chunks=self.safe_chunks,
align_chunks=self.align_chunks,
)
batch_is_processed = True
except (FutureCancelledError, P2PConsistencyError, RuntimeError) as e:
error_text = str(e)
SHUFFLE_KEYWORDS = [
"P2P",
"failed during transfer phase",
"failed during barrier phase",
"failed during shuffle phase",
"shuffle failure",
"No active shuffle",
"Unexpected error encountered during P2P",
]
CONNECTION_KEYWORDS = [
"Timed out trying to connect",
"CancelledError",
"Too many open files",
]
DESERIALISATION_KEYWORDS = [
"Error during deserialization",
"different environments",
]
# Determine if the error should trigger a retry (cluster reset)
retryable = any(
keyword in error_text
for keyword in (
SHUFFLE_KEYWORDS
+ CONNECTION_KEYWORDS
+ DESERIALISATION_KEYWORDS
)
)
# RuntimeError that is NOT retryable -> treat as normal exception
if isinstance(e, RuntimeError) and not retryable:
# Treat as a regular exception: fallback to individual processing
# code removed for simplification
else:
self._reset_cluster()
batch_is_processed = False
# ... the rest of the logic is in a while loop
I'd like to highlight that my workers/schedulers are only using 50% of cpu/mem max. I've tried many different settings, and ended up having to create workers only with 1 thread if I don't want to run into more race conditions.
Errors Observed
P2PConsistencyError: No active shuffle with id='868df3e600f2968e0cd678a103d355d2' found
RuntimeError: P2P 43b3a653ad1c988df817f013d7314141 failed during transfer phase
OSError: Timed out trying to connect to tls://<worker-ip>:8786 after 30 s
asyncio.exceptions.CancelledError
...
Minimal Complete Verifiable Example:
# Put your MCVE code here
Anything else we need to know?:
Environment:
- Dask version: 2025.11.0
- Python version:
- Operating System:
- Install method (conda, pip, source):
- Lenguaje dominante
- Python
- Estrellas
- 1.7k
- Forks
- 778
- Métricas de merge de PR
- Sin PR fusionados en 30 d
Preparar el entorno
Primeros pasos
- Lee el issue completo y luego la guía de contribución del proyecto.
- Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
- Haz un fork del repositorio y trabaja en una rama.
- Abre un pull request que haga referencia al número del issue.
Más de dask/distributed
-
needs triage
Dificultad 2/5 1-3 horas Aptitud para principiantes 72/100
dask/distributed#9366 ·
-
needs triage
Dificultad 2/5 1-3 horas Aptitud para principiantes 84/100
dask/distributed#9353 ·
-
documentation
Dificultad 1/5 1-3 horas Aptitud para principiantes 82/100
dask/distributed#8304 ·
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
dask/distributed#4816 · 2 comentarios ·
-
documentation good first issue
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
dask/distributed#2378 · 2 comentarios ·
Todos los issues de dask/distributed
Issues similares
-
[Bug] @deck.gl/arcgis dist import resolves to unpublished @deck.gl/core source path (9.3.11, 9.4.0)Abierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 72/100
Los mantenedores suelen responder en 1 día
-
workflow: a tick's dispatch counts as 'only this step', and no review self-grants a round unattendedAbiertoworkflow
Dificultad 2/5 1-3 horas Aptitud para principiantes 85/100
kristofdegrave/homeassistant-smart-charging#1505 ·
Los mantenedores suelen responder en 1 día
-
New Submission: TropWATERAbiertometadata submission
Dificultad 2/5 1-3 horas Aptitud para principiantes 82/100
-
Wrongly named dashboard variableAbiertobug
Dificultad 2/5 1-3 horas Aptitud para principiantes 65/100
canonical/content-cache-operator#163 · 1 comentario ·
Los mantenedores suelen responder en 1 día
-
[submission]Abiertosubmission
Dificultad 1/5 Menos de una hora Aptitud para principiantes 65/100
leanprover/lean-eval-submissions#1852 ·
Los mantenedores suelen responder en 1 día