python / python/cpython

multiprocessing race condition on flushing stdout, deadlocks child on exit

Offen
#91,776 12 Kommentare 3 Reaktionen 0 zugewiesene Personen Auf GitHub ansehen

Dieses Issue hat noch niemand übernommen.

topic-multiprocessing type-bug
Vorherrschende Sprache
Python
Sterne
77.2k
Forks
36k
Ø Merge
1 T. 9 Std.
Gemergte PRs (30 T.)
558

Beschreibung

I experienced deadlocks when using the logging example for multiprocessing, where a QueueHandler is used.

I found that sometimes a multiprocessing process is not terminating, when it has put an element in a Queue,
even if the parent process runs a thread to empty the queue and successfully retrieved the item.

The worker looks like this:

def worker(q):
    q.put(current_process().name)
    return

and while this works:

workers = []
for i in range(15):
    p = Process(target=worker, args=(logq,), name=f"Worker {i+1}")
    workers.append(p)

for p in workers:
    p.start()
    time.sleep(.1)

removing the sleep leads to a very high propability of deadlocking when I then try to join the processes, e.g,:

for w in workers:
    print(f'trying to join on {w.name}, alive={w.is_alive()}, exitcode={w.exitcode}', w.name, w.is_alive(), w.exitcode)
    w.join()

The Process is still marked as alive, hitting Ctrl+C gives this:

trying to join on Worker 14, alive=True, exitcode=None Worker 14 True None

^CTraceback (most recent call last):
  File "/home/fls/pybug/deadlock.py", line 43, in <module>
    w.join()
  File "/usr/lib/python3.10/multiprocessing/process.py", line 149, in join
    res = self._popen.wait(timeout)
  File "/usr/lib/python3.10/multiprocessing/popen_fork.py", line 43, in wait
    return self.poll(os.WNOHANG if timeout == 0.0 else 0)
  File "/usr/lib/python3.10/multiprocessing/popen_fork.py", line 27, in poll
    pid, sts = os.waitpid(self.pid, flag)
KeyboardInterrupt

Can reproduce using Python 3.9.10 and 3.10.4 on Linux:

import time
import threading
from multiprocessing import Process, Queue, current_process

def logger_thread(q: Queue):
    while True:
        record = q.get()
        if record is None:
            break
        print("logrecord: ", record)

def worker(q):
    q.put(current_process().name)
    return

logq = Queue()

lp = threading.Thread(target=logger_thread, args=(logq,), daemon=True)
lp.start()

workers = []
for i in range(15):
    p = Process(target=worker, args=(logq,), name=f"Worker {i+1}")
    workers.append(p)

print("starting workers")

for p in workers:
    p.start()
    # no deadlock when added:
    # time.sleep(.1)

print("waiting a bit")
time.sleep(1)
print("trying to join workers")

for w in workers:
    print(f'trying to join on {w.name}, alive={w.is_alive()}, exitcode={w.exitcode}', w.name, w.is_alive(), w.exitcode)
    w.join()
    print(f'joined on {w.name}', w.name)

logq.put(None)
lp.join()

Beitragsleitfaden

Beitragsleitfaden öffnen

Erste Schritte

  1. Lies das ganze Issue und danach den Beitragsleitfaden des Projekts.
  2. Schreib ins Issue, dass du es übernimmst — das erspart doppelte Arbeit.
  3. Forke das Repository und arbeite in einem Branch.
  4. Öffne einen Pull Request, der die Issue-Nummer nennt.

Rechercherichtung

Beginnen Sie damit, die bereitgestellte multiprocessing-Reproduktion unter Linux mit den gemeldeten Python-Versionen auszuführen, und untersuchen Sie anschließend die multiprocessing Queue und den Pfad zum Beenden der Kindprozesse, einschließlich der in der Traceback angezeigten Pfade multiprocessing/process.py und multiprocessing/popen_fork.py. Als abgeschlossen gilt die Arbeit, wenn die Worker zuverlässig joinen, nachdem Queue-Elemente abgerufen wurden, und eine Regressionstestabdeckung für den reproduzierten Fall vorhanden ist.

Vom Indexierungsmodell aus dem Issue-Text verfasst.

Bewertung

Tech-Stack
python
Bereich
distributed-systems
Issue-Typ
Bug
Schwierigkeit
4/5
Geschätzter Aufwand
3-5 Tage
Aktivitätsstatus
Veraltet
Klarheit
Muss geklärt werden
Anfängerfreundlichkeit
35/100

Neue Issues direkt in Ihr Postfach

Eine kurze Übersicht über anfängerfreundliche GitHub-Issues.