dask / dask/distributed

Worker profile limited to a short timespan

Aberta
#8,653 3 comentários 0 reações 0 responsáveis Ver no GitHub
bug
Linguagem predominante
Python
Estrelas
1.7k
Forks
778
Merge médio
2h 50min
PRs com merge (30d)
3

Descrição

**Describe the issue**:
The worker profile has a limited span and older data seems to be lost. For example, with the minimal example below, the total CPU time is 24 hours, but the profile never contains more than 4 h 26 min. At one point during the run, the profile looks like this:
![image](https://github.com/dask/distributed/assets/52597883/532af79a-0158-4035-91f6-fa76f663229c)

30 minutes later, it looks like this:
![image](https://github.com/dask/distributed/assets/52597883/37702c5a-e092-4b4a-8572-ddfb17d3a209)

The previous data is completely gone, as can be seen by the "activity over time" graph at the bottom.

This issue has been occurring for several months, most recently with dask and distributed 2024.4.2. You can look at the full discussion on Discourse: https://dask.discourse.group/t/measuring-the-overall-profile-of-long-runs/1859/11

**Minimal Complete Verifiable Example**:

```python
import time
import logging

from dask.distributed import (
Client,
LocalCluster,
get_client,
as_completed,
performance_report,
)

NUM_PROCESSES = 16

logging.basicConfig(
level=logging.DEBUG,
format="%(asctime)s %(levelname)-8s %(message)s",
)

class DummyManager:
def run(self):
logging.info("Starting the manager")

jobs = list(range(1, 97))

client = get_client()
futs = []
for j in jobs:
futs.append(client.submit(self.job, j))

asc = as_completed(futs, with_results=True)
for fut, ret in asc:
logging.info(f"Processing future {str(fut)}: ret={str(ret)}")
if ret > 0:
logging.info(f"Launching a subjob with time {ret}")
asc.add(client.submit(self.job, ret))
fut.release()

def job(self, n):
time.sleep(60 * 15)

return 0

if __name__ == "__main__":
cluster = LocalCluster(
n_workers=1,
threads_per_worker=NUM_PROCESSES,
processes=False,
)

client = Client(cluster)
with performance_report(filename=f"dask-performance_{time.time():.0f}.html"):
try:
manager = DummyManager()
manager.run()
except KeyboardInterrupt:
logging.info("Stopping the job...")
cluster.close()
exit(0)
client.close()
cluster.close()
```

**Environment**:

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

Guia de contribuição

Abrir o guia de contribuição

Direção de pesquisa

Execute o exemplo mínimo de Python fornecido com Client, LocalCluster e performance_report e, em seguida, compare o perfil do worker e a visualização da atividade ao longo do tempo durante a execução longa. Rastreie como os dados de perfil são coletados e mantidos, usando a discussão vinculada no Discourse como contexto. Considera-se concluído quando a atividade anterior continua disponível e o perfil cobre toda a execução, em vez de parar por volta de 4 h 26 min.

Escrita pelo modelo de indexação a partir do texto da issue.

Avaliação

Stack de tecnologia
python
Domínio
distributed-systems, observability-sre
Tipo de issue
Bug
Dificuldade
4/5
Tempo estimado
3-5 dias
Status de atividade
Estagnada
Clareza
Precisa de esclarecimento
Facilidade para iniciantes
35/100

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.