python / python/cpython

multiprocessing.Barrier does not return if waiting process is terminated (tested on windows)

オープン
#123,899 コメント 14 件 リアクション 0 件 担当者 1 名 GitHub で見る

@Zheaoli がすでに取り組んでいます。

2024年9月10日 から。

topic-multiprocessing type-bug
主要言語
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

コントリビューションガイド

コントリビューションガイドを開く

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

評価

この issue はまだ評価されていません。

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。