multiprocessing.Barrier does not return if waiting process is terminated (tested on windows)
@Zheaoli 已經在處理了。
開始於 2024年9月10日。
- 主要語言
- Python
- 星號
- 77.2k
- 分支
- 35.9k
- PR 合併指標
- PR 指標待擷取
描述
Bug report
Bug description:
I'm using multiprocessing.Barrier to synchronize processes. The individual processes might get terminated/killed by a GUI user interaction.
In case that a process, that already entered the waiting state, is killed, the barrier will never be released, a timeout is not taken into account.
When debugging this issue I realized that it is cause by https://github.com/python/cpython/blob/b52de7e02dba9e1f176d6d978d782fbd0509311e/Lib/multiprocessing/synchronize.py#L297
self._woken_count is not released by the terminated process and no timeout is specified in the call to self._woken_count.acquire()
If I add a (hardcoded...) timeout to self._woken_count.acquire(True, 2.0) the program continues as expected.
I hope that there a better approaches to fix this than using a hardcoded timeout...
Example script to reproduce error.
import logging
import time
import sys
import multiprocessing as mp
import multiprocessing.synchronize
from pathlib import Path
def process_main(process_index, barrier: multiprocessing.synchronize.Barrier):
logging.debug('Process %d: Started', process_index)
time.sleep(0.5)
logging.debug('Process %d: Waiting for barrier', process_index)
barrier.wait(timeout=5.0)
logging.debug('Process %d: Barrier passed', process_index)
time.sleep(0.5)
logging.debug('Process %d: Terminated', process_index)
if __name__ == '__main__':
# Set up logging
logfile_name = Path(__file__).with_suffix('.log').name
logging.basicConfig(
level=logging.DEBUG,
format='%(asctime)s %(levelname)s:%(name)s %(message)s',
handlers=[
logging.FileHandler(logfile_name, mode='w'),
logging.StreamHandler(sys.stdout)
]
)
instance_count = 4
barrier = mp.Barrier(instance_count)
processes = []
for i in range(instance_count):
runner_process = mp.Process(
target=process_main, args=(i, barrier), daemon=True)
processes.append(runner_process)
for i, process in enumerate(processes):
logging.debug('Starting process %d', i)
process.start()
time.sleep(0.200)
# Terminate already waiting process
logging.debug('Killing process 0')
processes[0].kill()
for process in processes:
process.join()
CPython versions tested on:
3.12
Operating systems tested on:
Windows
Linked PRs
- gh-125578
貢獻指南
從這裡開始
- 先讀完整個 Issue,再讀專案的貢獻指南。
- 在 Issue 下留言說明你要接手 —— 這能避免兩個人做同樣的事。
- Fork 儲存庫,在一個分支上完成修改。
- 送出 Pull Request,並在描述裡引用這個 Issue 編號。
評估
這個 Issue 還沒有評估資料。