Python Multiprocessing get hung

Viewed 317

Below is a min. reproducible code:

import multiprocessing

class E(Exception):
    def __init__(self, a1, a2):
        Exception.__init__(self, '{}{}'.format(a1, a2))

def f(_):
    raise E(1, 2)

multiprocessing.Pool(1).map(f, (1,))

This gives below error:

Exception in thread Thread-5:
Traceback (most recent call last):
  File "/usr/lib/python3.7/threading.py", line 917, in _bootstrap_inner
    self.run()
  File "/usr/lib/python3.7/threading.py", line 865, in run
    self._target(*self._args, **self._kwargs)
  File "/usr/lib/python3.7/multiprocessing/pool.py", line 496, in _handle_results
    task = get()
  File "/usr/lib/python3.7/multiprocessing/contrnection.py", line 251, in recv
    return _ForkingPickler.loads(buf.getbuffer())
TypeError: __init__() missing 1 required positional argument: 'a2'

Any way to fix this? Similar issue is mentioned here: https://bugs.python.org/issue39751

2 Answers

I think this has to do with what the __reduce__ function returns.

it is probably returning something like (self.__class__, self.args) basically saying that to recreate this object pass in everything stored in self.args which is equal to ("12",).

Since there is no what of taking in the one string and splitting it back into the original args you probably don't want to mess with formatting before storage and should just implement a __str__.

class E(Exception):
    def __init__(self, a1, a2):
        super().__init__(a1, a2)

    def __str__(self):
        return f"{self.args[0]}{self.args[1]}"

Edit: Added other possibilities I came up with but I don't recommend them.

# Change reduce
class E(Exception):
    def __init__(self, a1, a2):
        super().__init__(f"{a1}{a2}")

    def __reduce__(self, *args, **kwargs):
        return (self.__class__, (self.args[0], ""))

# Optional arg this one is similar to Booboo's answer.
class E(Exception):
    def __init__(self, a1, a2=""):
        # I am assuming you have a more complex format string when not using the toy example
        # If so you need an if here.
        if a2:
            super().__init__(f"{a1}{a2}")
        else:
            super().__init__(a1)
    

I tried to resolve this by defining pickle functions __getstate__ and __getstate__, but they are not even getting invoked. However, the following should sidestep the problem until the issue, which I do believe is a problem with the reduction.ForkingPickler class, is resolved:

import multiprocessing

class E(Exception):
    def __init__(self, *args):
        if len(args) == 2:
            # Normal instantiation:
            Exception.__init__(self, '{}{}'.format(args[0], args[1]))
        else:
            # We are being pickled by the reduction.ForkingPickler:
            assert(len(args) == 1)
            Exception.__init__(self, args[0])

def f(_):
    raise E(1, 2)

if __name__ == '__main__': # Required for Windows
    try:
        multiprocessing.Pool(1).map(f, (1,))
    except E as e:
        print('Got E Exception:', e)

Prints:

Got E Exception: 12
Related