From 4fcf7dc0fdbdaaeef61f11b8beedb121a21f180d Mon Sep 17 00:00:00 2001 From: Sam Cranford Date: Wed, 23 Sep 2026 20:34:14 +0000 Subject: [PATCH 1/3] FIX: Deadlock caused by fork in multiprocess --- podpac/core/managers/multi_process.py | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/podpac/core/managers/multi_process.py b/podpac/core/managers/multi_process.py index d7b8135f..1c56ef0c 100644 --- a/podpac/core/managers/multi_process.py +++ b/podpac/core/managers/multi_process.py @@ -1,7 +1,6 @@ from __future__ import division, unicode_literals, print_function, absolute_import -from multiprocessing import Process as mpProcess -from multiprocessing import Queue +import multiprocessing import traitlets as tl import logging @@ -32,12 +31,23 @@ def _f(definition, coords, q, outputkw): class Process(Node): """ Source node will be evaluated in another process, and it is blocking! + + Attributes + ---------- + start_method : str + The `multiprocessing` start method to use: "fork", "spawn", or "forkserver". + Default is "spawn". + + .. warning:: + "fork" is not safe: it can deadlock or crash if the parent process is + multithreaded, and neither "fork" nor "forkserver" is available on Windows. """ source = NodeTrait().tag(attr=True) output_format = tl.Dict(None, allow_none=True).tag(attr=True) timeout = tl.Int(None, allow_none=True) block = tl.Bool(True) + start_method = tl.Enum(["fork", "spawn", "forkserver"], default_value="spawn") @property def outputs(self): @@ -48,8 +58,9 @@ def eval(self, coordinates, **kwargs): # noqa: A003 definition = self.source.json coords = coordinates.json - q = Queue() - process = mpProcess(target=_f, args=(definition, coords, q, self.output_format)) + ctx = multiprocessing.get_context(self.start_method) + q = ctx.Queue() + process = ctx.Process(target=_f, args=(definition, coords, q, self.output_format)) # type: ignore[attr-defined] process.daemon = True _log.debug("Starting process.") process.start() From c1dd7ced7ac7ccf136619dfede3d124b387c2e80 Mon Sep 17 00:00:00 2001 From: Sam Cranford Date: Thu, 24 Sep 2026 09:20:46 -0400 Subject: [PATCH 2/3] Revert "FIX: Deadlock caused by fork in multiprocess" This reverts commit 4fcf7dc0fdbdaaeef61f11b8beedb121a21f180d. --- podpac/core/managers/multi_process.py | 19 ++++--------------- 1 file changed, 4 insertions(+), 15 deletions(-) diff --git a/podpac/core/managers/multi_process.py b/podpac/core/managers/multi_process.py index 1c56ef0c..d7b8135f 100644 --- a/podpac/core/managers/multi_process.py +++ b/podpac/core/managers/multi_process.py @@ -1,6 +1,7 @@ from __future__ import division, unicode_literals, print_function, absolute_import -import multiprocessing +from multiprocessing import Process as mpProcess +from multiprocessing import Queue import traitlets as tl import logging @@ -31,23 +32,12 @@ def _f(definition, coords, q, outputkw): class Process(Node): """ Source node will be evaluated in another process, and it is blocking! - - Attributes - ---------- - start_method : str - The `multiprocessing` start method to use: "fork", "spawn", or "forkserver". - Default is "spawn". - - .. warning:: - "fork" is not safe: it can deadlock or crash if the parent process is - multithreaded, and neither "fork" nor "forkserver" is available on Windows. """ source = NodeTrait().tag(attr=True) output_format = tl.Dict(None, allow_none=True).tag(attr=True) timeout = tl.Int(None, allow_none=True) block = tl.Bool(True) - start_method = tl.Enum(["fork", "spawn", "forkserver"], default_value="spawn") @property def outputs(self): @@ -58,9 +48,8 @@ def eval(self, coordinates, **kwargs): # noqa: A003 definition = self.source.json coords = coordinates.json - ctx = multiprocessing.get_context(self.start_method) - q = ctx.Queue() - process = ctx.Process(target=_f, args=(definition, coords, q, self.output_format)) # type: ignore[attr-defined] + q = Queue() + process = mpProcess(target=_f, args=(definition, coords, q, self.output_format)) process.daemon = True _log.debug("Starting process.") process.start() From 6fcd3d9f3cc17de191f715146c64747bbcd36292 Mon Sep 17 00:00:00 2001 From: Sam Cranford Date: Thu, 24 Sep 2026 16:25:53 +0000 Subject: [PATCH 3/3] Deprecate process node and skip related unit tests --- podpac/core/managers/multi_process.py | 9 +++++++++ podpac/core/managers/test/test_multiprocess.py | 2 ++ podpac/core/managers/test/test_parallel.py | 6 ++++++ 3 files changed, 17 insertions(+) diff --git a/podpac/core/managers/multi_process.py b/podpac/core/managers/multi_process.py index d7b8135f..19fb2261 100644 --- a/podpac/core/managers/multi_process.py +++ b/podpac/core/managers/multi_process.py @@ -4,6 +4,7 @@ from multiprocessing import Queue import traitlets as tl import logging +import warnings from podpac.core.node import Node from podpac.core.utils import NodeTrait @@ -39,6 +40,14 @@ class Process(Node): timeout = tl.Int(None, allow_none=True) block = tl.Bool(True) + def _first_init(self, **kwargs): + warnings.warn( + "Process node is deprecated and will be removed in a future version of podpac.", + DeprecationWarning, + stacklevel=1, + ) + return super(Process, self)._first_init(**kwargs) + @property def outputs(self): return self.source.outputs diff --git a/podpac/core/managers/test/test_multiprocess.py b/podpac/core/managers/test/test_multiprocess.py index 94842915..cd33da8b 100644 --- a/podpac/core/managers/test/test_multiprocess.py +++ b/podpac/core/managers/test/test_multiprocess.py @@ -1,3 +1,4 @@ +import pytest import numpy as np from multiprocessing import Queue @@ -7,6 +8,7 @@ from podpac.core.managers.multi_process import Process, _f +@pytest.mark.skip(reason="Process node is deprecated and will be removed.") class TestProcess(object): def test_mp_results_the_same(self): coords = Coordinates([[1, 2, 3, 4, 5]], ["time"]) diff --git a/podpac/core/managers/test/test_parallel.py b/podpac/core/managers/test/test_parallel.py index 5a229f10..e62b93c5 100644 --- a/podpac/core/managers/test/test_parallel.py +++ b/podpac/core/managers/test/test_parallel.py @@ -4,6 +4,7 @@ import numpy as np import tempfile import logging +import pytest from podpac.core.coordinates import Coordinates from podpac.core.algorithm.utility import CoordData @@ -35,6 +36,7 @@ def test_parallel_multi_thread_compute_fill_output2(self): np.testing.assert_array_equal(o, o_p) + @pytest.mark.skip(reason="Process node is deprecated and will be removed.") def test_parallel_process(self): node = Process(source=CoordData(coord_name="time")) coords = Coordinates([[1, 2, 3, 4, 5]], ["time"]) @@ -49,6 +51,7 @@ def test_parallel_process(self): class TestParallelAsync(object): + @pytest.mark.skip(reason="Process node is deprecated and will be removed.") def test_parallel_process_async(self): node = Process(source=CoordData(coord_name="time")) # , block=False) coords = Coordinates([[1, 2, 3, 4, 5]], ["time"]) @@ -59,6 +62,7 @@ def test_parallel_process_async(self): class TestParallelOutputZarr(object): + @pytest.mark.skip(reason="Process node is deprecated and will be removed.") def test_parallel_process_zarr(self): # Can't use tempfile.TemporaryDirectory because multiple processess need access to dir tmpdir = os.path.join(tempfile.gettempdir(), "test_parallel_process_zarr.zarr") @@ -74,6 +78,7 @@ def test_parallel_process_zarr(self): shutil.rmtree(tmpdir) + @pytest.mark.skip(reason="Process node is deprecated and will be removed.") def test_parallel_process_zarr_async(self): # Can't use tempfile.TemporaryDirectory because multiple processess need access to dir tmpdir = os.path.join(tempfile.gettempdir(), "test_parallel_process_zarr_async.zarr") @@ -89,6 +94,7 @@ def test_parallel_process_zarr_async(self): shutil.rmtree(tmpdir) + @pytest.mark.skip(reason="Process node is deprecated and will be removed.") def test_parallel_process_zarr_async_starti(self): # Can't use tempfile.TemporaryDirectory because multiple processess need access to dir tmpdir = os.path.join(tempfile.gettempdir(), "test_parallel_process_zarr_async_starti.zarr")