dask / dask/distributed

Worker profile limited to a short timespan

未關閉
#8,653 3 則留言 0 個 reaction 已指派 0 人 在 GitHub 檢視
bug
主要語言
Python
星號
1.7k
分支
778
平均合併
2 小時 50 分鐘
30 天內合併 PR
3

描述

**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

貢獻指南

開啟貢獻指南

研究方向

執行提供的、使用 Client、LocalCluster 和 performance_report 的 Python 最小範例,然後比較長時間執行期間的 worker profile 和 activity-over-time view。以連結的 Discourse 討論作為背景,追蹤 profiling 資料如何收集與保留。完成的定義是較早的活動仍然可用,且 profile 涵蓋整個執行過程,而不是在約 4 h 26 min 時停止。

由索引模型根據 Issue 內容生成。

評估

技術堆疊
python
領域
distributed-systems, observability-sre
Issue 類型
缺陷
難度
4/5
預估耗時
3-5 天
活躍度
停滯
描述清晰度
需要釐清
新手友好度
35/100

把新 issue 寄到你的電子郵件信箱

精選適合新手參與的 GitHub issue 摘要。