python / python/cpython

multiprocessing pool apply_async failure due to unable to pickle local dynamically created function

Aperta
#96,464 0 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

topic-multiprocessing type-bug
Lingua principale
Python
Stelle
77.2k
Fork
36k
Metriche di merge delle PR
Metriche PR in attesa

Descrizione

Bug report

This bug report is about multiprocessing module in standard library. When submit a local dynamically created function to pool executor in multiprocessing.Pool, it will failed. The details will be shown below.

Your environment

  • CPython versions tested on: 3.10.6
  • Operating system and architecture: Linux, x86_64

Details

Consider the following code snippet:

import multiprocessing
import pickle

pool = multiprocessing.Pool(4)


def error_callback(e):
    raise e


def go():
    for i in range(40):
        def hello(j):
            print(f"hello {i} {j}")

        pool.apply_async(hello, (i + 1,), error_callback=error_callback)

    pool.close()
    pool.join()


go()


The error is below:

Exception in thread Thread-2 (_handle_tasks):
Traceback (most recent call last):
  File "/usr/lib/python3.10/threading.py", line 1016, in _bootstrap_inner
    self.run()
  File "/usr/lib/python3.10/threading.py", line 953, in run
    self._target(*self._args, **self._kwargs)
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 544, in _handle_tasks
    cache[job]._set(idx, (False, e))
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 781, in _set
    self._error_callback(self._value)
  File "/.../xxx.py", line 8, in error_callback
    raise e
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 540, in _handle_tasks
    put(task)
  File "/usr/lib/python3.10/multiprocessing/connection.py", line 211, in send
    self._send_bytes(_ForkingPickler.dumps(obj))
  File "/usr/lib/python3.10/multiprocessing/reduction.py", line 51, in dumps
    cls(buf, protocol).dump(obj)
AttributeError: Can't pickle local object 'go.<locals>.hello'

analysis

Currently, the function in multiprocessing utilize pickle to transfer object between different process. When pickle.dumps() is applied to a function, only its reference information will be dumped. As a result, only global function which is defined in both sender and receiver end with same reference information will works.

The code object of function will not be dumped.

def hello():
    def hi(s):
        print(f"hi {s}")

    return pickle.dumps(hi)

hello()

Will cause same error:

Traceback (most recent call last):
  File "/.../.venv/lib/python3.10/site-packages/IPython/core/interactiveshell.py", line 3398, in run_code
    exec(code_obj, self.user_global_ns, self.user_ns)
  File "<ipython-input-8-a75d7781aaeb>", line 1, in <cell line: 1>
    hello()
  File "<ipython-input-7-3e50c99ad472>", line 4, in hello
    return pickle.dumps(hi)
AttributeError: Can't pickle local object 'hello.<locals>.hi'

potential solution

I am not meant to modify the behavior of pickle.dumps, but multiprocessing is supposed to utilize a enhanced version of pickle.

It is believed that the security issue is not significant in multiprocessing, because the serialized object which will be load at receiver end has already been executed in sender end. And the permissions of sender and receiver process is strictly the same.

So just for inspiration, marshal which can serialize code object can be mentioned here. Of course, a more complicated serializing function should be construct in multiprocessing which can rebuild a function at receiver end from scratch for local dynamically created function in sender end.

Guida per i contributori

Apri la guida per i contributori

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Direzione di ricerca

Riprodurre l'esempio fornito con multiprocessing.Pool e analizzare il percorso di errore in multiprocessing/pool.py, connection.py e reduction.py. Confrontarlo con la gestione delle funzioni locali da parte di pickle; l'issue propone una modifica ampia alla serializzazione, ma non definisce il comportamento supportato né un criterio concreto di completamento, quindi è necessario confermare prima l'ambito con i maintainer.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
python
Ambito
distributed-systems
Tipo di issue
Bug
Difficoltà
5/5
Tempo stimato
Più di una settimana
Stato di attività
Ferma
Chiarezza
Da chiarire
Idoneità per principianti
25/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.