Hacktoberfest 2026: những issue maintainer đã đánh dấu cho tháng Mười, đang mở và phù hợp người mới. Xem issue Hacktoberfest

Scheduler deadlock when using SSHCluster due to stderr blocking

Đang mở
#9,033 3 bình luận 0 reaction 0 người được giao Xem trên GitHub

Chưa có ai nhận issue này.

Đánh giá

Độ khó
3/5
Thời gian dự kiến
1-2 ngày
Mức phù hợp với người mới
52/100
Loại issue
Lỗi
Độ rõ ràng
Đặc tả rõ ràng
Mức độ hoạt động
Đình trệ
Công nghệ
python
Lĩnh vực
distributed-systems

Hướng nghiên cứu

Bắt đầu bằng cách tái hiện ví dụ được cung cấp, sau đó kiểm tra distributed/deploy/ssh.py, đặc biệt là các lớp Scheduler và Worker. Xác minh cách stdout và stderr từ xa được xử lý, đồng thời sử dụng ví dụ để xác nhận rằng cluster không còn bị deadlock khi pipe stderr của scheduler bị đầy.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Mô tả

bug

Describe the issue:
When using SSHCluster eventually my cluster deadlocks (and dashboard stops responding).
Scheduler stack shows it's stuck reporting an asyncio unhandled task exception. My real-world cluster will deadlock pretty consistently after an hour or so of normal task execution. I'm not using run_on_scheduler in the real world, it's just a convenient way to trigger the issue quickly.

The root of the problem seems to be that both distributed.deploy.ssh.Scheduler and distributed.deploy.ssh.Worker stop polling the stdout/stderr pipes shortly after startup, which eventually causes the remote processes pipes to fill up and block if anything in the process writes to them.

Minimal Complete Verifiable Example:

from dask.distributed import (
    SSHCluster,
    worker_client,
    fire_and_forget,
)

async def _raise():
    raise RuntimeError('broken')

def _thing(n):
    with worker_client() as client:
        fire_and_forget(client.run_on_scheduler(_raise, wait=False))
        fire_and_forget(client.submit(_thing, n=n+1))


def main():
    cluster = SSHCluster(
        ['localhost', 'localhost'],
        connect_options=dict(known_hosts=None),
        worker_options=dict(nthreads=1, n_workers=1),
        scheduler_options=dict(dashboard=True),
    )
    with cluster.get_client() as client:
        client.upload_file(__file__)
        client.submit(_thing, n=1).result()
        input('waiting')

if __name__ == '__main__':
    main()

(run with python -c 'import deadlock; deadlock.main()')

When left running this example will eventually deadlock once the scheduler stderr pipe buffer fills.

Below is the py-spy stack trace:

Thread 4175110 (idle): "MainThread"
    emit (logging/__init__.py:1113)
    handle (logging/__init__.py:978)
    callHandlers (logging/__init__.py:1714)
    handle (logging/__init__.py:1644)
    _log (logging/__init__.py:1634)
    error (logging/__init__.py:1518)
    default_exception_handler (asyncio/base_events.py:1785)
    call_exception_handler (asyncio/base_events.py:1811)
    _run_once (asyncio/base_events.py:1937)
    run_forever (asyncio/base_events.py:608)
    run_until_complete (asyncio/base_events.py:641)
    run (asyncio/runners.py:118)
    asyncio_run (distributed/compatibility.py:204)
    main (distributed/cli/dask_spec.py:63)
    invoke (click/core.py:788)
    invoke (click/core.py:1443)
    main (click/core.py:1082)
    __call__ (click/core.py:1161)
    <module> (distributed/cli/dask_spec.py:67)
    _run_code (<frozen runpy>:88)
    _run_module_as_main (<frozen runpy>:198)

And corresponding stack from the kernel side showing we're blocked in a pipe:

[<0>] pipe_wait+0x6f/0xc0
[<0>] pipe_write+0x17b/0x470
[<0>] new_sync_write+0x125/0x1c0
[<0>] __vfs_write+0x29/0x40
[<0>] vfs_write+0xb9/0x1a0
[<0>] ksys_write+0x67/0xe0
[<0>] __x64_sys_write+0x1a/0x20
[<0>] do_syscall_64+0x57/0x190
[<0>] entry_SYSCALL_64_after_hwframe+0x44/0xa9

Anything else we need to know?:

Environment:

  • Dask version: 2024.2.1
  • Python version: 3.11
  • Operating System: Ubuntu 20.04
  • Install method (conda, pip, source): conda
Ngôn ngữ chính
Python
Star
1.7k
Fork
778
Chỉ số merge pull request
Không có pull request nào được merge trong 30 ngày

Chuẩn bị môi trường

Bắt đầu từ đâu

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Issue khác của dask/distributed

Tất cả issue của dask/distributed

Issue tương tự

Thêm issue về Python

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.